1#![allow(unused_imports)]
15use anyhow::Context;
16use async_trait::async_trait;
17use derive_builder::Builder;
18use rust_decimal::prelude::*;
19use serde::{Deserialize, Serialize};
20use serde_json::Value;
21use std::{collections::BTreeMap, sync::Arc};
22
23use crate::common::{
24 errors::WebsocketError,
25 models::{ParamBuildError, WebsocketApiResponse},
26 utils::remove_empty_value,
27 websocket::{WebsocketApi, WebsocketMessageSendOptions},
28};
29use crate::spot::websocket_api::models;
30
31#[async_trait]
32pub trait MarketApi: Send + Sync {
33 async fn avg_price(
34 &self,
35 params: AvgPriceParams,
36 ) -> anyhow::Result<WebsocketApiResponse<Box<models::AvgPriceResponseResult>>>;
37 async fn block_trades_historical(
38 &self,
39 params: BlockTradesHistoricalParams,
40 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::BlockTradesHistoricalResponseResultInner>>>;
41 async fn depth(
42 &self,
43 params: DepthParams,
44 ) -> anyhow::Result<WebsocketApiResponse<Box<models::DepthResponseResult>>>;
45 async fn klines(
46 &self,
47 params: KlinesParams,
48 ) -> anyhow::Result<WebsocketApiResponse<Vec<Vec<models::KlinesResponseResultInnerInner>>>>;
49 async fn reference_price(
50 &self,
51 params: ReferencePriceParams,
52 ) -> anyhow::Result<WebsocketApiResponse<Box<models::ReferencePriceResponseResult>>>;
53 async fn reference_price_calculation(
54 &self,
55 params: ReferencePriceCalculationParams,
56 ) -> anyhow::Result<WebsocketApiResponse<Box<models::ReferencePriceCalculationResponseResult>>>;
57 async fn ticker(
58 &self,
59 params: TickerParams,
60 ) -> anyhow::Result<WebsocketApiResponse<models::TickerResponse>>;
61 async fn ticker24hr(
62 &self,
63 params: Ticker24hrParams,
64 ) -> anyhow::Result<WebsocketApiResponse<models::Ticker24hrResponse>>;
65 async fn ticker_book(
66 &self,
67 params: TickerBookParams,
68 ) -> anyhow::Result<WebsocketApiResponse<models::TickerBookResponse>>;
69 async fn ticker_price(
70 &self,
71 params: TickerPriceParams,
72 ) -> anyhow::Result<WebsocketApiResponse<models::TickerPriceResponse>>;
73 async fn ticker_trading_day(
74 &self,
75 params: TickerTradingDayParams,
76 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::TickerTradingDayResponseResultInner>>>;
77 async fn trades_aggregate(
78 &self,
79 params: TradesAggregateParams,
80 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::TradesAggregateResponseResultInner>>>;
81 async fn trades_historical(
82 &self,
83 params: TradesHistoricalParams,
84 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::TradesHistoricalResponseResultInner>>>;
85 async fn trades_recent(
86 &self,
87 params: TradesRecentParams,
88 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::TradesRecentResponseResultInner>>>;
89 async fn ui_klines(
90 &self,
91 params: UiKlinesParams,
92 ) -> anyhow::Result<WebsocketApiResponse<Vec<Vec<models::KlinesResponseResultInnerInner>>>>;
93}
94
95#[derive(Clone)]
96pub struct MarketApiClient {
97 websocket_api_base: Arc<WebsocketApi>,
98}
99
100impl MarketApiClient {
101 pub fn new(websocket_api_base: Arc<WebsocketApi>) -> Self {
102 Self { websocket_api_base }
103 }
104}
105
106#[allow(non_camel_case_types)]
107#[derive(Debug, Clone, Serialize, Deserialize)]
108pub enum DepthSymbolStatusEnum {
109 #[serde(rename = "TRADING")]
110 Trading,
111 #[serde(rename = "HALT")]
112 Halt,
113 #[serde(rename = "BREAK")]
114 Break,
115}
116
117impl DepthSymbolStatusEnum {
118 #[must_use]
119 pub fn as_str(&self) -> &'static str {
120 match self {
121 Self::Trading => "TRADING",
122 Self::Halt => "HALT",
123 Self::Break => "BREAK",
124 }
125 }
126}
127
128impl std::str::FromStr for DepthSymbolStatusEnum {
129 type Err = Box<dyn std::error::Error + Send + Sync>;
130
131 fn from_str(s: &str) -> Result<Self, Self::Err> {
132 match s {
133 "TRADING" => Ok(Self::Trading),
134 "HALT" => Ok(Self::Halt),
135 "BREAK" => Ok(Self::Break),
136 other => Err(format!("invalid DepthSymbolStatusEnum: {}", other).into()),
137 }
138 }
139}
140
141#[allow(non_camel_case_types)]
142#[derive(Debug, Clone, Serialize, Deserialize)]
143pub enum KlinesIntervalEnum {
144 #[serde(rename = "1s")]
145 Interval1s,
146 #[serde(rename = "1m")]
147 Interval1m,
148 #[serde(rename = "3m")]
149 Interval3m,
150 #[serde(rename = "5m")]
151 Interval5m,
152 #[serde(rename = "15m")]
153 Interval15m,
154 #[serde(rename = "30m")]
155 Interval30m,
156 #[serde(rename = "1h")]
157 Interval1h,
158 #[serde(rename = "2h")]
159 Interval2h,
160 #[serde(rename = "4h")]
161 Interval4h,
162 #[serde(rename = "6h")]
163 Interval6h,
164 #[serde(rename = "8h")]
165 Interval8h,
166 #[serde(rename = "12h")]
167 Interval12h,
168 #[serde(rename = "1d")]
169 Interval1d,
170 #[serde(rename = "3d")]
171 Interval3d,
172 #[serde(rename = "1w")]
173 Interval1w,
174 #[serde(rename = "1M")]
175 Interval1M,
176}
177
178impl KlinesIntervalEnum {
179 #[must_use]
180 pub fn as_str(&self) -> &'static str {
181 match self {
182 Self::Interval1s => "1s",
183 Self::Interval1m => "1m",
184 Self::Interval3m => "3m",
185 Self::Interval5m => "5m",
186 Self::Interval15m => "15m",
187 Self::Interval30m => "30m",
188 Self::Interval1h => "1h",
189 Self::Interval2h => "2h",
190 Self::Interval4h => "4h",
191 Self::Interval6h => "6h",
192 Self::Interval8h => "8h",
193 Self::Interval12h => "12h",
194 Self::Interval1d => "1d",
195 Self::Interval3d => "3d",
196 Self::Interval1w => "1w",
197 Self::Interval1M => "1M",
198 }
199 }
200}
201
202impl std::str::FromStr for KlinesIntervalEnum {
203 type Err = Box<dyn std::error::Error + Send + Sync>;
204
205 fn from_str(s: &str) -> Result<Self, Self::Err> {
206 match s {
207 "1s" => Ok(Self::Interval1s),
208 "1m" => Ok(Self::Interval1m),
209 "3m" => Ok(Self::Interval3m),
210 "5m" => Ok(Self::Interval5m),
211 "15m" => Ok(Self::Interval15m),
212 "30m" => Ok(Self::Interval30m),
213 "1h" => Ok(Self::Interval1h),
214 "2h" => Ok(Self::Interval2h),
215 "4h" => Ok(Self::Interval4h),
216 "6h" => Ok(Self::Interval6h),
217 "8h" => Ok(Self::Interval8h),
218 "12h" => Ok(Self::Interval12h),
219 "1d" => Ok(Self::Interval1d),
220 "3d" => Ok(Self::Interval3d),
221 "1w" => Ok(Self::Interval1w),
222 "1M" => Ok(Self::Interval1M),
223 other => Err(format!("invalid KlinesIntervalEnum: {}", other).into()),
224 }
225 }
226}
227
228#[allow(non_camel_case_types)]
229#[derive(Debug, Clone, Serialize, Deserialize)]
230pub enum ReferencePriceCalculationSymbolStatusEnum {
231 #[serde(rename = "TRADING")]
232 Trading,
233 #[serde(rename = "HALT")]
234 Halt,
235 #[serde(rename = "BREAK")]
236 Break,
237}
238
239impl ReferencePriceCalculationSymbolStatusEnum {
240 #[must_use]
241 pub fn as_str(&self) -> &'static str {
242 match self {
243 Self::Trading => "TRADING",
244 Self::Halt => "HALT",
245 Self::Break => "BREAK",
246 }
247 }
248}
249
250impl std::str::FromStr for ReferencePriceCalculationSymbolStatusEnum {
251 type Err = Box<dyn std::error::Error + Send + Sync>;
252
253 fn from_str(s: &str) -> Result<Self, Self::Err> {
254 match s {
255 "TRADING" => Ok(Self::Trading),
256 "HALT" => Ok(Self::Halt),
257 "BREAK" => Ok(Self::Break),
258 other => Err(format!(
259 "invalid ReferencePriceCalculationSymbolStatusEnum: {}",
260 other
261 )
262 .into()),
263 }
264 }
265}
266
267#[allow(non_camel_case_types)]
268#[derive(Debug, Clone, Serialize, Deserialize)]
269pub enum TickerTypeEnum {
270 #[serde(rename = "FULL")]
271 Full,
272 #[serde(rename = "MINI")]
273 Mini,
274}
275
276impl TickerTypeEnum {
277 #[must_use]
278 pub fn as_str(&self) -> &'static str {
279 match self {
280 Self::Full => "FULL",
281 Self::Mini => "MINI",
282 }
283 }
284}
285
286impl std::str::FromStr for TickerTypeEnum {
287 type Err = Box<dyn std::error::Error + Send + Sync>;
288
289 fn from_str(s: &str) -> Result<Self, Self::Err> {
290 match s {
291 "FULL" => Ok(Self::Full),
292 "MINI" => Ok(Self::Mini),
293 other => Err(format!("invalid TickerTypeEnum: {}", other).into()),
294 }
295 }
296}
297
298#[allow(non_camel_case_types)]
299#[derive(Debug, Clone, Serialize, Deserialize)]
300pub enum TickerWindowSizeEnum {
301 #[serde(rename = "1m")]
302 WindowSize1m,
303 #[serde(rename = "2m")]
304 WindowSize2m,
305 #[serde(rename = "3m")]
306 WindowSize3m,
307 #[serde(rename = "4m")]
308 WindowSize4m,
309 #[serde(rename = "5m")]
310 WindowSize5m,
311 #[serde(rename = "6m")]
312 WindowSize6m,
313 #[serde(rename = "7m")]
314 WindowSize7m,
315 #[serde(rename = "8m")]
316 WindowSize8m,
317 #[serde(rename = "9m")]
318 WindowSize9m,
319 #[serde(rename = "10m")]
320 WindowSize10m,
321 #[serde(rename = "11m")]
322 WindowSize11m,
323 #[serde(rename = "12m")]
324 WindowSize12m,
325 #[serde(rename = "13m")]
326 WindowSize13m,
327 #[serde(rename = "14m")]
328 WindowSize14m,
329 #[serde(rename = "15m")]
330 WindowSize15m,
331 #[serde(rename = "16m")]
332 WindowSize16m,
333 #[serde(rename = "17m")]
334 WindowSize17m,
335 #[serde(rename = "18m")]
336 WindowSize18m,
337 #[serde(rename = "19m")]
338 WindowSize19m,
339 #[serde(rename = "20m")]
340 WindowSize20m,
341 #[serde(rename = "21m")]
342 WindowSize21m,
343 #[serde(rename = "22m")]
344 WindowSize22m,
345 #[serde(rename = "23m")]
346 WindowSize23m,
347 #[serde(rename = "24m")]
348 WindowSize24m,
349 #[serde(rename = "25m")]
350 WindowSize25m,
351 #[serde(rename = "26m")]
352 WindowSize26m,
353 #[serde(rename = "27m")]
354 WindowSize27m,
355 #[serde(rename = "28m")]
356 WindowSize28m,
357 #[serde(rename = "29m")]
358 WindowSize29m,
359 #[serde(rename = "30m")]
360 WindowSize30m,
361 #[serde(rename = "31m")]
362 WindowSize31m,
363 #[serde(rename = "32m")]
364 WindowSize32m,
365 #[serde(rename = "33m")]
366 WindowSize33m,
367 #[serde(rename = "34m")]
368 WindowSize34m,
369 #[serde(rename = "35m")]
370 WindowSize35m,
371 #[serde(rename = "36m")]
372 WindowSize36m,
373 #[serde(rename = "37m")]
374 WindowSize37m,
375 #[serde(rename = "38m")]
376 WindowSize38m,
377 #[serde(rename = "39m")]
378 WindowSize39m,
379 #[serde(rename = "40m")]
380 WindowSize40m,
381 #[serde(rename = "41m")]
382 WindowSize41m,
383 #[serde(rename = "42m")]
384 WindowSize42m,
385 #[serde(rename = "43m")]
386 WindowSize43m,
387 #[serde(rename = "44m")]
388 WindowSize44m,
389 #[serde(rename = "45m")]
390 WindowSize45m,
391 #[serde(rename = "46m")]
392 WindowSize46m,
393 #[serde(rename = "47m")]
394 WindowSize47m,
395 #[serde(rename = "48m")]
396 WindowSize48m,
397 #[serde(rename = "49m")]
398 WindowSize49m,
399 #[serde(rename = "50m")]
400 WindowSize50m,
401 #[serde(rename = "51m")]
402 WindowSize51m,
403 #[serde(rename = "52m")]
404 WindowSize52m,
405 #[serde(rename = "53m")]
406 WindowSize53m,
407 #[serde(rename = "54m")]
408 WindowSize54m,
409 #[serde(rename = "55m")]
410 WindowSize55m,
411 #[serde(rename = "56m")]
412 WindowSize56m,
413 #[serde(rename = "57m")]
414 WindowSize57m,
415 #[serde(rename = "58m")]
416 WindowSize58m,
417 #[serde(rename = "59m")]
418 WindowSize59m,
419 #[serde(rename = "1h")]
420 WindowSize1h,
421 #[serde(rename = "2h")]
422 WindowSize2h,
423 #[serde(rename = "3h")]
424 WindowSize3h,
425 #[serde(rename = "4h")]
426 WindowSize4h,
427 #[serde(rename = "5h")]
428 WindowSize5h,
429 #[serde(rename = "6h")]
430 WindowSize6h,
431 #[serde(rename = "7h")]
432 WindowSize7h,
433 #[serde(rename = "8h")]
434 WindowSize8h,
435 #[serde(rename = "9h")]
436 WindowSize9h,
437 #[serde(rename = "10h")]
438 WindowSize10h,
439 #[serde(rename = "11h")]
440 WindowSize11h,
441 #[serde(rename = "12h")]
442 WindowSize12h,
443 #[serde(rename = "13h")]
444 WindowSize13h,
445 #[serde(rename = "14h")]
446 WindowSize14h,
447 #[serde(rename = "15h")]
448 WindowSize15h,
449 #[serde(rename = "16h")]
450 WindowSize16h,
451 #[serde(rename = "17h")]
452 WindowSize17h,
453 #[serde(rename = "18h")]
454 WindowSize18h,
455 #[serde(rename = "19h")]
456 WindowSize19h,
457 #[serde(rename = "20h")]
458 WindowSize20h,
459 #[serde(rename = "21h")]
460 WindowSize21h,
461 #[serde(rename = "22h")]
462 WindowSize22h,
463 #[serde(rename = "23h")]
464 WindowSize23h,
465 #[serde(rename = "1d")]
466 WindowSize1d,
467 #[serde(rename = "2d")]
468 WindowSize2d,
469 #[serde(rename = "3d")]
470 WindowSize3d,
471 #[serde(rename = "4d")]
472 WindowSize4d,
473 #[serde(rename = "5d")]
474 WindowSize5d,
475 #[serde(rename = "6d")]
476 WindowSize6d,
477 #[serde(rename = "7d")]
478 WindowSize7d,
479}
480
481impl TickerWindowSizeEnum {
482 #[must_use]
483 pub fn as_str(&self) -> &'static str {
484 match self {
485 Self::WindowSize1m => "1m",
486 Self::WindowSize2m => "2m",
487 Self::WindowSize3m => "3m",
488 Self::WindowSize4m => "4m",
489 Self::WindowSize5m => "5m",
490 Self::WindowSize6m => "6m",
491 Self::WindowSize7m => "7m",
492 Self::WindowSize8m => "8m",
493 Self::WindowSize9m => "9m",
494 Self::WindowSize10m => "10m",
495 Self::WindowSize11m => "11m",
496 Self::WindowSize12m => "12m",
497 Self::WindowSize13m => "13m",
498 Self::WindowSize14m => "14m",
499 Self::WindowSize15m => "15m",
500 Self::WindowSize16m => "16m",
501 Self::WindowSize17m => "17m",
502 Self::WindowSize18m => "18m",
503 Self::WindowSize19m => "19m",
504 Self::WindowSize20m => "20m",
505 Self::WindowSize21m => "21m",
506 Self::WindowSize22m => "22m",
507 Self::WindowSize23m => "23m",
508 Self::WindowSize24m => "24m",
509 Self::WindowSize25m => "25m",
510 Self::WindowSize26m => "26m",
511 Self::WindowSize27m => "27m",
512 Self::WindowSize28m => "28m",
513 Self::WindowSize29m => "29m",
514 Self::WindowSize30m => "30m",
515 Self::WindowSize31m => "31m",
516 Self::WindowSize32m => "32m",
517 Self::WindowSize33m => "33m",
518 Self::WindowSize34m => "34m",
519 Self::WindowSize35m => "35m",
520 Self::WindowSize36m => "36m",
521 Self::WindowSize37m => "37m",
522 Self::WindowSize38m => "38m",
523 Self::WindowSize39m => "39m",
524 Self::WindowSize40m => "40m",
525 Self::WindowSize41m => "41m",
526 Self::WindowSize42m => "42m",
527 Self::WindowSize43m => "43m",
528 Self::WindowSize44m => "44m",
529 Self::WindowSize45m => "45m",
530 Self::WindowSize46m => "46m",
531 Self::WindowSize47m => "47m",
532 Self::WindowSize48m => "48m",
533 Self::WindowSize49m => "49m",
534 Self::WindowSize50m => "50m",
535 Self::WindowSize51m => "51m",
536 Self::WindowSize52m => "52m",
537 Self::WindowSize53m => "53m",
538 Self::WindowSize54m => "54m",
539 Self::WindowSize55m => "55m",
540 Self::WindowSize56m => "56m",
541 Self::WindowSize57m => "57m",
542 Self::WindowSize58m => "58m",
543 Self::WindowSize59m => "59m",
544 Self::WindowSize1h => "1h",
545 Self::WindowSize2h => "2h",
546 Self::WindowSize3h => "3h",
547 Self::WindowSize4h => "4h",
548 Self::WindowSize5h => "5h",
549 Self::WindowSize6h => "6h",
550 Self::WindowSize7h => "7h",
551 Self::WindowSize8h => "8h",
552 Self::WindowSize9h => "9h",
553 Self::WindowSize10h => "10h",
554 Self::WindowSize11h => "11h",
555 Self::WindowSize12h => "12h",
556 Self::WindowSize13h => "13h",
557 Self::WindowSize14h => "14h",
558 Self::WindowSize15h => "15h",
559 Self::WindowSize16h => "16h",
560 Self::WindowSize17h => "17h",
561 Self::WindowSize18h => "18h",
562 Self::WindowSize19h => "19h",
563 Self::WindowSize20h => "20h",
564 Self::WindowSize21h => "21h",
565 Self::WindowSize22h => "22h",
566 Self::WindowSize23h => "23h",
567 Self::WindowSize1d => "1d",
568 Self::WindowSize2d => "2d",
569 Self::WindowSize3d => "3d",
570 Self::WindowSize4d => "4d",
571 Self::WindowSize5d => "5d",
572 Self::WindowSize6d => "6d",
573 Self::WindowSize7d => "7d",
574 }
575 }
576}
577
578impl std::str::FromStr for TickerWindowSizeEnum {
579 type Err = Box<dyn std::error::Error + Send + Sync>;
580
581 fn from_str(s: &str) -> Result<Self, Self::Err> {
582 match s {
583 "1m" => Ok(Self::WindowSize1m),
584 "2m" => Ok(Self::WindowSize2m),
585 "3m" => Ok(Self::WindowSize3m),
586 "4m" => Ok(Self::WindowSize4m),
587 "5m" => Ok(Self::WindowSize5m),
588 "6m" => Ok(Self::WindowSize6m),
589 "7m" => Ok(Self::WindowSize7m),
590 "8m" => Ok(Self::WindowSize8m),
591 "9m" => Ok(Self::WindowSize9m),
592 "10m" => Ok(Self::WindowSize10m),
593 "11m" => Ok(Self::WindowSize11m),
594 "12m" => Ok(Self::WindowSize12m),
595 "13m" => Ok(Self::WindowSize13m),
596 "14m" => Ok(Self::WindowSize14m),
597 "15m" => Ok(Self::WindowSize15m),
598 "16m" => Ok(Self::WindowSize16m),
599 "17m" => Ok(Self::WindowSize17m),
600 "18m" => Ok(Self::WindowSize18m),
601 "19m" => Ok(Self::WindowSize19m),
602 "20m" => Ok(Self::WindowSize20m),
603 "21m" => Ok(Self::WindowSize21m),
604 "22m" => Ok(Self::WindowSize22m),
605 "23m" => Ok(Self::WindowSize23m),
606 "24m" => Ok(Self::WindowSize24m),
607 "25m" => Ok(Self::WindowSize25m),
608 "26m" => Ok(Self::WindowSize26m),
609 "27m" => Ok(Self::WindowSize27m),
610 "28m" => Ok(Self::WindowSize28m),
611 "29m" => Ok(Self::WindowSize29m),
612 "30m" => Ok(Self::WindowSize30m),
613 "31m" => Ok(Self::WindowSize31m),
614 "32m" => Ok(Self::WindowSize32m),
615 "33m" => Ok(Self::WindowSize33m),
616 "34m" => Ok(Self::WindowSize34m),
617 "35m" => Ok(Self::WindowSize35m),
618 "36m" => Ok(Self::WindowSize36m),
619 "37m" => Ok(Self::WindowSize37m),
620 "38m" => Ok(Self::WindowSize38m),
621 "39m" => Ok(Self::WindowSize39m),
622 "40m" => Ok(Self::WindowSize40m),
623 "41m" => Ok(Self::WindowSize41m),
624 "42m" => Ok(Self::WindowSize42m),
625 "43m" => Ok(Self::WindowSize43m),
626 "44m" => Ok(Self::WindowSize44m),
627 "45m" => Ok(Self::WindowSize45m),
628 "46m" => Ok(Self::WindowSize46m),
629 "47m" => Ok(Self::WindowSize47m),
630 "48m" => Ok(Self::WindowSize48m),
631 "49m" => Ok(Self::WindowSize49m),
632 "50m" => Ok(Self::WindowSize50m),
633 "51m" => Ok(Self::WindowSize51m),
634 "52m" => Ok(Self::WindowSize52m),
635 "53m" => Ok(Self::WindowSize53m),
636 "54m" => Ok(Self::WindowSize54m),
637 "55m" => Ok(Self::WindowSize55m),
638 "56m" => Ok(Self::WindowSize56m),
639 "57m" => Ok(Self::WindowSize57m),
640 "58m" => Ok(Self::WindowSize58m),
641 "59m" => Ok(Self::WindowSize59m),
642 "1h" => Ok(Self::WindowSize1h),
643 "2h" => Ok(Self::WindowSize2h),
644 "3h" => Ok(Self::WindowSize3h),
645 "4h" => Ok(Self::WindowSize4h),
646 "5h" => Ok(Self::WindowSize5h),
647 "6h" => Ok(Self::WindowSize6h),
648 "7h" => Ok(Self::WindowSize7h),
649 "8h" => Ok(Self::WindowSize8h),
650 "9h" => Ok(Self::WindowSize9h),
651 "10h" => Ok(Self::WindowSize10h),
652 "11h" => Ok(Self::WindowSize11h),
653 "12h" => Ok(Self::WindowSize12h),
654 "13h" => Ok(Self::WindowSize13h),
655 "14h" => Ok(Self::WindowSize14h),
656 "15h" => Ok(Self::WindowSize15h),
657 "16h" => Ok(Self::WindowSize16h),
658 "17h" => Ok(Self::WindowSize17h),
659 "18h" => Ok(Self::WindowSize18h),
660 "19h" => Ok(Self::WindowSize19h),
661 "20h" => Ok(Self::WindowSize20h),
662 "21h" => Ok(Self::WindowSize21h),
663 "22h" => Ok(Self::WindowSize22h),
664 "23h" => Ok(Self::WindowSize23h),
665 "1d" => Ok(Self::WindowSize1d),
666 "2d" => Ok(Self::WindowSize2d),
667 "3d" => Ok(Self::WindowSize3d),
668 "4d" => Ok(Self::WindowSize4d),
669 "5d" => Ok(Self::WindowSize5d),
670 "6d" => Ok(Self::WindowSize6d),
671 "7d" => Ok(Self::WindowSize7d),
672 other => Err(format!("invalid TickerWindowSizeEnum: {}", other).into()),
673 }
674 }
675}
676
677#[allow(non_camel_case_types)]
678#[derive(Debug, Clone, Serialize, Deserialize)]
679pub enum TickerSymbolStatusEnum {
680 #[serde(rename = "TRADING")]
681 Trading,
682 #[serde(rename = "HALT")]
683 Halt,
684 #[serde(rename = "BREAK")]
685 Break,
686}
687
688impl TickerSymbolStatusEnum {
689 #[must_use]
690 pub fn as_str(&self) -> &'static str {
691 match self {
692 Self::Trading => "TRADING",
693 Self::Halt => "HALT",
694 Self::Break => "BREAK",
695 }
696 }
697}
698
699impl std::str::FromStr for TickerSymbolStatusEnum {
700 type Err = Box<dyn std::error::Error + Send + Sync>;
701
702 fn from_str(s: &str) -> Result<Self, Self::Err> {
703 match s {
704 "TRADING" => Ok(Self::Trading),
705 "HALT" => Ok(Self::Halt),
706 "BREAK" => Ok(Self::Break),
707 other => Err(format!("invalid TickerSymbolStatusEnum: {}", other).into()),
708 }
709 }
710}
711
712#[allow(non_camel_case_types)]
713#[derive(Debug, Clone, Serialize, Deserialize)]
714pub enum Ticker24hrTypeEnum {
715 #[serde(rename = "FULL")]
716 Full,
717 #[serde(rename = "MINI")]
718 Mini,
719}
720
721impl Ticker24hrTypeEnum {
722 #[must_use]
723 pub fn as_str(&self) -> &'static str {
724 match self {
725 Self::Full => "FULL",
726 Self::Mini => "MINI",
727 }
728 }
729}
730
731impl std::str::FromStr for Ticker24hrTypeEnum {
732 type Err = Box<dyn std::error::Error + Send + Sync>;
733
734 fn from_str(s: &str) -> Result<Self, Self::Err> {
735 match s {
736 "FULL" => Ok(Self::Full),
737 "MINI" => Ok(Self::Mini),
738 other => Err(format!("invalid Ticker24hrTypeEnum: {}", other).into()),
739 }
740 }
741}
742
743#[allow(non_camel_case_types)]
744#[derive(Debug, Clone, Serialize, Deserialize)]
745pub enum Ticker24hrSymbolStatusEnum {
746 #[serde(rename = "TRADING")]
747 Trading,
748 #[serde(rename = "HALT")]
749 Halt,
750 #[serde(rename = "BREAK")]
751 Break,
752}
753
754impl Ticker24hrSymbolStatusEnum {
755 #[must_use]
756 pub fn as_str(&self) -> &'static str {
757 match self {
758 Self::Trading => "TRADING",
759 Self::Halt => "HALT",
760 Self::Break => "BREAK",
761 }
762 }
763}
764
765impl std::str::FromStr for Ticker24hrSymbolStatusEnum {
766 type Err = Box<dyn std::error::Error + Send + Sync>;
767
768 fn from_str(s: &str) -> Result<Self, Self::Err> {
769 match s {
770 "TRADING" => Ok(Self::Trading),
771 "HALT" => Ok(Self::Halt),
772 "BREAK" => Ok(Self::Break),
773 other => Err(format!("invalid Ticker24hrSymbolStatusEnum: {}", other).into()),
774 }
775 }
776}
777
778#[allow(non_camel_case_types)]
779#[derive(Debug, Clone, Serialize, Deserialize)]
780pub enum TickerBookSymbolStatusEnum {
781 #[serde(rename = "TRADING")]
782 Trading,
783 #[serde(rename = "HALT")]
784 Halt,
785 #[serde(rename = "BREAK")]
786 Break,
787}
788
789impl TickerBookSymbolStatusEnum {
790 #[must_use]
791 pub fn as_str(&self) -> &'static str {
792 match self {
793 Self::Trading => "TRADING",
794 Self::Halt => "HALT",
795 Self::Break => "BREAK",
796 }
797 }
798}
799
800impl std::str::FromStr for TickerBookSymbolStatusEnum {
801 type Err = Box<dyn std::error::Error + Send + Sync>;
802
803 fn from_str(s: &str) -> Result<Self, Self::Err> {
804 match s {
805 "TRADING" => Ok(Self::Trading),
806 "HALT" => Ok(Self::Halt),
807 "BREAK" => Ok(Self::Break),
808 other => Err(format!("invalid TickerBookSymbolStatusEnum: {}", other).into()),
809 }
810 }
811}
812
813#[allow(non_camel_case_types)]
814#[derive(Debug, Clone, Serialize, Deserialize)]
815pub enum TickerPriceSymbolStatusEnum {
816 #[serde(rename = "TRADING")]
817 Trading,
818 #[serde(rename = "HALT")]
819 Halt,
820 #[serde(rename = "BREAK")]
821 Break,
822}
823
824impl TickerPriceSymbolStatusEnum {
825 #[must_use]
826 pub fn as_str(&self) -> &'static str {
827 match self {
828 Self::Trading => "TRADING",
829 Self::Halt => "HALT",
830 Self::Break => "BREAK",
831 }
832 }
833}
834
835impl std::str::FromStr for TickerPriceSymbolStatusEnum {
836 type Err = Box<dyn std::error::Error + Send + Sync>;
837
838 fn from_str(s: &str) -> Result<Self, Self::Err> {
839 match s {
840 "TRADING" => Ok(Self::Trading),
841 "HALT" => Ok(Self::Halt),
842 "BREAK" => Ok(Self::Break),
843 other => Err(format!("invalid TickerPriceSymbolStatusEnum: {}", other).into()),
844 }
845 }
846}
847
848#[allow(non_camel_case_types)]
849#[derive(Debug, Clone, Serialize, Deserialize)]
850pub enum TickerTradingDayTypeEnum {
851 #[serde(rename = "FULL")]
852 Full,
853 #[serde(rename = "MINI")]
854 Mini,
855}
856
857impl TickerTradingDayTypeEnum {
858 #[must_use]
859 pub fn as_str(&self) -> &'static str {
860 match self {
861 Self::Full => "FULL",
862 Self::Mini => "MINI",
863 }
864 }
865}
866
867impl std::str::FromStr for TickerTradingDayTypeEnum {
868 type Err = Box<dyn std::error::Error + Send + Sync>;
869
870 fn from_str(s: &str) -> Result<Self, Self::Err> {
871 match s {
872 "FULL" => Ok(Self::Full),
873 "MINI" => Ok(Self::Mini),
874 other => Err(format!("invalid TickerTradingDayTypeEnum: {}", other).into()),
875 }
876 }
877}
878
879#[allow(non_camel_case_types)]
880#[derive(Debug, Clone, Serialize, Deserialize)]
881pub enum TickerTradingDaySymbolStatusEnum {
882 #[serde(rename = "TRADING")]
883 Trading,
884 #[serde(rename = "HALT")]
885 Halt,
886 #[serde(rename = "BREAK")]
887 Break,
888}
889
890impl TickerTradingDaySymbolStatusEnum {
891 #[must_use]
892 pub fn as_str(&self) -> &'static str {
893 match self {
894 Self::Trading => "TRADING",
895 Self::Halt => "HALT",
896 Self::Break => "BREAK",
897 }
898 }
899}
900
901impl std::str::FromStr for TickerTradingDaySymbolStatusEnum {
902 type Err = Box<dyn std::error::Error + Send + Sync>;
903
904 fn from_str(s: &str) -> Result<Self, Self::Err> {
905 match s {
906 "TRADING" => Ok(Self::Trading),
907 "HALT" => Ok(Self::Halt),
908 "BREAK" => Ok(Self::Break),
909 other => Err(format!("invalid TickerTradingDaySymbolStatusEnum: {}", other).into()),
910 }
911 }
912}
913
914#[allow(non_camel_case_types)]
915#[derive(Debug, Clone, Serialize, Deserialize)]
916pub enum UiKlinesIntervalEnum {
917 #[serde(rename = "1s")]
918 Interval1s,
919 #[serde(rename = "1m")]
920 Interval1m,
921 #[serde(rename = "3m")]
922 Interval3m,
923 #[serde(rename = "5m")]
924 Interval5m,
925 #[serde(rename = "15m")]
926 Interval15m,
927 #[serde(rename = "30m")]
928 Interval30m,
929 #[serde(rename = "1h")]
930 Interval1h,
931 #[serde(rename = "2h")]
932 Interval2h,
933 #[serde(rename = "4h")]
934 Interval4h,
935 #[serde(rename = "6h")]
936 Interval6h,
937 #[serde(rename = "8h")]
938 Interval8h,
939 #[serde(rename = "12h")]
940 Interval12h,
941 #[serde(rename = "1d")]
942 Interval1d,
943 #[serde(rename = "3d")]
944 Interval3d,
945 #[serde(rename = "1w")]
946 Interval1w,
947 #[serde(rename = "1M")]
948 Interval1M,
949}
950
951impl UiKlinesIntervalEnum {
952 #[must_use]
953 pub fn as_str(&self) -> &'static str {
954 match self {
955 Self::Interval1s => "1s",
956 Self::Interval1m => "1m",
957 Self::Interval3m => "3m",
958 Self::Interval5m => "5m",
959 Self::Interval15m => "15m",
960 Self::Interval30m => "30m",
961 Self::Interval1h => "1h",
962 Self::Interval2h => "2h",
963 Self::Interval4h => "4h",
964 Self::Interval6h => "6h",
965 Self::Interval8h => "8h",
966 Self::Interval12h => "12h",
967 Self::Interval1d => "1d",
968 Self::Interval3d => "3d",
969 Self::Interval1w => "1w",
970 Self::Interval1M => "1M",
971 }
972 }
973}
974
975impl std::str::FromStr for UiKlinesIntervalEnum {
976 type Err = Box<dyn std::error::Error + Send + Sync>;
977
978 fn from_str(s: &str) -> Result<Self, Self::Err> {
979 match s {
980 "1s" => Ok(Self::Interval1s),
981 "1m" => Ok(Self::Interval1m),
982 "3m" => Ok(Self::Interval3m),
983 "5m" => Ok(Self::Interval5m),
984 "15m" => Ok(Self::Interval15m),
985 "30m" => Ok(Self::Interval30m),
986 "1h" => Ok(Self::Interval1h),
987 "2h" => Ok(Self::Interval2h),
988 "4h" => Ok(Self::Interval4h),
989 "6h" => Ok(Self::Interval6h),
990 "8h" => Ok(Self::Interval8h),
991 "12h" => Ok(Self::Interval12h),
992 "1d" => Ok(Self::Interval1d),
993 "3d" => Ok(Self::Interval3d),
994 "1w" => Ok(Self::Interval1w),
995 "1M" => Ok(Self::Interval1M),
996 other => Err(format!("invalid UiKlinesIntervalEnum: {}", other).into()),
997 }
998 }
999}
1000
1001#[derive(Clone, Debug, Builder, Deserialize)]
1006#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1007pub struct AvgPriceParams {
1008 #[builder(setter(into))]
1013 #[serde(rename = "symbol")]
1014 pub symbol: String,
1015 #[builder(setter(into), default)]
1019 #[serde(rename = "id", default)]
1020 pub id: Option<String>,
1021}
1022
1023impl AvgPriceParams {
1024 #[must_use]
1031 pub fn builder(symbol: String) -> AvgPriceParamsBuilder {
1032 AvgPriceParamsBuilder::default().symbol(symbol)
1033 }
1034}
1035#[derive(Clone, Debug, Builder, Deserialize)]
1040#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1041pub struct BlockTradesHistoricalParams {
1042 #[builder(setter(into))]
1047 #[serde(rename = "symbol")]
1048 pub symbol: String,
1049 #[builder(setter(into))]
1053 #[serde(rename = "fromId")]
1054 pub from_id: i64,
1055 #[builder(setter(into), default)]
1059 #[serde(rename = "id", default)]
1060 pub id: Option<String>,
1061 #[builder(setter(into), default)]
1065 #[serde(rename = "limit", default)]
1066 pub limit: Option<i64>,
1067}
1068
1069impl BlockTradesHistoricalParams {
1070 #[must_use]
1078 pub fn builder(symbol: String, from_id: i64) -> BlockTradesHistoricalParamsBuilder {
1079 BlockTradesHistoricalParamsBuilder::default()
1080 .symbol(symbol)
1081 .from_id(from_id)
1082 }
1083}
1084#[derive(Clone, Debug, Builder, Deserialize)]
1089#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1090pub struct DepthParams {
1091 #[builder(setter(into))]
1096 #[serde(rename = "symbol")]
1097 pub symbol: String,
1098 #[builder(setter(into), default)]
1102 #[serde(rename = "id", default)]
1103 pub id: Option<String>,
1104 #[builder(setter(into), default)]
1109 #[serde(rename = "limit", default)]
1110 pub limit: Option<i32>,
1111 #[builder(setter(into), default)]
1116 #[serde(rename = "symbolStatus", default)]
1117 pub symbol_status: Option<DepthSymbolStatusEnum>,
1118}
1119
1120impl DepthParams {
1121 #[must_use]
1128 pub fn builder(symbol: String) -> DepthParamsBuilder {
1129 DepthParamsBuilder::default().symbol(symbol)
1130 }
1131}
1132#[derive(Clone, Debug, Builder, Deserialize)]
1137#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1138pub struct KlinesParams {
1139 #[builder(setter(into))]
1144 #[serde(rename = "symbol")]
1145 pub symbol: String,
1146 #[builder(setter(into))]
1151 #[serde(rename = "interval")]
1152 pub interval: KlinesIntervalEnum,
1153 #[builder(setter(into), default)]
1157 #[serde(rename = "id", default)]
1158 pub id: Option<String>,
1159 #[builder(setter(into), default)]
1164 #[serde(rename = "startTime", default)]
1165 pub start_time: Option<i64>,
1166 #[builder(setter(into), default)]
1171 #[serde(rename = "endTime", default)]
1172 pub end_time: Option<i64>,
1173 #[builder(setter(into), default)]
1177 #[serde(rename = "timeZone", default)]
1178 pub time_zone: Option<String>,
1179 #[builder(setter(into), default)]
1184 #[serde(rename = "limit", default)]
1185 pub limit: Option<i32>,
1186}
1187
1188impl KlinesParams {
1189 #[must_use]
1197 pub fn builder(symbol: String, interval: KlinesIntervalEnum) -> KlinesParamsBuilder {
1198 KlinesParamsBuilder::default()
1199 .symbol(symbol)
1200 .interval(interval)
1201 }
1202}
1203#[derive(Clone, Debug, Builder, Deserialize)]
1208#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1209pub struct ReferencePriceParams {
1210 #[builder(setter(into))]
1215 #[serde(rename = "symbol")]
1216 pub symbol: String,
1217 #[builder(setter(into), default)]
1221 #[serde(rename = "id", default)]
1222 pub id: Option<String>,
1223}
1224
1225impl ReferencePriceParams {
1226 #[must_use]
1233 pub fn builder(symbol: String) -> ReferencePriceParamsBuilder {
1234 ReferencePriceParamsBuilder::default().symbol(symbol)
1235 }
1236}
1237#[derive(Clone, Debug, Builder, Deserialize)]
1242#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1243pub struct ReferencePriceCalculationParams {
1244 #[builder(setter(into))]
1249 #[serde(rename = "symbol")]
1250 pub symbol: String,
1251 #[builder(setter(into), default)]
1255 #[serde(rename = "id", default)]
1256 pub id: Option<String>,
1257 #[builder(setter(into), default)]
1262 #[serde(rename = "symbolStatus", default)]
1263 pub symbol_status: Option<ReferencePriceCalculationSymbolStatusEnum>,
1264}
1265
1266impl ReferencePriceCalculationParams {
1267 #[must_use]
1274 pub fn builder(symbol: String) -> ReferencePriceCalculationParamsBuilder {
1275 ReferencePriceCalculationParamsBuilder::default().symbol(symbol)
1276 }
1277}
1278#[derive(Clone, Debug, Builder, Deserialize, Default)]
1283#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1284pub struct TickerParams {
1285 #[builder(setter(into), default)]
1289 #[serde(rename = "id", default)]
1290 pub id: Option<String>,
1291 #[builder(setter(into), default)]
1295 #[serde(rename = "symbol", default)]
1296 pub symbol: Option<String>,
1297 #[builder(setter(into), default)]
1301 #[serde(rename = "symbols", default)]
1302 pub symbols: Option<Vec<String>>,
1303 #[builder(setter(into), default)]
1307 #[serde(rename = "type", default)]
1308 pub r#type: Option<TickerTypeEnum>,
1309 #[builder(setter(into), default)]
1313 #[serde(rename = "windowSize", default)]
1314 pub window_size: Option<TickerWindowSizeEnum>,
1315 #[builder(setter(into), default)]
1319 #[serde(rename = "symbolStatus", default)]
1320 pub symbol_status: Option<TickerSymbolStatusEnum>,
1321}
1322
1323impl TickerParams {
1324 #[must_use]
1327 pub fn builder() -> TickerParamsBuilder {
1328 TickerParamsBuilder::default()
1329 }
1330}
1331#[derive(Clone, Debug, Builder, Deserialize, Default)]
1336#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1337pub struct Ticker24hrParams {
1338 #[builder(setter(into), default)]
1342 #[serde(rename = "id", default)]
1343 pub id: Option<String>,
1344 #[builder(setter(into), default)]
1349 #[serde(rename = "symbol", default)]
1350 pub symbol: Option<String>,
1351 #[builder(setter(into), default)]
1356 #[serde(rename = "symbols", default)]
1357 pub symbols: Option<Vec<String>>,
1358 #[builder(setter(into), default)]
1362 #[serde(rename = "type", default)]
1363 pub r#type: Option<Ticker24hrTypeEnum>,
1364 #[builder(setter(into), default)]
1368 #[serde(rename = "symbolStatus", default)]
1369 pub symbol_status: Option<Ticker24hrSymbolStatusEnum>,
1370}
1371
1372impl Ticker24hrParams {
1373 #[must_use]
1376 pub fn builder() -> Ticker24hrParamsBuilder {
1377 Ticker24hrParamsBuilder::default()
1378 }
1379}
1380#[derive(Clone, Debug, Builder, Deserialize, Default)]
1385#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1386pub struct TickerBookParams {
1387 #[builder(setter(into), default)]
1391 #[serde(rename = "id", default)]
1392 pub id: Option<String>,
1393 #[builder(setter(into), default)]
1397 #[serde(rename = "symbol", default)]
1398 pub symbol: Option<String>,
1399 #[builder(setter(into), default)]
1403 #[serde(rename = "symbols", default)]
1404 pub symbols: Option<Vec<String>>,
1405 #[builder(setter(into), default)]
1409 #[serde(rename = "symbolStatus", default)]
1410 pub symbol_status: Option<TickerBookSymbolStatusEnum>,
1411}
1412
1413impl TickerBookParams {
1414 #[must_use]
1417 pub fn builder() -> TickerBookParamsBuilder {
1418 TickerBookParamsBuilder::default()
1419 }
1420}
1421#[derive(Clone, Debug, Builder, Deserialize, Default)]
1426#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1427pub struct TickerPriceParams {
1428 #[builder(setter(into), default)]
1432 #[serde(rename = "id", default)]
1433 pub id: Option<String>,
1434 #[builder(setter(into), default)]
1438 #[serde(rename = "symbol", default)]
1439 pub symbol: Option<String>,
1440 #[builder(setter(into), default)]
1444 #[serde(rename = "symbols", default)]
1445 pub symbols: Option<Vec<String>>,
1446 #[builder(setter(into), default)]
1450 #[serde(rename = "symbolStatus", default)]
1451 pub symbol_status: Option<TickerPriceSymbolStatusEnum>,
1452}
1453
1454impl TickerPriceParams {
1455 #[must_use]
1458 pub fn builder() -> TickerPriceParamsBuilder {
1459 TickerPriceParamsBuilder::default()
1460 }
1461}
1462#[derive(Clone, Debug, Builder, Deserialize, Default)]
1467#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1468pub struct TickerTradingDayParams {
1469 #[builder(setter(into), default)]
1473 #[serde(rename = "id", default)]
1474 pub id: Option<String>,
1475 #[builder(setter(into), default)]
1480 #[serde(rename = "symbol", default)]
1481 pub symbol: Option<String>,
1482 #[builder(setter(into), default)]
1487 #[serde(rename = "symbols", default)]
1488 pub symbols: Option<Vec<String>>,
1489 #[builder(setter(into), default)]
1493 #[serde(rename = "timeZone", default)]
1494 pub time_zone: Option<String>,
1495 #[builder(setter(into), default)]
1499 #[serde(rename = "type", default)]
1500 pub r#type: Option<TickerTradingDayTypeEnum>,
1501 #[builder(setter(into), default)]
1505 #[serde(rename = "symbolStatus", default)]
1506 pub symbol_status: Option<TickerTradingDaySymbolStatusEnum>,
1507}
1508
1509impl TickerTradingDayParams {
1510 #[must_use]
1513 pub fn builder() -> TickerTradingDayParamsBuilder {
1514 TickerTradingDayParamsBuilder::default()
1515 }
1516}
1517#[derive(Clone, Debug, Builder, Deserialize)]
1522#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1523pub struct TradesAggregateParams {
1524 #[builder(setter(into))]
1529 #[serde(rename = "symbol")]
1530 pub symbol: String,
1531 #[builder(setter(into), default)]
1535 #[serde(rename = "id", default)]
1536 pub id: Option<String>,
1537 #[builder(setter(into), default)]
1541 #[serde(rename = "fromId", default)]
1542 pub from_id: Option<i64>,
1543 #[builder(setter(into), default)]
1547 #[serde(rename = "startTime", default)]
1548 pub start_time: Option<i64>,
1549 #[builder(setter(into), default)]
1553 #[serde(rename = "endTime", default)]
1554 pub end_time: Option<i64>,
1555 #[builder(setter(into), default)]
1560 #[serde(rename = "limit", default)]
1561 pub limit: Option<i32>,
1562}
1563
1564impl TradesAggregateParams {
1565 #[must_use]
1572 pub fn builder(symbol: String) -> TradesAggregateParamsBuilder {
1573 TradesAggregateParamsBuilder::default().symbol(symbol)
1574 }
1575}
1576#[derive(Clone, Debug, Builder, Deserialize)]
1581#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1582pub struct TradesHistoricalParams {
1583 #[builder(setter(into))]
1588 #[serde(rename = "symbol")]
1589 pub symbol: String,
1590 #[builder(setter(into), default)]
1594 #[serde(rename = "id", default)]
1595 pub id: Option<String>,
1596 #[builder(setter(into), default)]
1600 #[serde(rename = "fromId", default)]
1601 pub from_id: Option<i64>,
1602 #[builder(setter(into), default)]
1607 #[serde(rename = "limit", default)]
1608 pub limit: Option<i32>,
1609}
1610
1611impl TradesHistoricalParams {
1612 #[must_use]
1619 pub fn builder(symbol: String) -> TradesHistoricalParamsBuilder {
1620 TradesHistoricalParamsBuilder::default().symbol(symbol)
1621 }
1622}
1623#[derive(Clone, Debug, Builder, Deserialize)]
1628#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1629pub struct TradesRecentParams {
1630 #[builder(setter(into))]
1635 #[serde(rename = "symbol")]
1636 pub symbol: String,
1637 #[builder(setter(into), default)]
1641 #[serde(rename = "id", default)]
1642 pub id: Option<String>,
1643 #[builder(setter(into), default)]
1648 #[serde(rename = "limit", default)]
1649 pub limit: Option<i32>,
1650}
1651
1652impl TradesRecentParams {
1653 #[must_use]
1660 pub fn builder(symbol: String) -> TradesRecentParamsBuilder {
1661 TradesRecentParamsBuilder::default().symbol(symbol)
1662 }
1663}
1664#[derive(Clone, Debug, Builder, Deserialize)]
1669#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
1670pub struct UiKlinesParams {
1671 #[builder(setter(into))]
1676 #[serde(rename = "symbol")]
1677 pub symbol: String,
1678 #[builder(setter(into))]
1683 #[serde(rename = "interval")]
1684 pub interval: UiKlinesIntervalEnum,
1685 #[builder(setter(into), default)]
1689 #[serde(rename = "id", default)]
1690 pub id: Option<String>,
1691 #[builder(setter(into), default)]
1696 #[serde(rename = "startTime", default)]
1697 pub start_time: Option<i64>,
1698 #[builder(setter(into), default)]
1703 #[serde(rename = "endTime", default)]
1704 pub end_time: Option<i64>,
1705 #[builder(setter(into), default)]
1709 #[serde(rename = "timeZone", default)]
1710 pub time_zone: Option<String>,
1711 #[builder(setter(into), default)]
1716 #[serde(rename = "limit", default)]
1717 pub limit: Option<i32>,
1718}
1719
1720impl UiKlinesParams {
1721 #[must_use]
1729 pub fn builder(symbol: String, interval: UiKlinesIntervalEnum) -> UiKlinesParamsBuilder {
1730 UiKlinesParamsBuilder::default()
1731 .symbol(symbol)
1732 .interval(interval)
1733 }
1734}
1735
1736#[async_trait]
1737impl MarketApi for MarketApiClient {
1738 async fn avg_price(
1739 &self,
1740 params: AvgPriceParams,
1741 ) -> anyhow::Result<WebsocketApiResponse<Box<models::AvgPriceResponseResult>>> {
1742 let AvgPriceParams { symbol, id } = params;
1743
1744 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1745 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1746 if let Some(value) = id {
1747 payload.insert("id".to_string(), serde_json::json!(value));
1748 }
1749 let payload = remove_empty_value(payload);
1750
1751 self.websocket_api_base
1752 .send_message::<Box<models::AvgPriceResponseResult>>(
1753 "/avgPrice".trim_start_matches('/'),
1754 payload,
1755 WebsocketMessageSendOptions::new(),
1756 )
1757 .await
1758 .map_err(anyhow::Error::from)?
1759 .into_iter()
1760 .next()
1761 .ok_or(WebsocketError::NoResponse)
1762 .map_err(anyhow::Error::from)
1763 }
1764
1765 async fn block_trades_historical(
1766 &self,
1767 params: BlockTradesHistoricalParams,
1768 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::BlockTradesHistoricalResponseResultInner>>>
1769 {
1770 let BlockTradesHistoricalParams {
1771 symbol,
1772 from_id,
1773 id,
1774 limit,
1775 } = params;
1776
1777 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1778 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1779 payload.insert("fromId".to_string(), serde_json::json!(from_id));
1780 if let Some(value) = id {
1781 payload.insert("id".to_string(), serde_json::json!(value));
1782 }
1783 if let Some(value) = limit {
1784 payload.insert("limit".to_string(), serde_json::json!(value));
1785 }
1786 let payload = remove_empty_value(payload);
1787
1788 self.websocket_api_base
1789 .send_message::<Vec<models::BlockTradesHistoricalResponseResultInner>>(
1790 "/blockTrades.historical".trim_start_matches('/'),
1791 payload,
1792 WebsocketMessageSendOptions::new(),
1793 )
1794 .await
1795 .map_err(anyhow::Error::from)?
1796 .into_iter()
1797 .next()
1798 .ok_or(WebsocketError::NoResponse)
1799 .map_err(anyhow::Error::from)
1800 }
1801
1802 async fn depth(
1803 &self,
1804 params: DepthParams,
1805 ) -> anyhow::Result<WebsocketApiResponse<Box<models::DepthResponseResult>>> {
1806 let DepthParams {
1807 symbol,
1808 id,
1809 limit,
1810 symbol_status,
1811 } = params;
1812
1813 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1814 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1815 if let Some(value) = id {
1816 payload.insert("id".to_string(), serde_json::json!(value));
1817 }
1818 if let Some(value) = limit {
1819 payload.insert("limit".to_string(), serde_json::json!(value));
1820 }
1821 if let Some(value) = symbol_status {
1822 payload.insert("symbolStatus".to_string(), serde_json::json!(value));
1823 }
1824 let payload = remove_empty_value(payload);
1825
1826 self.websocket_api_base
1827 .send_message::<Box<models::DepthResponseResult>>(
1828 "/depth".trim_start_matches('/'),
1829 payload,
1830 WebsocketMessageSendOptions::new(),
1831 )
1832 .await
1833 .map_err(anyhow::Error::from)?
1834 .into_iter()
1835 .next()
1836 .ok_or(WebsocketError::NoResponse)
1837 .map_err(anyhow::Error::from)
1838 }
1839
1840 async fn klines(
1841 &self,
1842 params: KlinesParams,
1843 ) -> anyhow::Result<WebsocketApiResponse<Vec<Vec<models::KlinesResponseResultInnerInner>>>>
1844 {
1845 let KlinesParams {
1846 symbol,
1847 interval,
1848 id,
1849 start_time,
1850 end_time,
1851 time_zone,
1852 limit,
1853 } = params;
1854
1855 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1856 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1857 payload.insert("interval".to_string(), serde_json::json!(interval));
1858 if let Some(value) = id {
1859 payload.insert("id".to_string(), serde_json::json!(value));
1860 }
1861 if let Some(value) = start_time {
1862 payload.insert("startTime".to_string(), serde_json::json!(value));
1863 }
1864 if let Some(value) = end_time {
1865 payload.insert("endTime".to_string(), serde_json::json!(value));
1866 }
1867 if let Some(value) = time_zone {
1868 payload.insert("timeZone".to_string(), serde_json::json!(value));
1869 }
1870 if let Some(value) = limit {
1871 payload.insert("limit".to_string(), serde_json::json!(value));
1872 }
1873 let payload = remove_empty_value(payload);
1874
1875 self.websocket_api_base
1876 .send_message::<Vec<Vec<models::KlinesResponseResultInnerInner>>>(
1877 "/klines".trim_start_matches('/'),
1878 payload,
1879 WebsocketMessageSendOptions::new(),
1880 )
1881 .await
1882 .map_err(anyhow::Error::from)?
1883 .into_iter()
1884 .next()
1885 .ok_or(WebsocketError::NoResponse)
1886 .map_err(anyhow::Error::from)
1887 }
1888
1889 async fn reference_price(
1890 &self,
1891 params: ReferencePriceParams,
1892 ) -> anyhow::Result<WebsocketApiResponse<Box<models::ReferencePriceResponseResult>>> {
1893 let ReferencePriceParams { symbol, id } = params;
1894
1895 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1896 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1897 if let Some(value) = id {
1898 payload.insert("id".to_string(), serde_json::json!(value));
1899 }
1900 let payload = remove_empty_value(payload);
1901
1902 self.websocket_api_base
1903 .send_message::<Box<models::ReferencePriceResponseResult>>(
1904 "/referencePrice".trim_start_matches('/'),
1905 payload,
1906 WebsocketMessageSendOptions::new(),
1907 )
1908 .await
1909 .map_err(anyhow::Error::from)?
1910 .into_iter()
1911 .next()
1912 .ok_or(WebsocketError::NoResponse)
1913 .map_err(anyhow::Error::from)
1914 }
1915
1916 async fn reference_price_calculation(
1917 &self,
1918 params: ReferencePriceCalculationParams,
1919 ) -> anyhow::Result<WebsocketApiResponse<Box<models::ReferencePriceCalculationResponseResult>>>
1920 {
1921 let ReferencePriceCalculationParams {
1922 symbol,
1923 id,
1924 symbol_status,
1925 } = params;
1926
1927 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1928 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1929 if let Some(value) = id {
1930 payload.insert("id".to_string(), serde_json::json!(value));
1931 }
1932 if let Some(value) = symbol_status {
1933 payload.insert("symbolStatus".to_string(), serde_json::json!(value));
1934 }
1935 let payload = remove_empty_value(payload);
1936
1937 self.websocket_api_base
1938 .send_message::<Box<models::ReferencePriceCalculationResponseResult>>(
1939 "/referencePrice.calculation".trim_start_matches('/'),
1940 payload,
1941 WebsocketMessageSendOptions::new(),
1942 )
1943 .await
1944 .map_err(anyhow::Error::from)?
1945 .into_iter()
1946 .next()
1947 .ok_or(WebsocketError::NoResponse)
1948 .map_err(anyhow::Error::from)
1949 }
1950
1951 async fn ticker(
1952 &self,
1953 params: TickerParams,
1954 ) -> anyhow::Result<WebsocketApiResponse<models::TickerResponse>> {
1955 let TickerParams {
1956 id,
1957 symbol,
1958 symbols,
1959 r#type,
1960 window_size,
1961 symbol_status,
1962 } = params;
1963
1964 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1965 if let Some(value) = id {
1966 payload.insert("id".to_string(), serde_json::json!(value));
1967 }
1968 if let Some(value) = symbol {
1969 payload.insert("symbol".to_string(), serde_json::json!(value));
1970 }
1971 if let Some(value) = symbols {
1972 payload.insert("symbols".to_string(), serde_json::json!(value));
1973 }
1974 if let Some(value) = r#type {
1975 payload.insert("type".to_string(), serde_json::json!(value));
1976 }
1977 if let Some(value) = window_size {
1978 payload.insert("windowSize".to_string(), serde_json::json!(value));
1979 }
1980 if let Some(value) = symbol_status {
1981 payload.insert("symbolStatus".to_string(), serde_json::json!(value));
1982 }
1983 let payload = remove_empty_value(payload);
1984
1985 self.websocket_api_base
1986 .send_message::<models::TickerResponse>(
1987 "/ticker".trim_start_matches('/'),
1988 payload,
1989 WebsocketMessageSendOptions::new(),
1990 )
1991 .await
1992 .map_err(anyhow::Error::from)?
1993 .into_iter()
1994 .next()
1995 .ok_or(WebsocketError::NoResponse)
1996 .map_err(anyhow::Error::from)
1997 }
1998
1999 async fn ticker24hr(
2000 &self,
2001 params: Ticker24hrParams,
2002 ) -> anyhow::Result<WebsocketApiResponse<models::Ticker24hrResponse>> {
2003 let Ticker24hrParams {
2004 id,
2005 symbol,
2006 symbols,
2007 r#type,
2008 symbol_status,
2009 } = params;
2010
2011 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
2012 if let Some(value) = id {
2013 payload.insert("id".to_string(), serde_json::json!(value));
2014 }
2015 if let Some(value) = symbol {
2016 payload.insert("symbol".to_string(), serde_json::json!(value));
2017 }
2018 if let Some(value) = symbols {
2019 payload.insert("symbols".to_string(), serde_json::json!(value));
2020 }
2021 if let Some(value) = r#type {
2022 payload.insert("type".to_string(), serde_json::json!(value));
2023 }
2024 if let Some(value) = symbol_status {
2025 payload.insert("symbolStatus".to_string(), serde_json::json!(value));
2026 }
2027 let payload = remove_empty_value(payload);
2028
2029 self.websocket_api_base
2030 .send_message::<models::Ticker24hrResponse>(
2031 "/ticker.24hr".trim_start_matches('/'),
2032 payload,
2033 WebsocketMessageSendOptions::new(),
2034 )
2035 .await
2036 .map_err(anyhow::Error::from)?
2037 .into_iter()
2038 .next()
2039 .ok_or(WebsocketError::NoResponse)
2040 .map_err(anyhow::Error::from)
2041 }
2042
2043 async fn ticker_book(
2044 &self,
2045 params: TickerBookParams,
2046 ) -> anyhow::Result<WebsocketApiResponse<models::TickerBookResponse>> {
2047 let TickerBookParams {
2048 id,
2049 symbol,
2050 symbols,
2051 symbol_status,
2052 } = params;
2053
2054 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
2055 if let Some(value) = id {
2056 payload.insert("id".to_string(), serde_json::json!(value));
2057 }
2058 if let Some(value) = symbol {
2059 payload.insert("symbol".to_string(), serde_json::json!(value));
2060 }
2061 if let Some(value) = symbols {
2062 payload.insert("symbols".to_string(), serde_json::json!(value));
2063 }
2064 if let Some(value) = symbol_status {
2065 payload.insert("symbolStatus".to_string(), serde_json::json!(value));
2066 }
2067 let payload = remove_empty_value(payload);
2068
2069 self.websocket_api_base
2070 .send_message::<models::TickerBookResponse>(
2071 "/ticker.book".trim_start_matches('/'),
2072 payload,
2073 WebsocketMessageSendOptions::new(),
2074 )
2075 .await
2076 .map_err(anyhow::Error::from)?
2077 .into_iter()
2078 .next()
2079 .ok_or(WebsocketError::NoResponse)
2080 .map_err(anyhow::Error::from)
2081 }
2082
2083 async fn ticker_price(
2084 &self,
2085 params: TickerPriceParams,
2086 ) -> anyhow::Result<WebsocketApiResponse<models::TickerPriceResponse>> {
2087 let TickerPriceParams {
2088 id,
2089 symbol,
2090 symbols,
2091 symbol_status,
2092 } = params;
2093
2094 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
2095 if let Some(value) = id {
2096 payload.insert("id".to_string(), serde_json::json!(value));
2097 }
2098 if let Some(value) = symbol {
2099 payload.insert("symbol".to_string(), serde_json::json!(value));
2100 }
2101 if let Some(value) = symbols {
2102 payload.insert("symbols".to_string(), serde_json::json!(value));
2103 }
2104 if let Some(value) = symbol_status {
2105 payload.insert("symbolStatus".to_string(), serde_json::json!(value));
2106 }
2107 let payload = remove_empty_value(payload);
2108
2109 self.websocket_api_base
2110 .send_message::<models::TickerPriceResponse>(
2111 "/ticker.price".trim_start_matches('/'),
2112 payload,
2113 WebsocketMessageSendOptions::new(),
2114 )
2115 .await
2116 .map_err(anyhow::Error::from)?
2117 .into_iter()
2118 .next()
2119 .ok_or(WebsocketError::NoResponse)
2120 .map_err(anyhow::Error::from)
2121 }
2122
2123 async fn ticker_trading_day(
2124 &self,
2125 params: TickerTradingDayParams,
2126 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::TickerTradingDayResponseResultInner>>>
2127 {
2128 let TickerTradingDayParams {
2129 id,
2130 symbol,
2131 symbols,
2132 time_zone,
2133 r#type,
2134 symbol_status,
2135 } = params;
2136
2137 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
2138 if let Some(value) = id {
2139 payload.insert("id".to_string(), serde_json::json!(value));
2140 }
2141 if let Some(value) = symbol {
2142 payload.insert("symbol".to_string(), serde_json::json!(value));
2143 }
2144 if let Some(value) = symbols {
2145 payload.insert("symbols".to_string(), serde_json::json!(value));
2146 }
2147 if let Some(value) = time_zone {
2148 payload.insert("timeZone".to_string(), serde_json::json!(value));
2149 }
2150 if let Some(value) = r#type {
2151 payload.insert("type".to_string(), serde_json::json!(value));
2152 }
2153 if let Some(value) = symbol_status {
2154 payload.insert("symbolStatus".to_string(), serde_json::json!(value));
2155 }
2156 let payload = remove_empty_value(payload);
2157
2158 self.websocket_api_base
2159 .send_message::<Vec<models::TickerTradingDayResponseResultInner>>(
2160 "/ticker.tradingDay".trim_start_matches('/'),
2161 payload,
2162 WebsocketMessageSendOptions::new(),
2163 )
2164 .await
2165 .map_err(anyhow::Error::from)?
2166 .into_iter()
2167 .next()
2168 .ok_or(WebsocketError::NoResponse)
2169 .map_err(anyhow::Error::from)
2170 }
2171
2172 async fn trades_aggregate(
2173 &self,
2174 params: TradesAggregateParams,
2175 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::TradesAggregateResponseResultInner>>> {
2176 let TradesAggregateParams {
2177 symbol,
2178 id,
2179 from_id,
2180 start_time,
2181 end_time,
2182 limit,
2183 } = params;
2184
2185 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
2186 payload.insert("symbol".to_string(), serde_json::json!(symbol));
2187 if let Some(value) = id {
2188 payload.insert("id".to_string(), serde_json::json!(value));
2189 }
2190 if let Some(value) = from_id {
2191 payload.insert("fromId".to_string(), serde_json::json!(value));
2192 }
2193 if let Some(value) = start_time {
2194 payload.insert("startTime".to_string(), serde_json::json!(value));
2195 }
2196 if let Some(value) = end_time {
2197 payload.insert("endTime".to_string(), serde_json::json!(value));
2198 }
2199 if let Some(value) = limit {
2200 payload.insert("limit".to_string(), serde_json::json!(value));
2201 }
2202 let payload = remove_empty_value(payload);
2203
2204 self.websocket_api_base
2205 .send_message::<Vec<models::TradesAggregateResponseResultInner>>(
2206 "/trades.aggregate".trim_start_matches('/'),
2207 payload,
2208 WebsocketMessageSendOptions::new(),
2209 )
2210 .await
2211 .map_err(anyhow::Error::from)?
2212 .into_iter()
2213 .next()
2214 .ok_or(WebsocketError::NoResponse)
2215 .map_err(anyhow::Error::from)
2216 }
2217
2218 async fn trades_historical(
2219 &self,
2220 params: TradesHistoricalParams,
2221 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::TradesHistoricalResponseResultInner>>>
2222 {
2223 let TradesHistoricalParams {
2224 symbol,
2225 id,
2226 from_id,
2227 limit,
2228 } = params;
2229
2230 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
2231 payload.insert("symbol".to_string(), serde_json::json!(symbol));
2232 if let Some(value) = id {
2233 payload.insert("id".to_string(), serde_json::json!(value));
2234 }
2235 if let Some(value) = from_id {
2236 payload.insert("fromId".to_string(), serde_json::json!(value));
2237 }
2238 if let Some(value) = limit {
2239 payload.insert("limit".to_string(), serde_json::json!(value));
2240 }
2241 let payload = remove_empty_value(payload);
2242
2243 self.websocket_api_base
2244 .send_message::<Vec<models::TradesHistoricalResponseResultInner>>(
2245 "/trades.historical".trim_start_matches('/'),
2246 payload,
2247 WebsocketMessageSendOptions::new(),
2248 )
2249 .await
2250 .map_err(anyhow::Error::from)?
2251 .into_iter()
2252 .next()
2253 .ok_or(WebsocketError::NoResponse)
2254 .map_err(anyhow::Error::from)
2255 }
2256
2257 async fn trades_recent(
2258 &self,
2259 params: TradesRecentParams,
2260 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::TradesRecentResponseResultInner>>> {
2261 let TradesRecentParams { symbol, id, limit } = params;
2262
2263 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
2264 payload.insert("symbol".to_string(), serde_json::json!(symbol));
2265 if let Some(value) = id {
2266 payload.insert("id".to_string(), serde_json::json!(value));
2267 }
2268 if let Some(value) = limit {
2269 payload.insert("limit".to_string(), serde_json::json!(value));
2270 }
2271 let payload = remove_empty_value(payload);
2272
2273 self.websocket_api_base
2274 .send_message::<Vec<models::TradesRecentResponseResultInner>>(
2275 "/trades.recent".trim_start_matches('/'),
2276 payload,
2277 WebsocketMessageSendOptions::new(),
2278 )
2279 .await
2280 .map_err(anyhow::Error::from)?
2281 .into_iter()
2282 .next()
2283 .ok_or(WebsocketError::NoResponse)
2284 .map_err(anyhow::Error::from)
2285 }
2286
2287 async fn ui_klines(
2288 &self,
2289 params: UiKlinesParams,
2290 ) -> anyhow::Result<WebsocketApiResponse<Vec<Vec<models::KlinesResponseResultInnerInner>>>>
2291 {
2292 let UiKlinesParams {
2293 symbol,
2294 interval,
2295 id,
2296 start_time,
2297 end_time,
2298 time_zone,
2299 limit,
2300 } = params;
2301
2302 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
2303 payload.insert("symbol".to_string(), serde_json::json!(symbol));
2304 payload.insert("interval".to_string(), serde_json::json!(interval));
2305 if let Some(value) = id {
2306 payload.insert("id".to_string(), serde_json::json!(value));
2307 }
2308 if let Some(value) = start_time {
2309 payload.insert("startTime".to_string(), serde_json::json!(value));
2310 }
2311 if let Some(value) = end_time {
2312 payload.insert("endTime".to_string(), serde_json::json!(value));
2313 }
2314 if let Some(value) = time_zone {
2315 payload.insert("timeZone".to_string(), serde_json::json!(value));
2316 }
2317 if let Some(value) = limit {
2318 payload.insert("limit".to_string(), serde_json::json!(value));
2319 }
2320 let payload = remove_empty_value(payload);
2321
2322 self.websocket_api_base
2323 .send_message::<Vec<Vec<models::KlinesResponseResultInnerInner>>>(
2324 "/uiKlines".trim_start_matches('/'),
2325 payload,
2326 WebsocketMessageSendOptions::new(),
2327 )
2328 .await
2329 .map_err(anyhow::Error::from)?
2330 .into_iter()
2331 .next()
2332 .ok_or(WebsocketError::NoResponse)
2333 .map_err(anyhow::Error::from)
2334 }
2335}
2336
2337#[cfg(all(test, feature = "spot"))]
2338mod tests {
2339 use super::*;
2340 use crate::TOKIO_SHARED_RT;
2341 use crate::common::websocket::{WebsocketApi, WebsocketConnection, WebsocketHandler};
2342 use crate::config::ConfigurationWebsocketApi;
2343 use crate::errors::WebsocketError;
2344 use crate::models::WebsocketApiRateLimit;
2345 use serde_json::{Value, json};
2346 use tokio::spawn;
2347 use tokio::sync::mpsc::{UnboundedReceiver, unbounded_channel};
2348 use tokio::time::{Duration, timeout};
2349 use tokio_tungstenite::tungstenite::Message;
2350
2351 async fn setup() -> (
2352 Arc<WebsocketApi>,
2353 Arc<WebsocketConnection>,
2354 UnboundedReceiver<Message>,
2355 ) {
2356 let conn = WebsocketConnection::new("test-conn");
2357 let (tx, rx) = unbounded_channel::<Message>();
2358 {
2359 let mut conn_state = conn.state.lock().await;
2360 conn_state.ws_write_tx = Some(tx);
2361 }
2362
2363 let config = ConfigurationWebsocketApi::builder()
2364 .api_key("key")
2365 .api_secret("secret")
2366 .build()
2367 .expect("Failed to build configuration");
2368 let ws_api = WebsocketApi::new(config, vec![conn.clone()]);
2369 conn.set_handler(ws_api.clone() as Arc<dyn WebsocketHandler>)
2370 .await;
2371 ws_api.clone().connect().await.unwrap();
2372
2373 (ws_api, conn, rx)
2374 }
2375
2376 #[test]
2377 fn avg_price_success() {
2378 TOKIO_SHARED_RT.block_on(async {
2379 let (ws_api, conn, mut rx) = setup().await;
2380 let client = MarketApiClient::new(ws_api.clone());
2381
2382 let handle = spawn(async move {
2383 let params = AvgPriceParams::builder("BNBUSDT".to_string(),).build().unwrap();
2384 client.avg_price(params).await
2385 });
2386
2387 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2388 let Message::Text(text) = sent else { panic!() };
2389 let v: Value = serde_json::from_str(&text).unwrap();
2390 let id = v["id"].as_str().unwrap();
2391 assert_eq!(v["method"], "/avgPrice".trim_start_matches('/'));
2392 let mut resp_json: Value = serde_json::from_str(r#"{"id":"ddbfb65f-9ebf-42ec-8240-8f0f91de0867","status":200,"result":{"mins":5,"closeTime":1694061154503},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2393 resp_json["id"] = id.into();
2394
2395 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2396 let expected_data: Box<models::AvgPriceResponseResult> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2397 let empty_array = Value::Array(vec![]);
2398 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2399 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2400 match raw_rate_limits.as_array() {
2401 Some(arr) if arr.is_empty() => None,
2402 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2403 None => None,
2404 };
2405
2406 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2407
2408 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2409
2410
2411 let response_rate_limits = response.rate_limits.clone();
2412 let response_data = response.data().expect("deserialize data");
2413
2414 assert_eq!(response_rate_limits, expected_rate_limits);
2415 assert_eq!(response_data, expected_data);
2416 });
2417 }
2418
2419 #[test]
2420 fn avg_price_error_response() {
2421 TOKIO_SHARED_RT.block_on(async {
2422 let (ws_api, conn, mut rx) = setup().await;
2423 let client = MarketApiClient::new(ws_api.clone());
2424
2425 let handle = tokio::spawn(async move {
2426 let params = AvgPriceParams::builder("BNBUSDT".to_string(),).build().unwrap();
2427 client.avg_price(params).await
2428 });
2429
2430 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2431 let Message::Text(text) = sent else { panic!() };
2432 let v: Value = serde_json::from_str(&text).unwrap();
2433 let id = v["id"].as_str().unwrap().to_string();
2434
2435 let resp_json = json!({
2436 "id": id,
2437 "status": 400,
2438 "error": {
2439 "code": -2010,
2440 "msg": "Account has insufficient balance for requested action.",
2441 },
2442 "rateLimits": [
2443 {
2444 "rateLimitType": "ORDERS",
2445 "interval": "SECOND",
2446 "intervalNum": 10,
2447 "limit": 50,
2448 "count": 13
2449 },
2450 ],
2451 });
2452 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2453
2454 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2455 match join {
2456 Ok(Err(e)) => {
2457 let msg = e.to_string();
2458 assert!(
2459 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2460 "Expected error msg to contain server error, got: {msg}"
2461 );
2462 }
2463 Ok(Ok(_)) => panic!("Expected error"),
2464 Err(_) => panic!("Task panicked"),
2465 }
2466 });
2467 }
2468
2469 #[test]
2470 fn avg_price_request_timeout() {
2471 TOKIO_SHARED_RT.block_on(async {
2472 let (ws_api, _conn, mut rx) = setup().await;
2473 let client = MarketApiClient::new(ws_api.clone());
2474
2475 let handle = spawn(async move {
2476 let params = AvgPriceParams::builder("BNBUSDT".to_string())
2477 .build()
2478 .unwrap();
2479 client.avg_price(params).await
2480 });
2481
2482 let sent = timeout(Duration::from_secs(1), rx.recv())
2483 .await
2484 .expect("send should occur")
2485 .expect("channel closed");
2486 let Message::Text(text) = sent else {
2487 panic!("expected Message Text")
2488 };
2489
2490 let _: Value = serde_json::from_str(&text).unwrap();
2491
2492 let result = handle.await.expect("task completed");
2493 match result {
2494 Err(e) => {
2495 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2496 assert!(matches!(inner, WebsocketError::Timeout));
2497 } else {
2498 panic!("Unexpected error type: {:?}", e);
2499 }
2500 }
2501 Ok(_) => panic!("Expected timeout error"),
2502 }
2503 });
2504 }
2505
2506 #[test]
2507 fn block_trades_historical_success() {
2508 TOKIO_SHARED_RT.block_on(async {
2509 let (ws_api, conn, mut rx) = setup().await;
2510 let client = MarketApiClient::new(ws_api.clone());
2511
2512 let handle = spawn(async move {
2513 let params = BlockTradesHistoricalParams::builder("BNBBTC".to_string(),582,).build().unwrap();
2514 client.block_trades_historical(params).await
2515 });
2516
2517 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2518 let Message::Text(text) = sent else { panic!() };
2519 let v: Value = serde_json::from_str(&text).unwrap();
2520 let id = v["id"].as_str().unwrap();
2521 assert_eq!(v["method"], "/blockTrades.historical".trim_start_matches('/'));
2522 let mut resp_json: Value = serde_json::from_str(r#"{"id":"cffc9c7d-4efc-4ce0-b587-6b87448f052a","status":200,"result":[{"id":582,"price":"0.052","qty":"5838","quoteQty":"303.576","time":1772506983321,"isBuyerMaker":true}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":10}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2523 resp_json["id"] = id.into();
2524
2525 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2526 let expected_data: Vec<models::BlockTradesHistoricalResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2527 let empty_array = Value::Array(vec![]);
2528 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2529 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2530 match raw_rate_limits.as_array() {
2531 Some(arr) if arr.is_empty() => None,
2532 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2533 None => None,
2534 };
2535
2536 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2537
2538 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2539
2540
2541 let response_rate_limits = response.rate_limits.clone();
2542 let response_data = response.data().expect("deserialize data");
2543
2544 assert_eq!(response_rate_limits, expected_rate_limits);
2545 assert_eq!(response_data, expected_data);
2546 });
2547 }
2548
2549 #[test]
2550 fn block_trades_historical_error_response() {
2551 TOKIO_SHARED_RT.block_on(async {
2552 let (ws_api, conn, mut rx) = setup().await;
2553 let client = MarketApiClient::new(ws_api.clone());
2554
2555 let handle = tokio::spawn(async move {
2556 let params = BlockTradesHistoricalParams::builder("BNBBTC".to_string(),582,).build().unwrap();
2557 client.block_trades_historical(params).await
2558 });
2559
2560 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2561 let Message::Text(text) = sent else { panic!() };
2562 let v: Value = serde_json::from_str(&text).unwrap();
2563 let id = v["id"].as_str().unwrap().to_string();
2564
2565 let resp_json = json!({
2566 "id": id,
2567 "status": 400,
2568 "error": {
2569 "code": -2010,
2570 "msg": "Account has insufficient balance for requested action.",
2571 },
2572 "rateLimits": [
2573 {
2574 "rateLimitType": "ORDERS",
2575 "interval": "SECOND",
2576 "intervalNum": 10,
2577 "limit": 50,
2578 "count": 13
2579 },
2580 ],
2581 });
2582 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2583
2584 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2585 match join {
2586 Ok(Err(e)) => {
2587 let msg = e.to_string();
2588 assert!(
2589 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2590 "Expected error msg to contain server error, got: {msg}"
2591 );
2592 }
2593 Ok(Ok(_)) => panic!("Expected error"),
2594 Err(_) => panic!("Task panicked"),
2595 }
2596 });
2597 }
2598
2599 #[test]
2600 fn block_trades_historical_request_timeout() {
2601 TOKIO_SHARED_RT.block_on(async {
2602 let (ws_api, _conn, mut rx) = setup().await;
2603 let client = MarketApiClient::new(ws_api.clone());
2604
2605 let handle = spawn(async move {
2606 let params = BlockTradesHistoricalParams::builder("BNBBTC".to_string(), 582)
2607 .build()
2608 .unwrap();
2609 client.block_trades_historical(params).await
2610 });
2611
2612 let sent = timeout(Duration::from_secs(1), rx.recv())
2613 .await
2614 .expect("send should occur")
2615 .expect("channel closed");
2616 let Message::Text(text) = sent else {
2617 panic!("expected Message Text")
2618 };
2619
2620 let _: Value = serde_json::from_str(&text).unwrap();
2621
2622 let result = handle.await.expect("task completed");
2623 match result {
2624 Err(e) => {
2625 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2626 assert!(matches!(inner, WebsocketError::Timeout));
2627 } else {
2628 panic!("Unexpected error type: {:?}", e);
2629 }
2630 }
2631 Ok(_) => panic!("Expected timeout error"),
2632 }
2633 });
2634 }
2635
2636 #[test]
2637 fn depth_success() {
2638 TOKIO_SHARED_RT.block_on(async {
2639 let (ws_api, conn, mut rx) = setup().await;
2640 let client = MarketApiClient::new(ws_api.clone());
2641
2642 let handle = spawn(async move {
2643 let params = DepthParams::builder("BNBUSDT".to_string(),).build().unwrap();
2644 client.depth(params).await
2645 });
2646
2647 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2648 let Message::Text(text) = sent else { panic!() };
2649 let v: Value = serde_json::from_str(&text).unwrap();
2650 let id = v["id"].as_str().unwrap();
2651 assert_eq!(v["method"], "/depth".trim_start_matches('/'));
2652 let mut resp_json: Value = serde_json::from_str(r#"{"id":"51e2affb-0aba-4821-ba75-f2625006eb43","status":200,"result":{"lastUpdateId":2731179239,"bids":[["0.01379900","3.43200000"]],"asks":[["0.01380000","5.91700000"]]},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2653 resp_json["id"] = id.into();
2654
2655 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2656 let expected_data: Box<models::DepthResponseResult> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2657 let empty_array = Value::Array(vec![]);
2658 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2659 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2660 match raw_rate_limits.as_array() {
2661 Some(arr) if arr.is_empty() => None,
2662 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2663 None => None,
2664 };
2665
2666 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2667
2668 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2669
2670
2671 let response_rate_limits = response.rate_limits.clone();
2672 let response_data = response.data().expect("deserialize data");
2673
2674 assert_eq!(response_rate_limits, expected_rate_limits);
2675 assert_eq!(response_data, expected_data);
2676 });
2677 }
2678
2679 #[test]
2680 fn depth_error_response() {
2681 TOKIO_SHARED_RT.block_on(async {
2682 let (ws_api, conn, mut rx) = setup().await;
2683 let client = MarketApiClient::new(ws_api.clone());
2684
2685 let handle = tokio::spawn(async move {
2686 let params = DepthParams::builder("BNBUSDT".to_string(),).build().unwrap();
2687 client.depth(params).await
2688 });
2689
2690 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2691 let Message::Text(text) = sent else { panic!() };
2692 let v: Value = serde_json::from_str(&text).unwrap();
2693 let id = v["id"].as_str().unwrap().to_string();
2694
2695 let resp_json = json!({
2696 "id": id,
2697 "status": 400,
2698 "error": {
2699 "code": -2010,
2700 "msg": "Account has insufficient balance for requested action.",
2701 },
2702 "rateLimits": [
2703 {
2704 "rateLimitType": "ORDERS",
2705 "interval": "SECOND",
2706 "intervalNum": 10,
2707 "limit": 50,
2708 "count": 13
2709 },
2710 ],
2711 });
2712 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2713
2714 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2715 match join {
2716 Ok(Err(e)) => {
2717 let msg = e.to_string();
2718 assert!(
2719 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2720 "Expected error msg to contain server error, got: {msg}"
2721 );
2722 }
2723 Ok(Ok(_)) => panic!("Expected error"),
2724 Err(_) => panic!("Task panicked"),
2725 }
2726 });
2727 }
2728
2729 #[test]
2730 fn depth_request_timeout() {
2731 TOKIO_SHARED_RT.block_on(async {
2732 let (ws_api, _conn, mut rx) = setup().await;
2733 let client = MarketApiClient::new(ws_api.clone());
2734
2735 let handle = spawn(async move {
2736 let params = DepthParams::builder("BNBUSDT".to_string()).build().unwrap();
2737 client.depth(params).await
2738 });
2739
2740 let sent = timeout(Duration::from_secs(1), rx.recv())
2741 .await
2742 .expect("send should occur")
2743 .expect("channel closed");
2744 let Message::Text(text) = sent else {
2745 panic!("expected Message Text")
2746 };
2747
2748 let _: Value = serde_json::from_str(&text).unwrap();
2749
2750 let result = handle.await.expect("task completed");
2751 match result {
2752 Err(e) => {
2753 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2754 assert!(matches!(inner, WebsocketError::Timeout));
2755 } else {
2756 panic!("Unexpected error type: {:?}", e);
2757 }
2758 }
2759 Ok(_) => panic!("Expected timeout error"),
2760 }
2761 });
2762 }
2763
2764 #[test]
2765 fn klines_success() {
2766 TOKIO_SHARED_RT.block_on(async {
2767 let (ws_api, conn, mut rx) = setup().await;
2768 let client = MarketApiClient::new(ws_api.clone());
2769
2770 let handle = spawn(async move {
2771 let params = KlinesParams::builder("BNBUSDT".to_string(),KlinesIntervalEnum::Interval1s,).build().unwrap();
2772 client.klines(params).await
2773 });
2774
2775 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2776 let Message::Text(text) = sent else { panic!() };
2777 let v: Value = serde_json::from_str(&text).unwrap();
2778 let id = v["id"].as_str().unwrap();
2779 assert_eq!(v["method"], "/klines".trim_start_matches('/'));
2780 let mut resp_json: Value = serde_json::from_str(r#"{"id":"1dbbeb56-8eea-466a-8f6e-86bdcfa2fc0b","status":200,"result":[[1499040000000]],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2781 resp_json["id"] = id.into();
2782
2783 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2784 let expected_data: Vec<Vec<models::KlinesResponseResultInnerInner>> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2785 let empty_array = Value::Array(vec![]);
2786 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2787 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2788 match raw_rate_limits.as_array() {
2789 Some(arr) if arr.is_empty() => None,
2790 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2791 None => None,
2792 };
2793
2794 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2795
2796 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2797
2798
2799 let response_rate_limits = response.rate_limits.clone();
2800 let response_data = response.data().expect("deserialize data");
2801
2802 assert_eq!(response_rate_limits, expected_rate_limits);
2803 assert_eq!(response_data, expected_data);
2804 });
2805 }
2806
2807 #[test]
2808 fn klines_error_response() {
2809 TOKIO_SHARED_RT.block_on(async {
2810 let (ws_api, conn, mut rx) = setup().await;
2811 let client = MarketApiClient::new(ws_api.clone());
2812
2813 let handle = tokio::spawn(async move {
2814 let params = KlinesParams::builder("BNBUSDT".to_string(),KlinesIntervalEnum::Interval1s,).build().unwrap();
2815 client.klines(params).await
2816 });
2817
2818 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2819 let Message::Text(text) = sent else { panic!() };
2820 let v: Value = serde_json::from_str(&text).unwrap();
2821 let id = v["id"].as_str().unwrap().to_string();
2822
2823 let resp_json = json!({
2824 "id": id,
2825 "status": 400,
2826 "error": {
2827 "code": -2010,
2828 "msg": "Account has insufficient balance for requested action.",
2829 },
2830 "rateLimits": [
2831 {
2832 "rateLimitType": "ORDERS",
2833 "interval": "SECOND",
2834 "intervalNum": 10,
2835 "limit": 50,
2836 "count": 13
2837 },
2838 ],
2839 });
2840 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2841
2842 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2843 match join {
2844 Ok(Err(e)) => {
2845 let msg = e.to_string();
2846 assert!(
2847 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2848 "Expected error msg to contain server error, got: {msg}"
2849 );
2850 }
2851 Ok(Ok(_)) => panic!("Expected error"),
2852 Err(_) => panic!("Task panicked"),
2853 }
2854 });
2855 }
2856
2857 #[test]
2858 fn klines_request_timeout() {
2859 TOKIO_SHARED_RT.block_on(async {
2860 let (ws_api, _conn, mut rx) = setup().await;
2861 let client = MarketApiClient::new(ws_api.clone());
2862
2863 let handle = spawn(async move {
2864 let params =
2865 KlinesParams::builder("BNBUSDT".to_string(), KlinesIntervalEnum::Interval1s)
2866 .build()
2867 .unwrap();
2868 client.klines(params).await
2869 });
2870
2871 let sent = timeout(Duration::from_secs(1), rx.recv())
2872 .await
2873 .expect("send should occur")
2874 .expect("channel closed");
2875 let Message::Text(text) = sent else {
2876 panic!("expected Message Text")
2877 };
2878
2879 let _: Value = serde_json::from_str(&text).unwrap();
2880
2881 let result = handle.await.expect("task completed");
2882 match result {
2883 Err(e) => {
2884 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2885 assert!(matches!(inner, WebsocketError::Timeout));
2886 } else {
2887 panic!("Unexpected error type: {:?}", e);
2888 }
2889 }
2890 Ok(_) => panic!("Expected timeout error"),
2891 }
2892 });
2893 }
2894
2895 #[test]
2896 fn reference_price_success() {
2897 TOKIO_SHARED_RT.block_on(async {
2898 let (ws_api, conn, mut rx) = setup().await;
2899 let client = MarketApiClient::new(ws_api.clone());
2900
2901 let handle = spawn(async move {
2902 let params = ReferencePriceParams::builder("BAZUSD".to_string(),).build().unwrap();
2903 client.reference_price(params).await
2904 });
2905
2906 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2907 let Message::Text(text) = sent else { panic!() };
2908 let v: Value = serde_json::from_str(&text).unwrap();
2909 let id = v["id"].as_str().unwrap();
2910 assert_eq!(v["method"], "/referencePrice".trim_start_matches('/'));
2911 let mut resp_json: Value = serde_json::from_str(r#"{"id":"ddbfb65f-9ebf-42ec-8240-8f0f91de0867","status":200,"result":{"symbol":"BAZUSD","referencePrice":"0.00501900","timestamp":1770946889251,"code":-2043,"msg":"This symbol doesn't have a reference price."},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2912 resp_json["id"] = id.into();
2913
2914 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2915 let expected_data: Box<models::ReferencePriceResponseResult> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2916 let empty_array = Value::Array(vec![]);
2917 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2918 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2919 match raw_rate_limits.as_array() {
2920 Some(arr) if arr.is_empty() => None,
2921 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2922 None => None,
2923 };
2924
2925 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2926
2927 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2928
2929
2930 let response_rate_limits = response.rate_limits.clone();
2931 let response_data = response.data().expect("deserialize data");
2932
2933 assert_eq!(response_rate_limits, expected_rate_limits);
2934 assert_eq!(response_data, expected_data);
2935 });
2936 }
2937
2938 #[test]
2939 fn reference_price_error_response() {
2940 TOKIO_SHARED_RT.block_on(async {
2941 let (ws_api, conn, mut rx) = setup().await;
2942 let client = MarketApiClient::new(ws_api.clone());
2943
2944 let handle = tokio::spawn(async move {
2945 let params = ReferencePriceParams::builder("BAZUSD".to_string(),).build().unwrap();
2946 client.reference_price(params).await
2947 });
2948
2949 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2950 let Message::Text(text) = sent else { panic!() };
2951 let v: Value = serde_json::from_str(&text).unwrap();
2952 let id = v["id"].as_str().unwrap().to_string();
2953
2954 let resp_json = json!({
2955 "id": id,
2956 "status": 400,
2957 "error": {
2958 "code": -2010,
2959 "msg": "Account has insufficient balance for requested action.",
2960 },
2961 "rateLimits": [
2962 {
2963 "rateLimitType": "ORDERS",
2964 "interval": "SECOND",
2965 "intervalNum": 10,
2966 "limit": 50,
2967 "count": 13
2968 },
2969 ],
2970 });
2971 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2972
2973 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2974 match join {
2975 Ok(Err(e)) => {
2976 let msg = e.to_string();
2977 assert!(
2978 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2979 "Expected error msg to contain server error, got: {msg}"
2980 );
2981 }
2982 Ok(Ok(_)) => panic!("Expected error"),
2983 Err(_) => panic!("Task panicked"),
2984 }
2985 });
2986 }
2987
2988 #[test]
2989 fn reference_price_request_timeout() {
2990 TOKIO_SHARED_RT.block_on(async {
2991 let (ws_api, _conn, mut rx) = setup().await;
2992 let client = MarketApiClient::new(ws_api.clone());
2993
2994 let handle = spawn(async move {
2995 let params = ReferencePriceParams::builder("BAZUSD".to_string())
2996 .build()
2997 .unwrap();
2998 client.reference_price(params).await
2999 });
3000
3001 let sent = timeout(Duration::from_secs(1), rx.recv())
3002 .await
3003 .expect("send should occur")
3004 .expect("channel closed");
3005 let Message::Text(text) = sent else {
3006 panic!("expected Message Text")
3007 };
3008
3009 let _: Value = serde_json::from_str(&text).unwrap();
3010
3011 let result = handle.await.expect("task completed");
3012 match result {
3013 Err(e) => {
3014 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3015 assert!(matches!(inner, WebsocketError::Timeout));
3016 } else {
3017 panic!("Unexpected error type: {:?}", e);
3018 }
3019 }
3020 Ok(_) => panic!("Expected timeout error"),
3021 }
3022 });
3023 }
3024
3025 #[test]
3026 fn reference_price_calculation_success() {
3027 TOKIO_SHARED_RT.block_on(async {
3028 let (ws_api, conn, mut rx) = setup().await;
3029 let client = MarketApiClient::new(ws_api.clone());
3030
3031 let handle = spawn(async move {
3032 let params = ReferencePriceCalculationParams::builder("BAZUSD".to_string(),).build().unwrap();
3033 client.reference_price_calculation(params).await
3034 });
3035
3036 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3037 let Message::Text(text) = sent else { panic!() };
3038 let v: Value = serde_json::from_str(&text).unwrap();
3039 let id = v["id"].as_str().unwrap();
3040 assert_eq!(v["method"], "/referencePrice.calculation".trim_start_matches('/'));
3041 let mut resp_json: Value = serde_json::from_str(r#"{"id":"ddbfb65f-9ebf-42ec-8240-8f0f91de0867","status":200,"result":{"symbol":"BAZUSD","calculationType":"ARITHMETIC_MEAN","bucketCount":10,"bucketWidthMs":1000,"externalCalculationId":42},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3042 resp_json["id"] = id.into();
3043
3044 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3045 let expected_data: Box<models::ReferencePriceCalculationResponseResult> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3046 let empty_array = Value::Array(vec![]);
3047 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3048 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3049 match raw_rate_limits.as_array() {
3050 Some(arr) if arr.is_empty() => None,
3051 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3052 None => None,
3053 };
3054
3055 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3056
3057 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3058
3059
3060 let response_rate_limits = response.rate_limits.clone();
3061 let response_data = response.data().expect("deserialize data");
3062
3063 assert_eq!(response_rate_limits, expected_rate_limits);
3064 assert_eq!(response_data, expected_data);
3065 });
3066 }
3067
3068 #[test]
3069 fn reference_price_calculation_error_response() {
3070 TOKIO_SHARED_RT.block_on(async {
3071 let (ws_api, conn, mut rx) = setup().await;
3072 let client = MarketApiClient::new(ws_api.clone());
3073
3074 let handle = tokio::spawn(async move {
3075 let params = ReferencePriceCalculationParams::builder("BAZUSD".to_string(),).build().unwrap();
3076 client.reference_price_calculation(params).await
3077 });
3078
3079 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3080 let Message::Text(text) = sent else { panic!() };
3081 let v: Value = serde_json::from_str(&text).unwrap();
3082 let id = v["id"].as_str().unwrap().to_string();
3083
3084 let resp_json = json!({
3085 "id": id,
3086 "status": 400,
3087 "error": {
3088 "code": -2010,
3089 "msg": "Account has insufficient balance for requested action.",
3090 },
3091 "rateLimits": [
3092 {
3093 "rateLimitType": "ORDERS",
3094 "interval": "SECOND",
3095 "intervalNum": 10,
3096 "limit": 50,
3097 "count": 13
3098 },
3099 ],
3100 });
3101 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3102
3103 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3104 match join {
3105 Ok(Err(e)) => {
3106 let msg = e.to_string();
3107 assert!(
3108 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3109 "Expected error msg to contain server error, got: {msg}"
3110 );
3111 }
3112 Ok(Ok(_)) => panic!("Expected error"),
3113 Err(_) => panic!("Task panicked"),
3114 }
3115 });
3116 }
3117
3118 #[test]
3119 fn reference_price_calculation_request_timeout() {
3120 TOKIO_SHARED_RT.block_on(async {
3121 let (ws_api, _conn, mut rx) = setup().await;
3122 let client = MarketApiClient::new(ws_api.clone());
3123
3124 let handle = spawn(async move {
3125 let params = ReferencePriceCalculationParams::builder("BAZUSD".to_string())
3126 .build()
3127 .unwrap();
3128 client.reference_price_calculation(params).await
3129 });
3130
3131 let sent = timeout(Duration::from_secs(1), rx.recv())
3132 .await
3133 .expect("send should occur")
3134 .expect("channel closed");
3135 let Message::Text(text) = sent else {
3136 panic!("expected Message Text")
3137 };
3138
3139 let _: Value = serde_json::from_str(&text).unwrap();
3140
3141 let result = handle.await.expect("task completed");
3142 match result {
3143 Err(e) => {
3144 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3145 assert!(matches!(inner, WebsocketError::Timeout));
3146 } else {
3147 panic!("Unexpected error type: {:?}", e);
3148 }
3149 }
3150 Ok(_) => panic!("Expected timeout error"),
3151 }
3152 });
3153 }
3154
3155 #[test]
3156 fn ticker_success() {
3157 TOKIO_SHARED_RT.block_on(async {
3158 let (ws_api, conn, mut rx) = setup().await;
3159 let client = MarketApiClient::new(ws_api.clone());
3160
3161 let handle = spawn(async move {
3162 let params = TickerParams::builder().build().unwrap();
3163 client.ticker(params).await
3164 });
3165
3166 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3167 let Message::Text(text) = sent else { panic!() };
3168 let v: Value = serde_json::from_str(&text).unwrap();
3169 let id = v["id"].as_str().unwrap();
3170 assert_eq!(v["method"], "/ticker".trim_start_matches('/'));
3171 let mut resp_json: Value = serde_json::from_str(r#"{"id":"bdb7c503-542c-495c-b797-4d2ee2e91173","status":200,"result":{"symbol":"BNBBTC","openTime":1659580020000,"closeTime":1660184865291,"firstId":192977765,"lastId":195365758,"count":2387994},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3172 resp_json["id"] = id.into();
3173
3174 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3175 let expected_data: models::TickerResponse = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3176 let empty_array = Value::Array(vec![]);
3177 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3178 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3179 match raw_rate_limits.as_array() {
3180 Some(arr) if arr.is_empty() => None,
3181 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3182 None => None,
3183 };
3184
3185 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3186
3187 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3188
3189
3190 let response_rate_limits = response.rate_limits.clone();
3191 let response_data = response.data().expect("deserialize data");
3192
3193 assert_eq!(response_rate_limits, expected_rate_limits);
3194 assert_eq!(response_data, expected_data);
3195 });
3196 }
3197
3198 #[test]
3199 fn ticker_error_response() {
3200 TOKIO_SHARED_RT.block_on(async {
3201 let (ws_api, conn, mut rx) = setup().await;
3202 let client = MarketApiClient::new(ws_api.clone());
3203
3204 let handle = tokio::spawn(async move {
3205 let params = TickerParams::builder().build().unwrap();
3206 client.ticker(params).await
3207 });
3208
3209 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3210 let Message::Text(text) = sent else { panic!() };
3211 let v: Value = serde_json::from_str(&text).unwrap();
3212 let id = v["id"].as_str().unwrap().to_string();
3213
3214 let resp_json = json!({
3215 "id": id,
3216 "status": 400,
3217 "error": {
3218 "code": -2010,
3219 "msg": "Account has insufficient balance for requested action.",
3220 },
3221 "rateLimits": [
3222 {
3223 "rateLimitType": "ORDERS",
3224 "interval": "SECOND",
3225 "intervalNum": 10,
3226 "limit": 50,
3227 "count": 13
3228 },
3229 ],
3230 });
3231 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3232
3233 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3234 match join {
3235 Ok(Err(e)) => {
3236 let msg = e.to_string();
3237 assert!(
3238 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3239 "Expected error msg to contain server error, got: {msg}"
3240 );
3241 }
3242 Ok(Ok(_)) => panic!("Expected error"),
3243 Err(_) => panic!("Task panicked"),
3244 }
3245 });
3246 }
3247
3248 #[test]
3249 fn ticker_request_timeout() {
3250 TOKIO_SHARED_RT.block_on(async {
3251 let (ws_api, _conn, mut rx) = setup().await;
3252 let client = MarketApiClient::new(ws_api.clone());
3253
3254 let handle = spawn(async move {
3255 let params = TickerParams::builder().build().unwrap();
3256 client.ticker(params).await
3257 });
3258
3259 let sent = timeout(Duration::from_secs(1), rx.recv())
3260 .await
3261 .expect("send should occur")
3262 .expect("channel closed");
3263 let Message::Text(text) = sent else {
3264 panic!("expected Message Text")
3265 };
3266
3267 let _: Value = serde_json::from_str(&text).unwrap();
3268
3269 let result = handle.await.expect("task completed");
3270 match result {
3271 Err(e) => {
3272 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3273 assert!(matches!(inner, WebsocketError::Timeout));
3274 } else {
3275 panic!("Unexpected error type: {:?}", e);
3276 }
3277 }
3278 Ok(_) => panic!("Expected timeout error"),
3279 }
3280 });
3281 }
3282
3283 #[test]
3284 fn ticker24hr_success() {
3285 TOKIO_SHARED_RT.block_on(async {
3286 let (ws_api, conn, mut rx) = setup().await;
3287 let client = MarketApiClient::new(ws_api.clone());
3288
3289 let handle = spawn(async move {
3290 let params = Ticker24hrParams::builder().build().unwrap();
3291 client.ticker24hr(params).await
3292 });
3293
3294 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3295 let Message::Text(text) = sent else { panic!() };
3296 let v: Value = serde_json::from_str(&text).unwrap();
3297 let id = v["id"].as_str().unwrap();
3298 assert_eq!(v["method"], "/ticker.24hr".trim_start_matches('/'));
3299 let mut resp_json: Value = serde_json::from_str(r#"{"id":"9fa2a91b-3fca-4ed7-a9ad-58e3b67483de","status":200,"result":{"symbol":"BNBBTC","openTime":1660014164909,"closeTime":1660100564909,"firstId":194696115,"lastId":194968287,"count":272173},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3300 resp_json["id"] = id.into();
3301
3302 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3303 let expected_data: models::Ticker24hrResponse = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3304 let empty_array = Value::Array(vec![]);
3305 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3306 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3307 match raw_rate_limits.as_array() {
3308 Some(arr) if arr.is_empty() => None,
3309 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3310 None => None,
3311 };
3312
3313 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3314
3315 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3316
3317
3318 let response_rate_limits = response.rate_limits.clone();
3319 let response_data = response.data().expect("deserialize data");
3320
3321 assert_eq!(response_rate_limits, expected_rate_limits);
3322 assert_eq!(response_data, expected_data);
3323 });
3324 }
3325
3326 #[test]
3327 fn ticker24hr_error_response() {
3328 TOKIO_SHARED_RT.block_on(async {
3329 let (ws_api, conn, mut rx) = setup().await;
3330 let client = MarketApiClient::new(ws_api.clone());
3331
3332 let handle = tokio::spawn(async move {
3333 let params = Ticker24hrParams::builder().build().unwrap();
3334 client.ticker24hr(params).await
3335 });
3336
3337 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3338 let Message::Text(text) = sent else { panic!() };
3339 let v: Value = serde_json::from_str(&text).unwrap();
3340 let id = v["id"].as_str().unwrap().to_string();
3341
3342 let resp_json = json!({
3343 "id": id,
3344 "status": 400,
3345 "error": {
3346 "code": -2010,
3347 "msg": "Account has insufficient balance for requested action.",
3348 },
3349 "rateLimits": [
3350 {
3351 "rateLimitType": "ORDERS",
3352 "interval": "SECOND",
3353 "intervalNum": 10,
3354 "limit": 50,
3355 "count": 13
3356 },
3357 ],
3358 });
3359 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3360
3361 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3362 match join {
3363 Ok(Err(e)) => {
3364 let msg = e.to_string();
3365 assert!(
3366 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3367 "Expected error msg to contain server error, got: {msg}"
3368 );
3369 }
3370 Ok(Ok(_)) => panic!("Expected error"),
3371 Err(_) => panic!("Task panicked"),
3372 }
3373 });
3374 }
3375
3376 #[test]
3377 fn ticker24hr_request_timeout() {
3378 TOKIO_SHARED_RT.block_on(async {
3379 let (ws_api, _conn, mut rx) = setup().await;
3380 let client = MarketApiClient::new(ws_api.clone());
3381
3382 let handle = spawn(async move {
3383 let params = Ticker24hrParams::builder().build().unwrap();
3384 client.ticker24hr(params).await
3385 });
3386
3387 let sent = timeout(Duration::from_secs(1), rx.recv())
3388 .await
3389 .expect("send should occur")
3390 .expect("channel closed");
3391 let Message::Text(text) = sent else {
3392 panic!("expected Message Text")
3393 };
3394
3395 let _: Value = serde_json::from_str(&text).unwrap();
3396
3397 let result = handle.await.expect("task completed");
3398 match result {
3399 Err(e) => {
3400 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3401 assert!(matches!(inner, WebsocketError::Timeout));
3402 } else {
3403 panic!("Unexpected error type: {:?}", e);
3404 }
3405 }
3406 Ok(_) => panic!("Expected timeout error"),
3407 }
3408 });
3409 }
3410
3411 #[test]
3412 fn ticker_book_success() {
3413 TOKIO_SHARED_RT.block_on(async {
3414 let (ws_api, conn, mut rx) = setup().await;
3415 let client = MarketApiClient::new(ws_api.clone());
3416
3417 let handle = spawn(async move {
3418 let params = TickerBookParams::builder().build().unwrap();
3419 client.ticker_book(params).await
3420 });
3421
3422 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3423 let Message::Text(text) = sent else { panic!() };
3424 let v: Value = serde_json::from_str(&text).unwrap();
3425 let id = v["id"].as_str().unwrap();
3426 assert_eq!(v["method"], "/ticker.book".trim_start_matches('/'));
3427 let mut resp_json: Value = serde_json::from_str(r#"{"id":"9d32157c-a556-4d27-9866-66760a174b57","status":200,"result":{"symbol":"BNBBTC"},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3428 resp_json["id"] = id.into();
3429
3430 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3431 let expected_data: models::TickerBookResponse = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3432 let empty_array = Value::Array(vec![]);
3433 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3434 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3435 match raw_rate_limits.as_array() {
3436 Some(arr) if arr.is_empty() => None,
3437 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3438 None => None,
3439 };
3440
3441 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3442
3443 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3444
3445
3446 let response_rate_limits = response.rate_limits.clone();
3447 let response_data = response.data().expect("deserialize data");
3448
3449 assert_eq!(response_rate_limits, expected_rate_limits);
3450 assert_eq!(response_data, expected_data);
3451 });
3452 }
3453
3454 #[test]
3455 fn ticker_book_error_response() {
3456 TOKIO_SHARED_RT.block_on(async {
3457 let (ws_api, conn, mut rx) = setup().await;
3458 let client = MarketApiClient::new(ws_api.clone());
3459
3460 let handle = tokio::spawn(async move {
3461 let params = TickerBookParams::builder().build().unwrap();
3462 client.ticker_book(params).await
3463 });
3464
3465 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3466 let Message::Text(text) = sent else { panic!() };
3467 let v: Value = serde_json::from_str(&text).unwrap();
3468 let id = v["id"].as_str().unwrap().to_string();
3469
3470 let resp_json = json!({
3471 "id": id,
3472 "status": 400,
3473 "error": {
3474 "code": -2010,
3475 "msg": "Account has insufficient balance for requested action.",
3476 },
3477 "rateLimits": [
3478 {
3479 "rateLimitType": "ORDERS",
3480 "interval": "SECOND",
3481 "intervalNum": 10,
3482 "limit": 50,
3483 "count": 13
3484 },
3485 ],
3486 });
3487 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3488
3489 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3490 match join {
3491 Ok(Err(e)) => {
3492 let msg = e.to_string();
3493 assert!(
3494 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3495 "Expected error msg to contain server error, got: {msg}"
3496 );
3497 }
3498 Ok(Ok(_)) => panic!("Expected error"),
3499 Err(_) => panic!("Task panicked"),
3500 }
3501 });
3502 }
3503
3504 #[test]
3505 fn ticker_book_request_timeout() {
3506 TOKIO_SHARED_RT.block_on(async {
3507 let (ws_api, _conn, mut rx) = setup().await;
3508 let client = MarketApiClient::new(ws_api.clone());
3509
3510 let handle = spawn(async move {
3511 let params = TickerBookParams::builder().build().unwrap();
3512 client.ticker_book(params).await
3513 });
3514
3515 let sent = timeout(Duration::from_secs(1), rx.recv())
3516 .await
3517 .expect("send should occur")
3518 .expect("channel closed");
3519 let Message::Text(text) = sent else {
3520 panic!("expected Message Text")
3521 };
3522
3523 let _: Value = serde_json::from_str(&text).unwrap();
3524
3525 let result = handle.await.expect("task completed");
3526 match result {
3527 Err(e) => {
3528 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3529 assert!(matches!(inner, WebsocketError::Timeout));
3530 } else {
3531 panic!("Unexpected error type: {:?}", e);
3532 }
3533 }
3534 Ok(_) => panic!("Expected timeout error"),
3535 }
3536 });
3537 }
3538
3539 #[test]
3540 fn ticker_price_success() {
3541 TOKIO_SHARED_RT.block_on(async {
3542 let (ws_api, conn, mut rx) = setup().await;
3543 let client = MarketApiClient::new(ws_api.clone());
3544
3545 let handle = spawn(async move {
3546 let params = TickerPriceParams::builder().build().unwrap();
3547 client.ticker_price(params).await
3548 });
3549
3550 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3551 let Message::Text(text) = sent else { panic!() };
3552 let v: Value = serde_json::from_str(&text).unwrap();
3553 let id = v["id"].as_str().unwrap();
3554 assert_eq!(v["method"], "/ticker.price".trim_start_matches('/'));
3555 let mut resp_json: Value = serde_json::from_str(r#"{"id":"043a7cf2-bde3-4888-9604-c8ac41fcba4d","status":200,"result":{"symbol":"BNBBTC"},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3556 resp_json["id"] = id.into();
3557
3558 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3559 let expected_data: models::TickerPriceResponse = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3560 let empty_array = Value::Array(vec![]);
3561 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3562 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3563 match raw_rate_limits.as_array() {
3564 Some(arr) if arr.is_empty() => None,
3565 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3566 None => None,
3567 };
3568
3569 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3570
3571 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3572
3573
3574 let response_rate_limits = response.rate_limits.clone();
3575 let response_data = response.data().expect("deserialize data");
3576
3577 assert_eq!(response_rate_limits, expected_rate_limits);
3578 assert_eq!(response_data, expected_data);
3579 });
3580 }
3581
3582 #[test]
3583 fn ticker_price_error_response() {
3584 TOKIO_SHARED_RT.block_on(async {
3585 let (ws_api, conn, mut rx) = setup().await;
3586 let client = MarketApiClient::new(ws_api.clone());
3587
3588 let handle = tokio::spawn(async move {
3589 let params = TickerPriceParams::builder().build().unwrap();
3590 client.ticker_price(params).await
3591 });
3592
3593 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3594 let Message::Text(text) = sent else { panic!() };
3595 let v: Value = serde_json::from_str(&text).unwrap();
3596 let id = v["id"].as_str().unwrap().to_string();
3597
3598 let resp_json = json!({
3599 "id": id,
3600 "status": 400,
3601 "error": {
3602 "code": -2010,
3603 "msg": "Account has insufficient balance for requested action.",
3604 },
3605 "rateLimits": [
3606 {
3607 "rateLimitType": "ORDERS",
3608 "interval": "SECOND",
3609 "intervalNum": 10,
3610 "limit": 50,
3611 "count": 13
3612 },
3613 ],
3614 });
3615 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3616
3617 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3618 match join {
3619 Ok(Err(e)) => {
3620 let msg = e.to_string();
3621 assert!(
3622 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3623 "Expected error msg to contain server error, got: {msg}"
3624 );
3625 }
3626 Ok(Ok(_)) => panic!("Expected error"),
3627 Err(_) => panic!("Task panicked"),
3628 }
3629 });
3630 }
3631
3632 #[test]
3633 fn ticker_price_request_timeout() {
3634 TOKIO_SHARED_RT.block_on(async {
3635 let (ws_api, _conn, mut rx) = setup().await;
3636 let client = MarketApiClient::new(ws_api.clone());
3637
3638 let handle = spawn(async move {
3639 let params = TickerPriceParams::builder().build().unwrap();
3640 client.ticker_price(params).await
3641 });
3642
3643 let sent = timeout(Duration::from_secs(1), rx.recv())
3644 .await
3645 .expect("send should occur")
3646 .expect("channel closed");
3647 let Message::Text(text) = sent else {
3648 panic!("expected Message Text")
3649 };
3650
3651 let _: Value = serde_json::from_str(&text).unwrap();
3652
3653 let result = handle.await.expect("task completed");
3654 match result {
3655 Err(e) => {
3656 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3657 assert!(matches!(inner, WebsocketError::Timeout));
3658 } else {
3659 panic!("Unexpected error type: {:?}", e);
3660 }
3661 }
3662 Ok(_) => panic!("Expected timeout error"),
3663 }
3664 });
3665 }
3666
3667 #[test]
3668 fn ticker_trading_day_success() {
3669 TOKIO_SHARED_RT.block_on(async {
3670 let (ws_api, conn, mut rx) = setup().await;
3671 let client = MarketApiClient::new(ws_api.clone());
3672
3673 let handle = spawn(async move {
3674 let params = TickerTradingDayParams::builder().build().unwrap();
3675 client.ticker_trading_day(params).await
3676 });
3677
3678 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3679 let Message::Text(text) = sent else { panic!() };
3680 let v: Value = serde_json::from_str(&text).unwrap();
3681 let id = v["id"].as_str().unwrap();
3682 assert_eq!(v["method"], "/ticker.tradingDay".trim_start_matches('/'));
3683 let mut resp_json: Value = serde_json::from_str(r#"{"id":"f4b3b507-c8f2-442a-81a6-b2f12daa030f","status":200,"result":[{"symbol":"BTCUSDT","openTime":1695686400000,"closeTime":1695772799999,"firstId":3220151555,"lastId":3220849281,"count":697727}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3684 resp_json["id"] = id.into();
3685
3686 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3687 let expected_data: Vec<models::TickerTradingDayResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3688 let empty_array = Value::Array(vec![]);
3689 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3690 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3691 match raw_rate_limits.as_array() {
3692 Some(arr) if arr.is_empty() => None,
3693 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3694 None => None,
3695 };
3696
3697 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3698
3699 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3700
3701
3702 let response_rate_limits = response.rate_limits.clone();
3703 let response_data = response.data().expect("deserialize data");
3704
3705 assert_eq!(response_rate_limits, expected_rate_limits);
3706 assert_eq!(response_data, expected_data);
3707 });
3708 }
3709
3710 #[test]
3711 fn ticker_trading_day_error_response() {
3712 TOKIO_SHARED_RT.block_on(async {
3713 let (ws_api, conn, mut rx) = setup().await;
3714 let client = MarketApiClient::new(ws_api.clone());
3715
3716 let handle = tokio::spawn(async move {
3717 let params = TickerTradingDayParams::builder().build().unwrap();
3718 client.ticker_trading_day(params).await
3719 });
3720
3721 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3722 let Message::Text(text) = sent else { panic!() };
3723 let v: Value = serde_json::from_str(&text).unwrap();
3724 let id = v["id"].as_str().unwrap().to_string();
3725
3726 let resp_json = json!({
3727 "id": id,
3728 "status": 400,
3729 "error": {
3730 "code": -2010,
3731 "msg": "Account has insufficient balance for requested action.",
3732 },
3733 "rateLimits": [
3734 {
3735 "rateLimitType": "ORDERS",
3736 "interval": "SECOND",
3737 "intervalNum": 10,
3738 "limit": 50,
3739 "count": 13
3740 },
3741 ],
3742 });
3743 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3744
3745 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3746 match join {
3747 Ok(Err(e)) => {
3748 let msg = e.to_string();
3749 assert!(
3750 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3751 "Expected error msg to contain server error, got: {msg}"
3752 );
3753 }
3754 Ok(Ok(_)) => panic!("Expected error"),
3755 Err(_) => panic!("Task panicked"),
3756 }
3757 });
3758 }
3759
3760 #[test]
3761 fn ticker_trading_day_request_timeout() {
3762 TOKIO_SHARED_RT.block_on(async {
3763 let (ws_api, _conn, mut rx) = setup().await;
3764 let client = MarketApiClient::new(ws_api.clone());
3765
3766 let handle = spawn(async move {
3767 let params = TickerTradingDayParams::builder().build().unwrap();
3768 client.ticker_trading_day(params).await
3769 });
3770
3771 let sent = timeout(Duration::from_secs(1), rx.recv())
3772 .await
3773 .expect("send should occur")
3774 .expect("channel closed");
3775 let Message::Text(text) = sent else {
3776 panic!("expected Message Text")
3777 };
3778
3779 let _: Value = serde_json::from_str(&text).unwrap();
3780
3781 let result = handle.await.expect("task completed");
3782 match result {
3783 Err(e) => {
3784 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3785 assert!(matches!(inner, WebsocketError::Timeout));
3786 } else {
3787 panic!("Unexpected error type: {:?}", e);
3788 }
3789 }
3790 Ok(_) => panic!("Expected timeout error"),
3791 }
3792 });
3793 }
3794
3795 #[test]
3796 fn trades_aggregate_success() {
3797 TOKIO_SHARED_RT.block_on(async {
3798 let (ws_api, conn, mut rx) = setup().await;
3799 let client = MarketApiClient::new(ws_api.clone());
3800
3801 let handle = spawn(async move {
3802 let params = TradesAggregateParams::builder("BNBUSDT".to_string(),).build().unwrap();
3803 client.trades_aggregate(params).await
3804 });
3805
3806 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3807 let Message::Text(text) = sent else { panic!() };
3808 let v: Value = serde_json::from_str(&text).unwrap();
3809 let id = v["id"].as_str().unwrap();
3810 assert_eq!(v["method"], "/trades.aggregate".trim_start_matches('/'));
3811 let mut resp_json: Value = serde_json::from_str(r#"{"id":"189da436-d4bd-48ca-9f95-9f613d621717","status":200,"result":[{"a":50000000,"f":59120167,"l":59120170,"T":1565877971222,"m":true,"M":true}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3812 resp_json["id"] = id.into();
3813
3814 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3815 let expected_data: Vec<models::TradesAggregateResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3816 let empty_array = Value::Array(vec![]);
3817 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3818 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3819 match raw_rate_limits.as_array() {
3820 Some(arr) if arr.is_empty() => None,
3821 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3822 None => None,
3823 };
3824
3825 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3826
3827 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3828
3829
3830 let response_rate_limits = response.rate_limits.clone();
3831 let response_data = response.data().expect("deserialize data");
3832
3833 assert_eq!(response_rate_limits, expected_rate_limits);
3834 assert_eq!(response_data, expected_data);
3835 });
3836 }
3837
3838 #[test]
3839 fn trades_aggregate_error_response() {
3840 TOKIO_SHARED_RT.block_on(async {
3841 let (ws_api, conn, mut rx) = setup().await;
3842 let client = MarketApiClient::new(ws_api.clone());
3843
3844 let handle = tokio::spawn(async move {
3845 let params = TradesAggregateParams::builder("BNBUSDT".to_string(),).build().unwrap();
3846 client.trades_aggregate(params).await
3847 });
3848
3849 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3850 let Message::Text(text) = sent else { panic!() };
3851 let v: Value = serde_json::from_str(&text).unwrap();
3852 let id = v["id"].as_str().unwrap().to_string();
3853
3854 let resp_json = json!({
3855 "id": id,
3856 "status": 400,
3857 "error": {
3858 "code": -2010,
3859 "msg": "Account has insufficient balance for requested action.",
3860 },
3861 "rateLimits": [
3862 {
3863 "rateLimitType": "ORDERS",
3864 "interval": "SECOND",
3865 "intervalNum": 10,
3866 "limit": 50,
3867 "count": 13
3868 },
3869 ],
3870 });
3871 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3872
3873 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3874 match join {
3875 Ok(Err(e)) => {
3876 let msg = e.to_string();
3877 assert!(
3878 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3879 "Expected error msg to contain server error, got: {msg}"
3880 );
3881 }
3882 Ok(Ok(_)) => panic!("Expected error"),
3883 Err(_) => panic!("Task panicked"),
3884 }
3885 });
3886 }
3887
3888 #[test]
3889 fn trades_aggregate_request_timeout() {
3890 TOKIO_SHARED_RT.block_on(async {
3891 let (ws_api, _conn, mut rx) = setup().await;
3892 let client = MarketApiClient::new(ws_api.clone());
3893
3894 let handle = spawn(async move {
3895 let params = TradesAggregateParams::builder("BNBUSDT".to_string())
3896 .build()
3897 .unwrap();
3898 client.trades_aggregate(params).await
3899 });
3900
3901 let sent = timeout(Duration::from_secs(1), rx.recv())
3902 .await
3903 .expect("send should occur")
3904 .expect("channel closed");
3905 let Message::Text(text) = sent else {
3906 panic!("expected Message Text")
3907 };
3908
3909 let _: Value = serde_json::from_str(&text).unwrap();
3910
3911 let result = handle.await.expect("task completed");
3912 match result {
3913 Err(e) => {
3914 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3915 assert!(matches!(inner, WebsocketError::Timeout));
3916 } else {
3917 panic!("Unexpected error type: {:?}", e);
3918 }
3919 }
3920 Ok(_) => panic!("Expected timeout error"),
3921 }
3922 });
3923 }
3924
3925 #[test]
3926 fn trades_historical_success() {
3927 TOKIO_SHARED_RT.block_on(async {
3928 let (ws_api, conn, mut rx) = setup().await;
3929 let client = MarketApiClient::new(ws_api.clone());
3930
3931 let handle = spawn(async move {
3932 let params = TradesHistoricalParams::builder("BNBUSDT".to_string(),).build().unwrap();
3933 client.trades_historical(params).await
3934 });
3935
3936 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3937 let Message::Text(text) = sent else { panic!() };
3938 let v: Value = serde_json::from_str(&text).unwrap();
3939 let id = v["id"].as_str().unwrap();
3940 assert_eq!(v["method"], "/trades.historical".trim_start_matches('/'));
3941 let mut resp_json: Value = serde_json::from_str(r#"{"id":"cffc9c7d-4efc-4ce0-b587-6b87448f052a","status":200,"result":[{"id":0,"time":1500004800376,"isBuyerMaker":true,"isBestMatch":true}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3942 resp_json["id"] = id.into();
3943
3944 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3945 let expected_data: Vec<models::TradesHistoricalResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3946 let empty_array = Value::Array(vec![]);
3947 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3948 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3949 match raw_rate_limits.as_array() {
3950 Some(arr) if arr.is_empty() => None,
3951 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3952 None => None,
3953 };
3954
3955 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3956
3957 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3958
3959
3960 let response_rate_limits = response.rate_limits.clone();
3961 let response_data = response.data().expect("deserialize data");
3962
3963 assert_eq!(response_rate_limits, expected_rate_limits);
3964 assert_eq!(response_data, expected_data);
3965 });
3966 }
3967
3968 #[test]
3969 fn trades_historical_error_response() {
3970 TOKIO_SHARED_RT.block_on(async {
3971 let (ws_api, conn, mut rx) = setup().await;
3972 let client = MarketApiClient::new(ws_api.clone());
3973
3974 let handle = tokio::spawn(async move {
3975 let params = TradesHistoricalParams::builder("BNBUSDT".to_string(),).build().unwrap();
3976 client.trades_historical(params).await
3977 });
3978
3979 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3980 let Message::Text(text) = sent else { panic!() };
3981 let v: Value = serde_json::from_str(&text).unwrap();
3982 let id = v["id"].as_str().unwrap().to_string();
3983
3984 let resp_json = json!({
3985 "id": id,
3986 "status": 400,
3987 "error": {
3988 "code": -2010,
3989 "msg": "Account has insufficient balance for requested action.",
3990 },
3991 "rateLimits": [
3992 {
3993 "rateLimitType": "ORDERS",
3994 "interval": "SECOND",
3995 "intervalNum": 10,
3996 "limit": 50,
3997 "count": 13
3998 },
3999 ],
4000 });
4001 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
4002
4003 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
4004 match join {
4005 Ok(Err(e)) => {
4006 let msg = e.to_string();
4007 assert!(
4008 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
4009 "Expected error msg to contain server error, got: {msg}"
4010 );
4011 }
4012 Ok(Ok(_)) => panic!("Expected error"),
4013 Err(_) => panic!("Task panicked"),
4014 }
4015 });
4016 }
4017
4018 #[test]
4019 fn trades_historical_request_timeout() {
4020 TOKIO_SHARED_RT.block_on(async {
4021 let (ws_api, _conn, mut rx) = setup().await;
4022 let client = MarketApiClient::new(ws_api.clone());
4023
4024 let handle = spawn(async move {
4025 let params = TradesHistoricalParams::builder("BNBUSDT".to_string())
4026 .build()
4027 .unwrap();
4028 client.trades_historical(params).await
4029 });
4030
4031 let sent = timeout(Duration::from_secs(1), rx.recv())
4032 .await
4033 .expect("send should occur")
4034 .expect("channel closed");
4035 let Message::Text(text) = sent else {
4036 panic!("expected Message Text")
4037 };
4038
4039 let _: Value = serde_json::from_str(&text).unwrap();
4040
4041 let result = handle.await.expect("task completed");
4042 match result {
4043 Err(e) => {
4044 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
4045 assert!(matches!(inner, WebsocketError::Timeout));
4046 } else {
4047 panic!("Unexpected error type: {:?}", e);
4048 }
4049 }
4050 Ok(_) => panic!("Expected timeout error"),
4051 }
4052 });
4053 }
4054
4055 #[test]
4056 fn trades_recent_success() {
4057 TOKIO_SHARED_RT.block_on(async {
4058 let (ws_api, conn, mut rx) = setup().await;
4059 let client = MarketApiClient::new(ws_api.clone());
4060
4061 let handle = spawn(async move {
4062 let params = TradesRecentParams::builder("BNBUSDT".to_string(),).build().unwrap();
4063 client.trades_recent(params).await
4064 });
4065
4066 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
4067 let Message::Text(text) = sent else { panic!() };
4068 let v: Value = serde_json::from_str(&text).unwrap();
4069 let id = v["id"].as_str().unwrap();
4070 assert_eq!(v["method"], "/trades.recent".trim_start_matches('/'));
4071 let mut resp_json: Value = serde_json::from_str(r#"{"id":"409a20bd-253d-41db-a6dd-687862a5882f","status":200,"result":[{"id":194686783,"time":1660009530807,"isBuyerMaker":true,"isBestMatch":true}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
4072 resp_json["id"] = id.into();
4073
4074 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
4075 let expected_data: Vec<models::TradesRecentResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
4076 let empty_array = Value::Array(vec![]);
4077 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
4078 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
4079 match raw_rate_limits.as_array() {
4080 Some(arr) if arr.is_empty() => None,
4081 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
4082 None => None,
4083 };
4084
4085 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
4086
4087 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
4088
4089
4090 let response_rate_limits = response.rate_limits.clone();
4091 let response_data = response.data().expect("deserialize data");
4092
4093 assert_eq!(response_rate_limits, expected_rate_limits);
4094 assert_eq!(response_data, expected_data);
4095 });
4096 }
4097
4098 #[test]
4099 fn trades_recent_error_response() {
4100 TOKIO_SHARED_RT.block_on(async {
4101 let (ws_api, conn, mut rx) = setup().await;
4102 let client = MarketApiClient::new(ws_api.clone());
4103
4104 let handle = tokio::spawn(async move {
4105 let params = TradesRecentParams::builder("BNBUSDT".to_string(),).build().unwrap();
4106 client.trades_recent(params).await
4107 });
4108
4109 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
4110 let Message::Text(text) = sent else { panic!() };
4111 let v: Value = serde_json::from_str(&text).unwrap();
4112 let id = v["id"].as_str().unwrap().to_string();
4113
4114 let resp_json = json!({
4115 "id": id,
4116 "status": 400,
4117 "error": {
4118 "code": -2010,
4119 "msg": "Account has insufficient balance for requested action.",
4120 },
4121 "rateLimits": [
4122 {
4123 "rateLimitType": "ORDERS",
4124 "interval": "SECOND",
4125 "intervalNum": 10,
4126 "limit": 50,
4127 "count": 13
4128 },
4129 ],
4130 });
4131 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
4132
4133 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
4134 match join {
4135 Ok(Err(e)) => {
4136 let msg = e.to_string();
4137 assert!(
4138 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
4139 "Expected error msg to contain server error, got: {msg}"
4140 );
4141 }
4142 Ok(Ok(_)) => panic!("Expected error"),
4143 Err(_) => panic!("Task panicked"),
4144 }
4145 });
4146 }
4147
4148 #[test]
4149 fn trades_recent_request_timeout() {
4150 TOKIO_SHARED_RT.block_on(async {
4151 let (ws_api, _conn, mut rx) = setup().await;
4152 let client = MarketApiClient::new(ws_api.clone());
4153
4154 let handle = spawn(async move {
4155 let params = TradesRecentParams::builder("BNBUSDT".to_string())
4156 .build()
4157 .unwrap();
4158 client.trades_recent(params).await
4159 });
4160
4161 let sent = timeout(Duration::from_secs(1), rx.recv())
4162 .await
4163 .expect("send should occur")
4164 .expect("channel closed");
4165 let Message::Text(text) = sent else {
4166 panic!("expected Message Text")
4167 };
4168
4169 let _: Value = serde_json::from_str(&text).unwrap();
4170
4171 let result = handle.await.expect("task completed");
4172 match result {
4173 Err(e) => {
4174 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
4175 assert!(matches!(inner, WebsocketError::Timeout));
4176 } else {
4177 panic!("Unexpected error type: {:?}", e);
4178 }
4179 }
4180 Ok(_) => panic!("Expected timeout error"),
4181 }
4182 });
4183 }
4184
4185 #[test]
4186 fn ui_klines_success() {
4187 TOKIO_SHARED_RT.block_on(async {
4188 let (ws_api, conn, mut rx) = setup().await;
4189 let client = MarketApiClient::new(ws_api.clone());
4190
4191 let handle = spawn(async move {
4192 let params = UiKlinesParams::builder("BNBUSDT".to_string(),UiKlinesIntervalEnum::Interval1s,).build().unwrap();
4193 client.ui_klines(params).await
4194 });
4195
4196 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
4197 let Message::Text(text) = sent else { panic!() };
4198 let v: Value = serde_json::from_str(&text).unwrap();
4199 let id = v["id"].as_str().unwrap();
4200 assert_eq!(v["method"], "/uiKlines".trim_start_matches('/'));
4201 let mut resp_json: Value = serde_json::from_str(r#"{"id":"b137468a-fb20-4c06-bd6b-625148eec958","status":200,"result":[[1499040000000]],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
4202 resp_json["id"] = id.into();
4203
4204 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
4205 let expected_data: Vec<Vec<models::KlinesResponseResultInnerInner>> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
4206 let empty_array = Value::Array(vec![]);
4207 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
4208 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
4209 match raw_rate_limits.as_array() {
4210 Some(arr) if arr.is_empty() => None,
4211 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
4212 None => None,
4213 };
4214
4215 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
4216
4217 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
4218
4219
4220 let response_rate_limits = response.rate_limits.clone();
4221 let response_data = response.data().expect("deserialize data");
4222
4223 assert_eq!(response_rate_limits, expected_rate_limits);
4224 assert_eq!(response_data, expected_data);
4225 });
4226 }
4227
4228 #[test]
4229 fn ui_klines_error_response() {
4230 TOKIO_SHARED_RT.block_on(async {
4231 let (ws_api, conn, mut rx) = setup().await;
4232 let client = MarketApiClient::new(ws_api.clone());
4233
4234 let handle = tokio::spawn(async move {
4235 let params = UiKlinesParams::builder("BNBUSDT".to_string(),UiKlinesIntervalEnum::Interval1s,).build().unwrap();
4236 client.ui_klines(params).await
4237 });
4238
4239 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
4240 let Message::Text(text) = sent else { panic!() };
4241 let v: Value = serde_json::from_str(&text).unwrap();
4242 let id = v["id"].as_str().unwrap().to_string();
4243
4244 let resp_json = json!({
4245 "id": id,
4246 "status": 400,
4247 "error": {
4248 "code": -2010,
4249 "msg": "Account has insufficient balance for requested action.",
4250 },
4251 "rateLimits": [
4252 {
4253 "rateLimitType": "ORDERS",
4254 "interval": "SECOND",
4255 "intervalNum": 10,
4256 "limit": 50,
4257 "count": 13
4258 },
4259 ],
4260 });
4261 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
4262
4263 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
4264 match join {
4265 Ok(Err(e)) => {
4266 let msg = e.to_string();
4267 assert!(
4268 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
4269 "Expected error msg to contain server error, got: {msg}"
4270 );
4271 }
4272 Ok(Ok(_)) => panic!("Expected error"),
4273 Err(_) => panic!("Task panicked"),
4274 }
4275 });
4276 }
4277
4278 #[test]
4279 fn ui_klines_request_timeout() {
4280 TOKIO_SHARED_RT.block_on(async {
4281 let (ws_api, _conn, mut rx) = setup().await;
4282 let client = MarketApiClient::new(ws_api.clone());
4283
4284 let handle = spawn(async move {
4285 let params = UiKlinesParams::builder(
4286 "BNBUSDT".to_string(),
4287 UiKlinesIntervalEnum::Interval1s,
4288 )
4289 .build()
4290 .unwrap();
4291 client.ui_klines(params).await
4292 });
4293
4294 let sent = timeout(Duration::from_secs(1), rx.recv())
4295 .await
4296 .expect("send should occur")
4297 .expect("channel closed");
4298 let Message::Text(text) = sent else {
4299 panic!("expected Message Text")
4300 };
4301
4302 let _: Value = serde_json::from_str(&text).unwrap();
4303
4304 let result = handle.await.expect("task completed");
4305 match result {
4306 Err(e) => {
4307 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
4308 assert!(matches!(inner, WebsocketError::Timeout));
4309 } else {
4310 panic!("Unexpected error type: {:?}", e);
4311 }
4312 }
4313 Ok(_) => panic!("Expected timeout error"),
4314 }
4315 });
4316 }
4317}