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 AccountApi: Send + Sync {
33 async fn account_commission(
34 &self,
35 params: AccountCommissionParams,
36 ) -> anyhow::Result<WebsocketApiResponse<Box<models::AccountCommissionResponseResult>>>;
37 async fn account_rate_limits_orders(
38 &self,
39 params: AccountRateLimitsOrdersParams,
40 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::AccountRateLimitsOrdersResponseResultInner>>>;
41 async fn account_status(
42 &self,
43 params: AccountStatusParams,
44 ) -> anyhow::Result<WebsocketApiResponse<Box<models::AccountStatusResponseResult>>>;
45 async fn all_order_lists(
46 &self,
47 params: AllOrderListsParams,
48 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::AllOrderListsResponseResultInner>>>;
49 async fn all_orders(
50 &self,
51 params: AllOrdersParams,
52 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::AllOrdersResponseResultInner>>>;
53 async fn my_allocations(
54 &self,
55 params: MyAllocationsParams,
56 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::MyAllocationsResponseResultInner>>>;
57 async fn my_filters(
58 &self,
59 params: MyFiltersParams,
60 ) -> anyhow::Result<WebsocketApiResponse<models::MyFiltersResponse>>;
61 async fn my_prevented_matches(
62 &self,
63 params: MyPreventedMatchesParams,
64 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::MyPreventedMatchesResponseResultInner>>>;
65 async fn my_trades(
66 &self,
67 params: MyTradesParams,
68 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::MyTradesResponseResultInner>>>;
69 async fn open_order_lists_status(
70 &self,
71 params: OpenOrderListsStatusParams,
72 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::OpenOrderListsStatusResponseResultInner>>>;
73 async fn open_orders_status(
74 &self,
75 params: OpenOrdersStatusParams,
76 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::OpenOrdersStatusResponseResultInner>>>;
77 async fn order_amendments(
78 &self,
79 params: OrderAmendmentsParams,
80 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::OrderAmendmentsResponseResultInner>>>;
81 async fn order_list_status(
82 &self,
83 params: OrderListStatusParams,
84 ) -> anyhow::Result<WebsocketApiResponse<Box<models::OrderListStatusResponseResult>>>;
85 async fn order_status(
86 &self,
87 params: OrderStatusParams,
88 ) -> anyhow::Result<WebsocketApiResponse<Box<models::OrderStatusResponseResult>>>;
89}
90
91#[derive(Clone)]
92pub struct AccountApiClient {
93 websocket_api_base: Arc<WebsocketApi>,
94}
95
96impl AccountApiClient {
97 pub fn new(websocket_api_base: Arc<WebsocketApi>) -> Self {
98 Self { websocket_api_base }
99 }
100}
101
102#[derive(Clone, Debug, Builder, Deserialize)]
107#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
108pub struct AccountCommissionParams {
109 #[builder(setter(into))]
114 #[serde(rename = "symbol")]
115 pub symbol: String,
116 #[builder(setter(into), default)]
120 #[serde(rename = "id", default)]
121 pub id: Option<String>,
122}
123
124impl AccountCommissionParams {
125 #[must_use]
132 pub fn builder(symbol: String) -> AccountCommissionParamsBuilder {
133 AccountCommissionParamsBuilder::default().symbol(symbol)
134 }
135}
136#[derive(Clone, Debug, Builder, Deserialize, Default)]
141#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
142pub struct AccountRateLimitsOrdersParams {
143 #[builder(setter(into), default)]
147 #[serde(rename = "id", default)]
148 pub id: Option<String>,
149 #[builder(setter(into), default)]
153 #[serde(rename = "recvWindow", default)]
154 pub recv_window: Option<rust_decimal::Decimal>,
155}
156
157impl AccountRateLimitsOrdersParams {
158 #[must_use]
161 pub fn builder() -> AccountRateLimitsOrdersParamsBuilder {
162 AccountRateLimitsOrdersParamsBuilder::default()
163 }
164}
165#[derive(Clone, Debug, Builder, Deserialize, Default)]
170#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
171pub struct AccountStatusParams {
172 #[builder(setter(into), default)]
176 #[serde(rename = "id", default)]
177 pub id: Option<String>,
178 #[builder(setter(into), default)]
182 #[serde(rename = "omitZeroBalances", default)]
183 pub omit_zero_balances: Option<bool>,
184 #[builder(setter(into), default)]
188 #[serde(rename = "recvWindow", default)]
189 pub recv_window: Option<rust_decimal::Decimal>,
190}
191
192impl AccountStatusParams {
193 #[must_use]
196 pub fn builder() -> AccountStatusParamsBuilder {
197 AccountStatusParamsBuilder::default()
198 }
199}
200#[derive(Clone, Debug, Builder, Deserialize, Default)]
205#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
206pub struct AllOrderListsParams {
207 #[builder(setter(into), default)]
211 #[serde(rename = "id", default)]
212 pub id: Option<String>,
213 #[builder(setter(into), default)]
217 #[serde(rename = "fromId", default)]
218 pub from_id: Option<i32>,
219 #[builder(setter(into), default)]
223 #[serde(rename = "startTime", default)]
224 pub start_time: Option<i64>,
225 #[builder(setter(into), default)]
229 #[serde(rename = "endTime", default)]
230 pub end_time: Option<i64>,
231 #[builder(setter(into), default)]
235 #[serde(rename = "limit", default)]
236 pub limit: Option<i32>,
237 #[builder(setter(into), default)]
241 #[serde(rename = "recvWindow", default)]
242 pub recv_window: Option<rust_decimal::Decimal>,
243}
244
245impl AllOrderListsParams {
246 #[must_use]
249 pub fn builder() -> AllOrderListsParamsBuilder {
250 AllOrderListsParamsBuilder::default()
251 }
252}
253#[derive(Clone, Debug, Builder, Deserialize)]
258#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
259pub struct AllOrdersParams {
260 #[builder(setter(into))]
265 #[serde(rename = "symbol")]
266 pub symbol: String,
267 #[builder(setter(into), default)]
271 #[serde(rename = "id", default)]
272 pub id: Option<String>,
273 #[builder(setter(into), default)]
277 #[serde(rename = "orderId", default)]
278 pub order_id: Option<i64>,
279 #[builder(setter(into), default)]
283 #[serde(rename = "startTime", default)]
284 pub start_time: Option<i64>,
285 #[builder(setter(into), default)]
289 #[serde(rename = "endTime", default)]
290 pub end_time: Option<i64>,
291 #[builder(setter(into), default)]
295 #[serde(rename = "limit", default)]
296 pub limit: Option<i32>,
297 #[builder(setter(into), default)]
301 #[serde(rename = "recvWindow", default)]
302 pub recv_window: Option<rust_decimal::Decimal>,
303}
304
305impl AllOrdersParams {
306 #[must_use]
313 pub fn builder(symbol: String) -> AllOrdersParamsBuilder {
314 AllOrdersParamsBuilder::default().symbol(symbol)
315 }
316}
317#[derive(Clone, Debug, Builder, Deserialize)]
322#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
323pub struct MyAllocationsParams {
324 #[builder(setter(into))]
329 #[serde(rename = "symbol")]
330 pub symbol: String,
331 #[builder(setter(into), default)]
335 #[serde(rename = "id", default)]
336 pub id: Option<String>,
337 #[builder(setter(into), default)]
341 #[serde(rename = "startTime", default)]
342 pub start_time: Option<i64>,
343 #[builder(setter(into), default)]
347 #[serde(rename = "endTime", default)]
348 pub end_time: Option<i64>,
349 #[builder(setter(into), default)]
353 #[serde(rename = "fromAllocationId", default)]
354 pub from_allocation_id: Option<i32>,
355 #[builder(setter(into), default)]
359 #[serde(rename = "limit", default)]
360 pub limit: Option<i32>,
361 #[builder(setter(into), default)]
365 #[serde(rename = "orderId", default)]
366 pub order_id: Option<i64>,
367 #[builder(setter(into), default)]
371 #[serde(rename = "recvWindow", default)]
372 pub recv_window: Option<rust_decimal::Decimal>,
373}
374
375impl MyAllocationsParams {
376 #[must_use]
383 pub fn builder(symbol: String) -> MyAllocationsParamsBuilder {
384 MyAllocationsParamsBuilder::default().symbol(symbol)
385 }
386}
387#[derive(Clone, Debug, Builder, Deserialize)]
392#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
393pub struct MyFiltersParams {
394 #[builder(setter(into))]
399 #[serde(rename = "symbol")]
400 pub symbol: String,
401 #[builder(setter(into), default)]
405 #[serde(rename = "id", default)]
406 pub id: Option<String>,
407 #[builder(setter(into), default)]
411 #[serde(rename = "recvWindow", default)]
412 pub recv_window: Option<rust_decimal::Decimal>,
413}
414
415impl MyFiltersParams {
416 #[must_use]
423 pub fn builder(symbol: String) -> MyFiltersParamsBuilder {
424 MyFiltersParamsBuilder::default().symbol(symbol)
425 }
426}
427#[derive(Clone, Debug, Builder, Deserialize)]
432#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
433pub struct MyPreventedMatchesParams {
434 #[builder(setter(into))]
439 #[serde(rename = "symbol")]
440 pub symbol: String,
441 #[builder(setter(into), default)]
445 #[serde(rename = "id", default)]
446 pub id: Option<String>,
447 #[builder(setter(into), default)]
451 #[serde(rename = "preventedMatchId", default)]
452 pub prevented_match_id: Option<i64>,
453 #[builder(setter(into), default)]
457 #[serde(rename = "orderId", default)]
458 pub order_id: Option<i64>,
459 #[builder(setter(into), default)]
463 #[serde(rename = "fromPreventedMatchId", default)]
464 pub from_prevented_match_id: Option<i64>,
465 #[builder(setter(into), default)]
469 #[serde(rename = "limit", default)]
470 pub limit: Option<i32>,
471 #[builder(setter(into), default)]
475 #[serde(rename = "recvWindow", default)]
476 pub recv_window: Option<rust_decimal::Decimal>,
477}
478
479impl MyPreventedMatchesParams {
480 #[must_use]
487 pub fn builder(symbol: String) -> MyPreventedMatchesParamsBuilder {
488 MyPreventedMatchesParamsBuilder::default().symbol(symbol)
489 }
490}
491#[derive(Clone, Debug, Builder, Deserialize)]
496#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
497pub struct MyTradesParams {
498 #[builder(setter(into))]
503 #[serde(rename = "symbol")]
504 pub symbol: String,
505 #[builder(setter(into), default)]
509 #[serde(rename = "id", default)]
510 pub id: Option<String>,
511 #[builder(setter(into), default)]
515 #[serde(rename = "orderId", default)]
516 pub order_id: Option<i64>,
517 #[builder(setter(into), default)]
521 #[serde(rename = "startTime", default)]
522 pub start_time: Option<i64>,
523 #[builder(setter(into), default)]
527 #[serde(rename = "endTime", default)]
528 pub end_time: Option<i64>,
529 #[builder(setter(into), default)]
533 #[serde(rename = "fromId", default)]
534 pub from_id: Option<i32>,
535 #[builder(setter(into), default)]
539 #[serde(rename = "limit", default)]
540 pub limit: Option<i32>,
541 #[builder(setter(into), default)]
545 #[serde(rename = "recvWindow", default)]
546 pub recv_window: Option<rust_decimal::Decimal>,
547}
548
549impl MyTradesParams {
550 #[must_use]
557 pub fn builder(symbol: String) -> MyTradesParamsBuilder {
558 MyTradesParamsBuilder::default().symbol(symbol)
559 }
560}
561#[derive(Clone, Debug, Builder, Deserialize, Default)]
566#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
567pub struct OpenOrderListsStatusParams {
568 #[builder(setter(into), default)]
572 #[serde(rename = "id", default)]
573 pub id: Option<String>,
574 #[builder(setter(into), default)]
578 #[serde(rename = "recvWindow", default)]
579 pub recv_window: Option<rust_decimal::Decimal>,
580}
581
582impl OpenOrderListsStatusParams {
583 #[must_use]
586 pub fn builder() -> OpenOrderListsStatusParamsBuilder {
587 OpenOrderListsStatusParamsBuilder::default()
588 }
589}
590#[derive(Clone, Debug, Builder, Deserialize, Default)]
595#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
596pub struct OpenOrdersStatusParams {
597 #[builder(setter(into), default)]
601 #[serde(rename = "id", default)]
602 pub id: Option<String>,
603 #[builder(setter(into), default)]
607 #[serde(rename = "symbol", default)]
608 pub symbol: Option<String>,
609 #[builder(setter(into), default)]
613 #[serde(rename = "recvWindow", default)]
614 pub recv_window: Option<rust_decimal::Decimal>,
615}
616
617impl OpenOrdersStatusParams {
618 #[must_use]
621 pub fn builder() -> OpenOrdersStatusParamsBuilder {
622 OpenOrdersStatusParamsBuilder::default()
623 }
624}
625#[derive(Clone, Debug, Builder, Deserialize)]
630#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
631pub struct OrderAmendmentsParams {
632 #[builder(setter(into))]
637 #[serde(rename = "symbol")]
638 pub symbol: String,
639 #[builder(setter(into))]
643 #[serde(rename = "orderId")]
644 pub order_id: i64,
645 #[builder(setter(into), default)]
649 #[serde(rename = "id", default)]
650 pub id: Option<String>,
651 #[builder(setter(into), default)]
655 #[serde(rename = "fromExecutionId", default)]
656 pub from_execution_id: Option<i64>,
657 #[builder(setter(into), default)]
661 #[serde(rename = "limit", default)]
662 pub limit: Option<i64>,
663 #[builder(setter(into), default)]
667 #[serde(rename = "recvWindow", default)]
668 pub recv_window: Option<rust_decimal::Decimal>,
669}
670
671impl OrderAmendmentsParams {
672 #[must_use]
680 pub fn builder(symbol: String, order_id: i64) -> OrderAmendmentsParamsBuilder {
681 OrderAmendmentsParamsBuilder::default()
682 .symbol(symbol)
683 .order_id(order_id)
684 }
685}
686#[derive(Clone, Debug, Builder, Deserialize, Default)]
691#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
692pub struct OrderListStatusParams {
693 #[builder(setter(into), default)]
697 #[serde(rename = "id", default)]
698 pub id: Option<String>,
699 #[builder(setter(into), default)]
703 #[serde(rename = "origClientOrderId", default)]
704 pub orig_client_order_id: Option<String>,
705 #[builder(setter(into), default)]
709 #[serde(rename = "orderListId", default)]
710 pub order_list_id: Option<i32>,
711 #[builder(setter(into), default)]
715 #[serde(rename = "recvWindow", default)]
716 pub recv_window: Option<rust_decimal::Decimal>,
717}
718
719impl OrderListStatusParams {
720 #[must_use]
723 pub fn builder() -> OrderListStatusParamsBuilder {
724 OrderListStatusParamsBuilder::default()
725 }
726}
727#[derive(Clone, Debug, Builder, Deserialize)]
732#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
733pub struct OrderStatusParams {
734 #[builder(setter(into))]
739 #[serde(rename = "symbol")]
740 pub symbol: String,
741 #[builder(setter(into), default)]
745 #[serde(rename = "id", default)]
746 pub id: Option<String>,
747 #[builder(setter(into), default)]
751 #[serde(rename = "orderId", default)]
752 pub order_id: Option<i64>,
753 #[builder(setter(into), default)]
757 #[serde(rename = "origClientOrderId", default)]
758 pub orig_client_order_id: Option<String>,
759 #[builder(setter(into), default)]
763 #[serde(rename = "recvWindow", default)]
764 pub recv_window: Option<rust_decimal::Decimal>,
765}
766
767impl OrderStatusParams {
768 #[must_use]
775 pub fn builder(symbol: String) -> OrderStatusParamsBuilder {
776 OrderStatusParamsBuilder::default().symbol(symbol)
777 }
778}
779
780#[async_trait]
781impl AccountApi for AccountApiClient {
782 async fn account_commission(
783 &self,
784 params: AccountCommissionParams,
785 ) -> anyhow::Result<WebsocketApiResponse<Box<models::AccountCommissionResponseResult>>> {
786 let AccountCommissionParams { symbol, id } = params;
787
788 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
789 payload.insert("symbol".to_string(), serde_json::json!(symbol));
790 if let Some(value) = id {
791 payload.insert("id".to_string(), serde_json::json!(value));
792 }
793 let payload = remove_empty_value(payload);
794
795 self.websocket_api_base
796 .send_message::<Box<models::AccountCommissionResponseResult>>(
797 "/account.commission".trim_start_matches('/'),
798 payload,
799 WebsocketMessageSendOptions::new().signed(),
800 )
801 .await
802 .map_err(anyhow::Error::from)?
803 .into_iter()
804 .next()
805 .ok_or(WebsocketError::NoResponse)
806 .map_err(anyhow::Error::from)
807 }
808
809 async fn account_rate_limits_orders(
810 &self,
811 params: AccountRateLimitsOrdersParams,
812 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::AccountRateLimitsOrdersResponseResultInner>>>
813 {
814 let AccountRateLimitsOrdersParams { id, recv_window } = params;
815
816 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
817 if let Some(value) = id {
818 payload.insert("id".to_string(), serde_json::json!(value));
819 }
820 if let Some(value) = recv_window {
821 payload.insert("recvWindow".to_string(), serde_json::json!(value));
822 }
823 let payload = remove_empty_value(payload);
824
825 self.websocket_api_base
826 .send_message::<Vec<models::AccountRateLimitsOrdersResponseResultInner>>(
827 "/account.rateLimits.orders".trim_start_matches('/'),
828 payload,
829 WebsocketMessageSendOptions::new().signed(),
830 )
831 .await
832 .map_err(anyhow::Error::from)?
833 .into_iter()
834 .next()
835 .ok_or(WebsocketError::NoResponse)
836 .map_err(anyhow::Error::from)
837 }
838
839 async fn account_status(
840 &self,
841 params: AccountStatusParams,
842 ) -> anyhow::Result<WebsocketApiResponse<Box<models::AccountStatusResponseResult>>> {
843 let AccountStatusParams {
844 id,
845 omit_zero_balances,
846 recv_window,
847 } = params;
848
849 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
850 if let Some(value) = id {
851 payload.insert("id".to_string(), serde_json::json!(value));
852 }
853 if let Some(value) = omit_zero_balances {
854 payload.insert("omitZeroBalances".to_string(), serde_json::json!(value));
855 }
856 if let Some(value) = recv_window {
857 payload.insert("recvWindow".to_string(), serde_json::json!(value));
858 }
859 let payload = remove_empty_value(payload);
860
861 self.websocket_api_base
862 .send_message::<Box<models::AccountStatusResponseResult>>(
863 "/account.status".trim_start_matches('/'),
864 payload,
865 WebsocketMessageSendOptions::new().signed(),
866 )
867 .await
868 .map_err(anyhow::Error::from)?
869 .into_iter()
870 .next()
871 .ok_or(WebsocketError::NoResponse)
872 .map_err(anyhow::Error::from)
873 }
874
875 async fn all_order_lists(
876 &self,
877 params: AllOrderListsParams,
878 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::AllOrderListsResponseResultInner>>> {
879 let AllOrderListsParams {
880 id,
881 from_id,
882 start_time,
883 end_time,
884 limit,
885 recv_window,
886 } = params;
887
888 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
889 if let Some(value) = id {
890 payload.insert("id".to_string(), serde_json::json!(value));
891 }
892 if let Some(value) = from_id {
893 payload.insert("fromId".to_string(), serde_json::json!(value));
894 }
895 if let Some(value) = start_time {
896 payload.insert("startTime".to_string(), serde_json::json!(value));
897 }
898 if let Some(value) = end_time {
899 payload.insert("endTime".to_string(), serde_json::json!(value));
900 }
901 if let Some(value) = limit {
902 payload.insert("limit".to_string(), serde_json::json!(value));
903 }
904 if let Some(value) = recv_window {
905 payload.insert("recvWindow".to_string(), serde_json::json!(value));
906 }
907 let payload = remove_empty_value(payload);
908
909 self.websocket_api_base
910 .send_message::<Vec<models::AllOrderListsResponseResultInner>>(
911 "/allOrderLists".trim_start_matches('/'),
912 payload,
913 WebsocketMessageSendOptions::new().signed(),
914 )
915 .await
916 .map_err(anyhow::Error::from)?
917 .into_iter()
918 .next()
919 .ok_or(WebsocketError::NoResponse)
920 .map_err(anyhow::Error::from)
921 }
922
923 async fn all_orders(
924 &self,
925 params: AllOrdersParams,
926 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::AllOrdersResponseResultInner>>> {
927 let AllOrdersParams {
928 symbol,
929 id,
930 order_id,
931 start_time,
932 end_time,
933 limit,
934 recv_window,
935 } = params;
936
937 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
938 payload.insert("symbol".to_string(), serde_json::json!(symbol));
939 if let Some(value) = id {
940 payload.insert("id".to_string(), serde_json::json!(value));
941 }
942 if let Some(value) = order_id {
943 payload.insert("orderId".to_string(), serde_json::json!(value));
944 }
945 if let Some(value) = start_time {
946 payload.insert("startTime".to_string(), serde_json::json!(value));
947 }
948 if let Some(value) = end_time {
949 payload.insert("endTime".to_string(), serde_json::json!(value));
950 }
951 if let Some(value) = limit {
952 payload.insert("limit".to_string(), serde_json::json!(value));
953 }
954 if let Some(value) = recv_window {
955 payload.insert("recvWindow".to_string(), serde_json::json!(value));
956 }
957 let payload = remove_empty_value(payload);
958
959 self.websocket_api_base
960 .send_message::<Vec<models::AllOrdersResponseResultInner>>(
961 "/allOrders".trim_start_matches('/'),
962 payload,
963 WebsocketMessageSendOptions::new().signed(),
964 )
965 .await
966 .map_err(anyhow::Error::from)?
967 .into_iter()
968 .next()
969 .ok_or(WebsocketError::NoResponse)
970 .map_err(anyhow::Error::from)
971 }
972
973 async fn my_allocations(
974 &self,
975 params: MyAllocationsParams,
976 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::MyAllocationsResponseResultInner>>> {
977 let MyAllocationsParams {
978 symbol,
979 id,
980 start_time,
981 end_time,
982 from_allocation_id,
983 limit,
984 order_id,
985 recv_window,
986 } = params;
987
988 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
989 payload.insert("symbol".to_string(), serde_json::json!(symbol));
990 if let Some(value) = id {
991 payload.insert("id".to_string(), serde_json::json!(value));
992 }
993 if let Some(value) = start_time {
994 payload.insert("startTime".to_string(), serde_json::json!(value));
995 }
996 if let Some(value) = end_time {
997 payload.insert("endTime".to_string(), serde_json::json!(value));
998 }
999 if let Some(value) = from_allocation_id {
1000 payload.insert("fromAllocationId".to_string(), serde_json::json!(value));
1001 }
1002 if let Some(value) = limit {
1003 payload.insert("limit".to_string(), serde_json::json!(value));
1004 }
1005 if let Some(value) = order_id {
1006 payload.insert("orderId".to_string(), serde_json::json!(value));
1007 }
1008 if let Some(value) = recv_window {
1009 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1010 }
1011 let payload = remove_empty_value(payload);
1012
1013 self.websocket_api_base
1014 .send_message::<Vec<models::MyAllocationsResponseResultInner>>(
1015 "/myAllocations".trim_start_matches('/'),
1016 payload,
1017 WebsocketMessageSendOptions::new().signed(),
1018 )
1019 .await
1020 .map_err(anyhow::Error::from)?
1021 .into_iter()
1022 .next()
1023 .ok_or(WebsocketError::NoResponse)
1024 .map_err(anyhow::Error::from)
1025 }
1026
1027 async fn my_filters(
1028 &self,
1029 params: MyFiltersParams,
1030 ) -> anyhow::Result<WebsocketApiResponse<models::MyFiltersResponse>> {
1031 let MyFiltersParams {
1032 symbol,
1033 id,
1034 recv_window,
1035 } = params;
1036
1037 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1038 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1039 if let Some(value) = id {
1040 payload.insert("id".to_string(), serde_json::json!(value));
1041 }
1042 if let Some(value) = recv_window {
1043 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1044 }
1045 let payload = remove_empty_value(payload);
1046
1047 self.websocket_api_base
1048 .send_message::<models::MyFiltersResponse>(
1049 "/myFilters".trim_start_matches('/'),
1050 payload,
1051 WebsocketMessageSendOptions::new().signed(),
1052 )
1053 .await
1054 .map_err(anyhow::Error::from)?
1055 .into_iter()
1056 .next()
1057 .ok_or(WebsocketError::NoResponse)
1058 .map_err(anyhow::Error::from)
1059 }
1060
1061 async fn my_prevented_matches(
1062 &self,
1063 params: MyPreventedMatchesParams,
1064 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::MyPreventedMatchesResponseResultInner>>>
1065 {
1066 let MyPreventedMatchesParams {
1067 symbol,
1068 id,
1069 prevented_match_id,
1070 order_id,
1071 from_prevented_match_id,
1072 limit,
1073 recv_window,
1074 } = params;
1075
1076 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1077 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1078 if let Some(value) = id {
1079 payload.insert("id".to_string(), serde_json::json!(value));
1080 }
1081 if let Some(value) = prevented_match_id {
1082 payload.insert("preventedMatchId".to_string(), serde_json::json!(value));
1083 }
1084 if let Some(value) = order_id {
1085 payload.insert("orderId".to_string(), serde_json::json!(value));
1086 }
1087 if let Some(value) = from_prevented_match_id {
1088 payload.insert("fromPreventedMatchId".to_string(), serde_json::json!(value));
1089 }
1090 if let Some(value) = limit {
1091 payload.insert("limit".to_string(), serde_json::json!(value));
1092 }
1093 if let Some(value) = recv_window {
1094 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1095 }
1096 let payload = remove_empty_value(payload);
1097
1098 self.websocket_api_base
1099 .send_message::<Vec<models::MyPreventedMatchesResponseResultInner>>(
1100 "/myPreventedMatches".trim_start_matches('/'),
1101 payload,
1102 WebsocketMessageSendOptions::new().signed(),
1103 )
1104 .await
1105 .map_err(anyhow::Error::from)?
1106 .into_iter()
1107 .next()
1108 .ok_or(WebsocketError::NoResponse)
1109 .map_err(anyhow::Error::from)
1110 }
1111
1112 async fn my_trades(
1113 &self,
1114 params: MyTradesParams,
1115 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::MyTradesResponseResultInner>>> {
1116 let MyTradesParams {
1117 symbol,
1118 id,
1119 order_id,
1120 start_time,
1121 end_time,
1122 from_id,
1123 limit,
1124 recv_window,
1125 } = params;
1126
1127 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1128 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1129 if let Some(value) = id {
1130 payload.insert("id".to_string(), serde_json::json!(value));
1131 }
1132 if let Some(value) = order_id {
1133 payload.insert("orderId".to_string(), serde_json::json!(value));
1134 }
1135 if let Some(value) = start_time {
1136 payload.insert("startTime".to_string(), serde_json::json!(value));
1137 }
1138 if let Some(value) = end_time {
1139 payload.insert("endTime".to_string(), serde_json::json!(value));
1140 }
1141 if let Some(value) = from_id {
1142 payload.insert("fromId".to_string(), serde_json::json!(value));
1143 }
1144 if let Some(value) = limit {
1145 payload.insert("limit".to_string(), serde_json::json!(value));
1146 }
1147 if let Some(value) = recv_window {
1148 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1149 }
1150 let payload = remove_empty_value(payload);
1151
1152 self.websocket_api_base
1153 .send_message::<Vec<models::MyTradesResponseResultInner>>(
1154 "/myTrades".trim_start_matches('/'),
1155 payload,
1156 WebsocketMessageSendOptions::new().signed(),
1157 )
1158 .await
1159 .map_err(anyhow::Error::from)?
1160 .into_iter()
1161 .next()
1162 .ok_or(WebsocketError::NoResponse)
1163 .map_err(anyhow::Error::from)
1164 }
1165
1166 async fn open_order_lists_status(
1167 &self,
1168 params: OpenOrderListsStatusParams,
1169 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::OpenOrderListsStatusResponseResultInner>>>
1170 {
1171 let OpenOrderListsStatusParams { id, recv_window } = params;
1172
1173 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1174 if let Some(value) = id {
1175 payload.insert("id".to_string(), serde_json::json!(value));
1176 }
1177 if let Some(value) = recv_window {
1178 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1179 }
1180 let payload = remove_empty_value(payload);
1181
1182 self.websocket_api_base
1183 .send_message::<Vec<models::OpenOrderListsStatusResponseResultInner>>(
1184 "/openOrderLists.status".trim_start_matches('/'),
1185 payload,
1186 WebsocketMessageSendOptions::new().signed(),
1187 )
1188 .await
1189 .map_err(anyhow::Error::from)?
1190 .into_iter()
1191 .next()
1192 .ok_or(WebsocketError::NoResponse)
1193 .map_err(anyhow::Error::from)
1194 }
1195
1196 async fn open_orders_status(
1197 &self,
1198 params: OpenOrdersStatusParams,
1199 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::OpenOrdersStatusResponseResultInner>>>
1200 {
1201 let OpenOrdersStatusParams {
1202 id,
1203 symbol,
1204 recv_window,
1205 } = params;
1206
1207 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1208 if let Some(value) = id {
1209 payload.insert("id".to_string(), serde_json::json!(value));
1210 }
1211 if let Some(value) = symbol {
1212 payload.insert("symbol".to_string(), serde_json::json!(value));
1213 }
1214 if let Some(value) = recv_window {
1215 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1216 }
1217 let payload = remove_empty_value(payload);
1218
1219 self.websocket_api_base
1220 .send_message::<Vec<models::OpenOrdersStatusResponseResultInner>>(
1221 "/openOrders.status".trim_start_matches('/'),
1222 payload,
1223 WebsocketMessageSendOptions::new().signed(),
1224 )
1225 .await
1226 .map_err(anyhow::Error::from)?
1227 .into_iter()
1228 .next()
1229 .ok_or(WebsocketError::NoResponse)
1230 .map_err(anyhow::Error::from)
1231 }
1232
1233 async fn order_amendments(
1234 &self,
1235 params: OrderAmendmentsParams,
1236 ) -> anyhow::Result<WebsocketApiResponse<Vec<models::OrderAmendmentsResponseResultInner>>> {
1237 let OrderAmendmentsParams {
1238 symbol,
1239 order_id,
1240 id,
1241 from_execution_id,
1242 limit,
1243 recv_window,
1244 } = params;
1245
1246 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1247 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1248 payload.insert("orderId".to_string(), serde_json::json!(order_id));
1249 if let Some(value) = id {
1250 payload.insert("id".to_string(), serde_json::json!(value));
1251 }
1252 if let Some(value) = from_execution_id {
1253 payload.insert("fromExecutionId".to_string(), serde_json::json!(value));
1254 }
1255 if let Some(value) = limit {
1256 payload.insert("limit".to_string(), serde_json::json!(value));
1257 }
1258 if let Some(value) = recv_window {
1259 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1260 }
1261 let payload = remove_empty_value(payload);
1262
1263 self.websocket_api_base
1264 .send_message::<Vec<models::OrderAmendmentsResponseResultInner>>(
1265 "/order.amendments".trim_start_matches('/'),
1266 payload,
1267 WebsocketMessageSendOptions::new().signed(),
1268 )
1269 .await
1270 .map_err(anyhow::Error::from)?
1271 .into_iter()
1272 .next()
1273 .ok_or(WebsocketError::NoResponse)
1274 .map_err(anyhow::Error::from)
1275 }
1276
1277 async fn order_list_status(
1278 &self,
1279 params: OrderListStatusParams,
1280 ) -> anyhow::Result<WebsocketApiResponse<Box<models::OrderListStatusResponseResult>>> {
1281 let OrderListStatusParams {
1282 id,
1283 orig_client_order_id,
1284 order_list_id,
1285 recv_window,
1286 } = params;
1287
1288 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1289 if let Some(value) = id {
1290 payload.insert("id".to_string(), serde_json::json!(value));
1291 }
1292 if let Some(value) = orig_client_order_id {
1293 payload.insert("origClientOrderId".to_string(), serde_json::json!(value));
1294 }
1295 if let Some(value) = order_list_id {
1296 payload.insert("orderListId".to_string(), serde_json::json!(value));
1297 }
1298 if let Some(value) = recv_window {
1299 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1300 }
1301 let payload = remove_empty_value(payload);
1302
1303 self.websocket_api_base
1304 .send_message::<Box<models::OrderListStatusResponseResult>>(
1305 "/orderList.status".trim_start_matches('/'),
1306 payload,
1307 WebsocketMessageSendOptions::new().signed(),
1308 )
1309 .await
1310 .map_err(anyhow::Error::from)?
1311 .into_iter()
1312 .next()
1313 .ok_or(WebsocketError::NoResponse)
1314 .map_err(anyhow::Error::from)
1315 }
1316
1317 async fn order_status(
1318 &self,
1319 params: OrderStatusParams,
1320 ) -> anyhow::Result<WebsocketApiResponse<Box<models::OrderStatusResponseResult>>> {
1321 let OrderStatusParams {
1322 symbol,
1323 id,
1324 order_id,
1325 orig_client_order_id,
1326 recv_window,
1327 } = params;
1328
1329 let mut payload: BTreeMap<String, Value> = BTreeMap::new();
1330 payload.insert("symbol".to_string(), serde_json::json!(symbol));
1331 if let Some(value) = id {
1332 payload.insert("id".to_string(), serde_json::json!(value));
1333 }
1334 if let Some(value) = order_id {
1335 payload.insert("orderId".to_string(), serde_json::json!(value));
1336 }
1337 if let Some(value) = orig_client_order_id {
1338 payload.insert("origClientOrderId".to_string(), serde_json::json!(value));
1339 }
1340 if let Some(value) = recv_window {
1341 payload.insert("recvWindow".to_string(), serde_json::json!(value));
1342 }
1343 let payload = remove_empty_value(payload);
1344
1345 self.websocket_api_base
1346 .send_message::<Box<models::OrderStatusResponseResult>>(
1347 "/order.status".trim_start_matches('/'),
1348 payload,
1349 WebsocketMessageSendOptions::new().signed(),
1350 )
1351 .await
1352 .map_err(anyhow::Error::from)?
1353 .into_iter()
1354 .next()
1355 .ok_or(WebsocketError::NoResponse)
1356 .map_err(anyhow::Error::from)
1357 }
1358}
1359
1360#[cfg(all(test, feature = "spot"))]
1361mod tests {
1362 use super::*;
1363 use crate::TOKIO_SHARED_RT;
1364 use crate::common::websocket::{WebsocketApi, WebsocketConnection, WebsocketHandler};
1365 use crate::config::ConfigurationWebsocketApi;
1366 use crate::errors::WebsocketError;
1367 use crate::models::WebsocketApiRateLimit;
1368 use serde_json::{Value, json};
1369 use tokio::spawn;
1370 use tokio::sync::mpsc::{UnboundedReceiver, unbounded_channel};
1371 use tokio::time::{Duration, timeout};
1372 use tokio_tungstenite::tungstenite::Message;
1373
1374 async fn setup() -> (
1375 Arc<WebsocketApi>,
1376 Arc<WebsocketConnection>,
1377 UnboundedReceiver<Message>,
1378 ) {
1379 let conn = WebsocketConnection::new("test-conn");
1380 let (tx, rx) = unbounded_channel::<Message>();
1381 {
1382 let mut conn_state = conn.state.lock().await;
1383 conn_state.ws_write_tx = Some(tx);
1384 }
1385
1386 let config = ConfigurationWebsocketApi::builder()
1387 .api_key("key")
1388 .api_secret("secret")
1389 .build()
1390 .expect("Failed to build configuration");
1391 let ws_api = WebsocketApi::new(config, vec![conn.clone()]);
1392 conn.set_handler(ws_api.clone() as Arc<dyn WebsocketHandler>)
1393 .await;
1394 ws_api.clone().connect().await.unwrap();
1395
1396 (ws_api, conn, rx)
1397 }
1398
1399 #[test]
1400 fn account_commission_success() {
1401 TOKIO_SHARED_RT.block_on(async {
1402 let (ws_api, conn, mut rx) = setup().await;
1403 let client = AccountApiClient::new(ws_api.clone());
1404
1405 let handle = spawn(async move {
1406 let params = AccountCommissionParams::builder("BNBUSDT".to_string(),).build().unwrap();
1407 client.account_commission(params).await
1408 });
1409
1410 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
1411 let Message::Text(text) = sent else { panic!() };
1412 let v: Value = serde_json::from_str(&text).unwrap();
1413 let id = v["id"].as_str().unwrap();
1414 assert_eq!(v["method"], "/account.commission".trim_start_matches('/'));
1415 let mut resp_json: Value = serde_json::from_str(r#"{"id":"d3df8a61-98ea-4fe0-8f4e-0fcea5d418b0","status":200,"result":{"symbol":"BTCUSDT","standardCommission":{"maker":"0.00000010","taker":"0.00000020","buyer":"0.00000030","seller":"0.00000040"},"specialCommission":{"maker":"0.01000000","taker":"0.02000000","buyer":"0.03000000","seller":"0.04000000"},"discount":{"enabledForAccount":true,"enabledForSymbol":true,"discountAsset":"BNB","discount":"0.75000000"}},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
1416 resp_json["id"] = id.into();
1417
1418 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
1419 let expected_data: Box<models::AccountCommissionResponseResult> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
1420 let empty_array = Value::Array(vec![]);
1421 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
1422 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
1423 match raw_rate_limits.as_array() {
1424 Some(arr) if arr.is_empty() => None,
1425 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
1426 None => None,
1427 };
1428
1429 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1430
1431 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
1432
1433
1434 let response_rate_limits = response.rate_limits.clone();
1435 let response_data = response.data().expect("deserialize data");
1436
1437 assert_eq!(response_rate_limits, expected_rate_limits);
1438 assert_eq!(response_data, expected_data);
1439 });
1440 }
1441
1442 #[test]
1443 fn account_commission_error_response() {
1444 TOKIO_SHARED_RT.block_on(async {
1445 let (ws_api, conn, mut rx) = setup().await;
1446 let client = AccountApiClient::new(ws_api.clone());
1447
1448 let handle = tokio::spawn(async move {
1449 let params = AccountCommissionParams::builder("BNBUSDT".to_string(),).build().unwrap();
1450 client.account_commission(params).await
1451 });
1452
1453 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
1454 let Message::Text(text) = sent else { panic!() };
1455 let v: Value = serde_json::from_str(&text).unwrap();
1456 let id = v["id"].as_str().unwrap().to_string();
1457
1458 let resp_json = json!({
1459 "id": id,
1460 "status": 400,
1461 "error": {
1462 "code": -2010,
1463 "msg": "Account has insufficient balance for requested action.",
1464 },
1465 "rateLimits": [
1466 {
1467 "rateLimitType": "ORDERS",
1468 "interval": "SECOND",
1469 "intervalNum": 10,
1470 "limit": 50,
1471 "count": 13
1472 },
1473 ],
1474 });
1475 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1476
1477 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
1478 match join {
1479 Ok(Err(e)) => {
1480 let msg = e.to_string();
1481 assert!(
1482 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
1483 "Expected error msg to contain server error, got: {msg}"
1484 );
1485 }
1486 Ok(Ok(_)) => panic!("Expected error"),
1487 Err(_) => panic!("Task panicked"),
1488 }
1489 });
1490 }
1491
1492 #[test]
1493 fn account_commission_request_timeout() {
1494 TOKIO_SHARED_RT.block_on(async {
1495 let (ws_api, _conn, mut rx) = setup().await;
1496 let client = AccountApiClient::new(ws_api.clone());
1497
1498 let handle = spawn(async move {
1499 let params = AccountCommissionParams::builder("BNBUSDT".to_string())
1500 .build()
1501 .unwrap();
1502 client.account_commission(params).await
1503 });
1504
1505 let sent = timeout(Duration::from_secs(1), rx.recv())
1506 .await
1507 .expect("send should occur")
1508 .expect("channel closed");
1509 let Message::Text(text) = sent else {
1510 panic!("expected Message Text")
1511 };
1512
1513 let _: Value = serde_json::from_str(&text).unwrap();
1514
1515 let result = handle.await.expect("task completed");
1516 match result {
1517 Err(e) => {
1518 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
1519 assert!(matches!(inner, WebsocketError::Timeout));
1520 } else {
1521 panic!("Unexpected error type: {:?}", e);
1522 }
1523 }
1524 Ok(_) => panic!("Expected timeout error"),
1525 }
1526 });
1527 }
1528
1529 #[test]
1530 fn account_rate_limits_orders_success() {
1531 TOKIO_SHARED_RT.block_on(async {
1532 let (ws_api, conn, mut rx) = setup().await;
1533 let client = AccountApiClient::new(ws_api.clone());
1534
1535 let handle = spawn(async move {
1536 let params = AccountRateLimitsOrdersParams::builder().build().unwrap();
1537 client.account_rate_limits_orders(params).await
1538 });
1539
1540 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
1541 let Message::Text(text) = sent else { panic!() };
1542 let v: Value = serde_json::from_str(&text).unwrap();
1543 let id = v["id"].as_str().unwrap();
1544 assert_eq!(v["method"], "/account.rateLimits.orders".trim_start_matches('/'));
1545 let mut resp_json: Value = serde_json::from_str(r#"{"id":"d3783d8d-f8d1-4d2c-b8a0-b7596af5a664","status":200,"result":[{"rateLimitType":"ORDERS","interval":"SECOND","intervalNum":10,"limit":50,"count":0}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
1546 resp_json["id"] = id.into();
1547
1548 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
1549 let expected_data: Vec<models::AccountRateLimitsOrdersResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
1550 let empty_array = Value::Array(vec![]);
1551 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
1552 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
1553 match raw_rate_limits.as_array() {
1554 Some(arr) if arr.is_empty() => None,
1555 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
1556 None => None,
1557 };
1558
1559 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1560
1561 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
1562
1563
1564 let response_rate_limits = response.rate_limits.clone();
1565 let response_data = response.data().expect("deserialize data");
1566
1567 assert_eq!(response_rate_limits, expected_rate_limits);
1568 assert_eq!(response_data, expected_data);
1569 });
1570 }
1571
1572 #[test]
1573 fn account_rate_limits_orders_error_response() {
1574 TOKIO_SHARED_RT.block_on(async {
1575 let (ws_api, conn, mut rx) = setup().await;
1576 let client = AccountApiClient::new(ws_api.clone());
1577
1578 let handle = tokio::spawn(async move {
1579 let params = AccountRateLimitsOrdersParams::builder().build().unwrap();
1580 client.account_rate_limits_orders(params).await
1581 });
1582
1583 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
1584 let Message::Text(text) = sent else { panic!() };
1585 let v: Value = serde_json::from_str(&text).unwrap();
1586 let id = v["id"].as_str().unwrap().to_string();
1587
1588 let resp_json = json!({
1589 "id": id,
1590 "status": 400,
1591 "error": {
1592 "code": -2010,
1593 "msg": "Account has insufficient balance for requested action.",
1594 },
1595 "rateLimits": [
1596 {
1597 "rateLimitType": "ORDERS",
1598 "interval": "SECOND",
1599 "intervalNum": 10,
1600 "limit": 50,
1601 "count": 13
1602 },
1603 ],
1604 });
1605 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1606
1607 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
1608 match join {
1609 Ok(Err(e)) => {
1610 let msg = e.to_string();
1611 assert!(
1612 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
1613 "Expected error msg to contain server error, got: {msg}"
1614 );
1615 }
1616 Ok(Ok(_)) => panic!("Expected error"),
1617 Err(_) => panic!("Task panicked"),
1618 }
1619 });
1620 }
1621
1622 #[test]
1623 fn account_rate_limits_orders_request_timeout() {
1624 TOKIO_SHARED_RT.block_on(async {
1625 let (ws_api, _conn, mut rx) = setup().await;
1626 let client = AccountApiClient::new(ws_api.clone());
1627
1628 let handle = spawn(async move {
1629 let params = AccountRateLimitsOrdersParams::builder().build().unwrap();
1630 client.account_rate_limits_orders(params).await
1631 });
1632
1633 let sent = timeout(Duration::from_secs(1), rx.recv())
1634 .await
1635 .expect("send should occur")
1636 .expect("channel closed");
1637 let Message::Text(text) = sent else {
1638 panic!("expected Message Text")
1639 };
1640
1641 let _: Value = serde_json::from_str(&text).unwrap();
1642
1643 let result = handle.await.expect("task completed");
1644 match result {
1645 Err(e) => {
1646 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
1647 assert!(matches!(inner, WebsocketError::Timeout));
1648 } else {
1649 panic!("Unexpected error type: {:?}", e);
1650 }
1651 }
1652 Ok(_) => panic!("Expected timeout error"),
1653 }
1654 });
1655 }
1656
1657 #[test]
1658 fn account_status_success() {
1659 TOKIO_SHARED_RT.block_on(async {
1660 let (ws_api, conn, mut rx) = setup().await;
1661 let client = AccountApiClient::new(ws_api.clone());
1662
1663 let handle = spawn(async move {
1664 let params = AccountStatusParams::builder().build().unwrap();
1665 client.account_status(params).await
1666 });
1667
1668 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
1669 let Message::Text(text) = sent else { panic!() };
1670 let v: Value = serde_json::from_str(&text).unwrap();
1671 let id = v["id"].as_str().unwrap();
1672 assert_eq!(v["method"], "/account.status".trim_start_matches('/'));
1673 let mut resp_json: Value = serde_json::from_str(r#"{"id":"605a6d20-6588-4cb9-afa0-b0ab087507ba","status":200,"result":{"makerCommission":15,"takerCommission":15,"buyerCommission":0,"sellerCommission":0,"canTrade":true,"canWithdraw":true,"canDeposit":true,"brokered":false,"requireSelfTradePrevention":false,"preventSor":false,"updateTime":1660801833000,"accountType":"SPOT","balances":[{"asset":"BNB"}],"permissions":["SPOT"],"uid":354937868},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
1674 resp_json["id"] = id.into();
1675
1676 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
1677 let expected_data: Box<models::AccountStatusResponseResult> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
1678 let empty_array = Value::Array(vec![]);
1679 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
1680 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
1681 match raw_rate_limits.as_array() {
1682 Some(arr) if arr.is_empty() => None,
1683 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
1684 None => None,
1685 };
1686
1687 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1688
1689 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
1690
1691
1692 let response_rate_limits = response.rate_limits.clone();
1693 let response_data = response.data().expect("deserialize data");
1694
1695 assert_eq!(response_rate_limits, expected_rate_limits);
1696 assert_eq!(response_data, expected_data);
1697 });
1698 }
1699
1700 #[test]
1701 fn account_status_error_response() {
1702 TOKIO_SHARED_RT.block_on(async {
1703 let (ws_api, conn, mut rx) = setup().await;
1704 let client = AccountApiClient::new(ws_api.clone());
1705
1706 let handle = tokio::spawn(async move {
1707 let params = AccountStatusParams::builder().build().unwrap();
1708 client.account_status(params).await
1709 });
1710
1711 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
1712 let Message::Text(text) = sent else { panic!() };
1713 let v: Value = serde_json::from_str(&text).unwrap();
1714 let id = v["id"].as_str().unwrap().to_string();
1715
1716 let resp_json = json!({
1717 "id": id,
1718 "status": 400,
1719 "error": {
1720 "code": -2010,
1721 "msg": "Account has insufficient balance for requested action.",
1722 },
1723 "rateLimits": [
1724 {
1725 "rateLimitType": "ORDERS",
1726 "interval": "SECOND",
1727 "intervalNum": 10,
1728 "limit": 50,
1729 "count": 13
1730 },
1731 ],
1732 });
1733 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1734
1735 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
1736 match join {
1737 Ok(Err(e)) => {
1738 let msg = e.to_string();
1739 assert!(
1740 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
1741 "Expected error msg to contain server error, got: {msg}"
1742 );
1743 }
1744 Ok(Ok(_)) => panic!("Expected error"),
1745 Err(_) => panic!("Task panicked"),
1746 }
1747 });
1748 }
1749
1750 #[test]
1751 fn account_status_request_timeout() {
1752 TOKIO_SHARED_RT.block_on(async {
1753 let (ws_api, _conn, mut rx) = setup().await;
1754 let client = AccountApiClient::new(ws_api.clone());
1755
1756 let handle = spawn(async move {
1757 let params = AccountStatusParams::builder().build().unwrap();
1758 client.account_status(params).await
1759 });
1760
1761 let sent = timeout(Duration::from_secs(1), rx.recv())
1762 .await
1763 .expect("send should occur")
1764 .expect("channel closed");
1765 let Message::Text(text) = sent else {
1766 panic!("expected Message Text")
1767 };
1768
1769 let _: Value = serde_json::from_str(&text).unwrap();
1770
1771 let result = handle.await.expect("task completed");
1772 match result {
1773 Err(e) => {
1774 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
1775 assert!(matches!(inner, WebsocketError::Timeout));
1776 } else {
1777 panic!("Unexpected error type: {:?}", e);
1778 }
1779 }
1780 Ok(_) => panic!("Expected timeout error"),
1781 }
1782 });
1783 }
1784
1785 #[test]
1786 fn all_order_lists_success() {
1787 TOKIO_SHARED_RT.block_on(async {
1788 let (ws_api, conn, mut rx) = setup().await;
1789 let client = AccountApiClient::new(ws_api.clone());
1790
1791 let handle = spawn(async move {
1792 let params = AllOrderListsParams::builder().build().unwrap();
1793 client.all_order_lists(params).await
1794 });
1795
1796 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
1797 let Message::Text(text) = sent else { panic!() };
1798 let v: Value = serde_json::from_str(&text).unwrap();
1799 let id = v["id"].as_str().unwrap();
1800 assert_eq!(v["method"], "/allOrderLists".trim_start_matches('/'));
1801 let mut resp_json: Value = serde_json::from_str(r#"{"id":"8617b7b3-1b3d-4dec-94cd-eefd929b8ceb","status":200,"result":[{"orderListId":1274512,"contingencyType":"OCO","listStatusType":"EXEC_STARTED","listOrderStatus":"EXECUTING","listClientOrderId":"08985fedd9ea2cf6b28996","transactionTime":1660801713793,"symbol":"BTCUSDT","orders":[{"symbol":"BTCUSDT","orderId":12569138901,"clientOrderId":"BqtFCj5odMoWtSqGk2X9tU"}]}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
1802 resp_json["id"] = id.into();
1803
1804 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
1805 let expected_data: Vec<models::AllOrderListsResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
1806 let empty_array = Value::Array(vec![]);
1807 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
1808 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
1809 match raw_rate_limits.as_array() {
1810 Some(arr) if arr.is_empty() => None,
1811 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
1812 None => None,
1813 };
1814
1815 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1816
1817 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
1818
1819
1820 let response_rate_limits = response.rate_limits.clone();
1821 let response_data = response.data().expect("deserialize data");
1822
1823 assert_eq!(response_rate_limits, expected_rate_limits);
1824 assert_eq!(response_data, expected_data);
1825 });
1826 }
1827
1828 #[test]
1829 fn all_order_lists_error_response() {
1830 TOKIO_SHARED_RT.block_on(async {
1831 let (ws_api, conn, mut rx) = setup().await;
1832 let client = AccountApiClient::new(ws_api.clone());
1833
1834 let handle = tokio::spawn(async move {
1835 let params = AllOrderListsParams::builder().build().unwrap();
1836 client.all_order_lists(params).await
1837 });
1838
1839 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
1840 let Message::Text(text) = sent else { panic!() };
1841 let v: Value = serde_json::from_str(&text).unwrap();
1842 let id = v["id"].as_str().unwrap().to_string();
1843
1844 let resp_json = json!({
1845 "id": id,
1846 "status": 400,
1847 "error": {
1848 "code": -2010,
1849 "msg": "Account has insufficient balance for requested action.",
1850 },
1851 "rateLimits": [
1852 {
1853 "rateLimitType": "ORDERS",
1854 "interval": "SECOND",
1855 "intervalNum": 10,
1856 "limit": 50,
1857 "count": 13
1858 },
1859 ],
1860 });
1861 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1862
1863 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
1864 match join {
1865 Ok(Err(e)) => {
1866 let msg = e.to_string();
1867 assert!(
1868 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
1869 "Expected error msg to contain server error, got: {msg}"
1870 );
1871 }
1872 Ok(Ok(_)) => panic!("Expected error"),
1873 Err(_) => panic!("Task panicked"),
1874 }
1875 });
1876 }
1877
1878 #[test]
1879 fn all_order_lists_request_timeout() {
1880 TOKIO_SHARED_RT.block_on(async {
1881 let (ws_api, _conn, mut rx) = setup().await;
1882 let client = AccountApiClient::new(ws_api.clone());
1883
1884 let handle = spawn(async move {
1885 let params = AllOrderListsParams::builder().build().unwrap();
1886 client.all_order_lists(params).await
1887 });
1888
1889 let sent = timeout(Duration::from_secs(1), rx.recv())
1890 .await
1891 .expect("send should occur")
1892 .expect("channel closed");
1893 let Message::Text(text) = sent else {
1894 panic!("expected Message Text")
1895 };
1896
1897 let _: Value = serde_json::from_str(&text).unwrap();
1898
1899 let result = handle.await.expect("task completed");
1900 match result {
1901 Err(e) => {
1902 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
1903 assert!(matches!(inner, WebsocketError::Timeout));
1904 } else {
1905 panic!("Unexpected error type: {:?}", e);
1906 }
1907 }
1908 Ok(_) => panic!("Expected timeout error"),
1909 }
1910 });
1911 }
1912
1913 #[test]
1914 fn all_orders_success() {
1915 TOKIO_SHARED_RT.block_on(async {
1916 let (ws_api, conn, mut rx) = setup().await;
1917 let client = AccountApiClient::new(ws_api.clone());
1918
1919 let handle = spawn(async move {
1920 let params = AllOrdersParams::builder("BNBUSDT".to_string(),).build().unwrap();
1921 client.all_orders(params).await
1922 });
1923
1924 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
1925 let Message::Text(text) = sent else { panic!() };
1926 let v: Value = serde_json::from_str(&text).unwrap();
1927 let id = v["id"].as_str().unwrap();
1928 assert_eq!(v["method"], "/allOrders".trim_start_matches('/'));
1929 let mut resp_json: Value = serde_json::from_str(r#"{"id":"734235c2-13d2-4574-be68-723e818c08f3","status":200,"result":[{"symbol":"BTCUSDT","orderId":12569099453,"orderListId":-1,"clientOrderId":"4d96324ff9d44481926157","status":"FILLED","timeInForce":"GTC","type":"LIMIT","side":"SELL","time":1660801715639,"updateTime":1660801717945,"isWorking":true,"workingTime":1660801715639,"selfTradePreventionMode":"NONE","preventedMatchId":0,"icebergQty":"0.00000000","stopPrice":"0.00000000","strategyId":1,"strategyType":1000000,"trailingDelta":10,"trailingTime":-1,"usedSor":true,"workingFloor":"SOR","pegPriceType":"PRIMARY_PEG","pegOffsetType":"PRICE_LEVEL","pegOffsetValue":5,"peggedPrice":"87523.83710000","expiryReason":"INSUFFICIENT_LIQUIDITY"}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
1930 resp_json["id"] = id.into();
1931
1932 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
1933 let expected_data: Vec<models::AllOrdersResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
1934 let empty_array = Value::Array(vec![]);
1935 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
1936 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
1937 match raw_rate_limits.as_array() {
1938 Some(arr) if arr.is_empty() => None,
1939 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
1940 None => None,
1941 };
1942
1943 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1944
1945 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
1946
1947
1948 let response_rate_limits = response.rate_limits.clone();
1949 let response_data = response.data().expect("deserialize data");
1950
1951 assert_eq!(response_rate_limits, expected_rate_limits);
1952 assert_eq!(response_data, expected_data);
1953 });
1954 }
1955
1956 #[test]
1957 fn all_orders_error_response() {
1958 TOKIO_SHARED_RT.block_on(async {
1959 let (ws_api, conn, mut rx) = setup().await;
1960 let client = AccountApiClient::new(ws_api.clone());
1961
1962 let handle = tokio::spawn(async move {
1963 let params = AllOrdersParams::builder("BNBUSDT".to_string(),).build().unwrap();
1964 client.all_orders(params).await
1965 });
1966
1967 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
1968 let Message::Text(text) = sent else { panic!() };
1969 let v: Value = serde_json::from_str(&text).unwrap();
1970 let id = v["id"].as_str().unwrap().to_string();
1971
1972 let resp_json = json!({
1973 "id": id,
1974 "status": 400,
1975 "error": {
1976 "code": -2010,
1977 "msg": "Account has insufficient balance for requested action.",
1978 },
1979 "rateLimits": [
1980 {
1981 "rateLimitType": "ORDERS",
1982 "interval": "SECOND",
1983 "intervalNum": 10,
1984 "limit": 50,
1985 "count": 13
1986 },
1987 ],
1988 });
1989 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
1990
1991 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
1992 match join {
1993 Ok(Err(e)) => {
1994 let msg = e.to_string();
1995 assert!(
1996 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
1997 "Expected error msg to contain server error, got: {msg}"
1998 );
1999 }
2000 Ok(Ok(_)) => panic!("Expected error"),
2001 Err(_) => panic!("Task panicked"),
2002 }
2003 });
2004 }
2005
2006 #[test]
2007 fn all_orders_request_timeout() {
2008 TOKIO_SHARED_RT.block_on(async {
2009 let (ws_api, _conn, mut rx) = setup().await;
2010 let client = AccountApiClient::new(ws_api.clone());
2011
2012 let handle = spawn(async move {
2013 let params = AllOrdersParams::builder("BNBUSDT".to_string())
2014 .build()
2015 .unwrap();
2016 client.all_orders(params).await
2017 });
2018
2019 let sent = timeout(Duration::from_secs(1), rx.recv())
2020 .await
2021 .expect("send should occur")
2022 .expect("channel closed");
2023 let Message::Text(text) = sent else {
2024 panic!("expected Message Text")
2025 };
2026
2027 let _: Value = serde_json::from_str(&text).unwrap();
2028
2029 let result = handle.await.expect("task completed");
2030 match result {
2031 Err(e) => {
2032 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2033 assert!(matches!(inner, WebsocketError::Timeout));
2034 } else {
2035 panic!("Unexpected error type: {:?}", e);
2036 }
2037 }
2038 Ok(_) => panic!("Expected timeout error"),
2039 }
2040 });
2041 }
2042
2043 #[test]
2044 fn my_allocations_success() {
2045 TOKIO_SHARED_RT.block_on(async {
2046 let (ws_api, conn, mut rx) = setup().await;
2047 let client = AccountApiClient::new(ws_api.clone());
2048
2049 let handle = spawn(async move {
2050 let params = MyAllocationsParams::builder("BNBUSDT".to_string(),).build().unwrap();
2051 client.my_allocations(params).await
2052 });
2053
2054 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2055 let Message::Text(text) = sent else { panic!() };
2056 let v: Value = serde_json::from_str(&text).unwrap();
2057 let id = v["id"].as_str().unwrap();
2058 assert_eq!(v["method"], "/myAllocations".trim_start_matches('/'));
2059 let mut resp_json: Value = serde_json::from_str(r#"{"id":"g4ce6a53-a39d-4f71-823b-4ab5r391d6y8","status":200,"result":[{"symbol":"BTCUSDT","allocationId":0,"allocationType":"SOR","orderId":500,"orderListId":-1,"commissionAsset":"BTC","time":1687319487614,"isBuyer":false,"isMaker":false,"isAllocator":false}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2060 resp_json["id"] = id.into();
2061
2062 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2063 let expected_data: Vec<models::MyAllocationsResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2064 let empty_array = Value::Array(vec![]);
2065 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2066 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2067 match raw_rate_limits.as_array() {
2068 Some(arr) if arr.is_empty() => None,
2069 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2070 None => None,
2071 };
2072
2073 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2074
2075 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2076
2077
2078 let response_rate_limits = response.rate_limits.clone();
2079 let response_data = response.data().expect("deserialize data");
2080
2081 assert_eq!(response_rate_limits, expected_rate_limits);
2082 assert_eq!(response_data, expected_data);
2083 });
2084 }
2085
2086 #[test]
2087 fn my_allocations_error_response() {
2088 TOKIO_SHARED_RT.block_on(async {
2089 let (ws_api, conn, mut rx) = setup().await;
2090 let client = AccountApiClient::new(ws_api.clone());
2091
2092 let handle = tokio::spawn(async move {
2093 let params = MyAllocationsParams::builder("BNBUSDT".to_string(),).build().unwrap();
2094 client.my_allocations(params).await
2095 });
2096
2097 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2098 let Message::Text(text) = sent else { panic!() };
2099 let v: Value = serde_json::from_str(&text).unwrap();
2100 let id = v["id"].as_str().unwrap().to_string();
2101
2102 let resp_json = json!({
2103 "id": id,
2104 "status": 400,
2105 "error": {
2106 "code": -2010,
2107 "msg": "Account has insufficient balance for requested action.",
2108 },
2109 "rateLimits": [
2110 {
2111 "rateLimitType": "ORDERS",
2112 "interval": "SECOND",
2113 "intervalNum": 10,
2114 "limit": 50,
2115 "count": 13
2116 },
2117 ],
2118 });
2119 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2120
2121 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2122 match join {
2123 Ok(Err(e)) => {
2124 let msg = e.to_string();
2125 assert!(
2126 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2127 "Expected error msg to contain server error, got: {msg}"
2128 );
2129 }
2130 Ok(Ok(_)) => panic!("Expected error"),
2131 Err(_) => panic!("Task panicked"),
2132 }
2133 });
2134 }
2135
2136 #[test]
2137 fn my_allocations_request_timeout() {
2138 TOKIO_SHARED_RT.block_on(async {
2139 let (ws_api, _conn, mut rx) = setup().await;
2140 let client = AccountApiClient::new(ws_api.clone());
2141
2142 let handle = spawn(async move {
2143 let params = MyAllocationsParams::builder("BNBUSDT".to_string())
2144 .build()
2145 .unwrap();
2146 client.my_allocations(params).await
2147 });
2148
2149 let sent = timeout(Duration::from_secs(1), rx.recv())
2150 .await
2151 .expect("send should occur")
2152 .expect("channel closed");
2153 let Message::Text(text) = sent else {
2154 panic!("expected Message Text")
2155 };
2156
2157 let _: Value = serde_json::from_str(&text).unwrap();
2158
2159 let result = handle.await.expect("task completed");
2160 match result {
2161 Err(e) => {
2162 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2163 assert!(matches!(inner, WebsocketError::Timeout));
2164 } else {
2165 panic!("Unexpected error type: {:?}", e);
2166 }
2167 }
2168 Ok(_) => panic!("Expected timeout error"),
2169 }
2170 });
2171 }
2172
2173 #[test]
2174 fn my_filters_success() {
2175 TOKIO_SHARED_RT.block_on(async {
2176 let (ws_api, conn, mut rx) = setup().await;
2177 let client = AccountApiClient::new(ws_api.clone());
2178
2179 let handle = spawn(async move {
2180 let params = MyFiltersParams::builder("BNBUSDT".to_string(),).build().unwrap();
2181 client.my_filters(params).await
2182 });
2183
2184 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2185 let Message::Text(text) = sent else { panic!() };
2186 let v: Value = serde_json::from_str(&text).unwrap();
2187 let id = v["id"].as_str().unwrap();
2188 assert_eq!(v["method"], "/myFilters".trim_start_matches('/'));
2189 let mut resp_json: Value = serde_json::from_str(r#"{"status":200,"result":{"exchangeFilters":[{"filterType":"EXCHANGE_MAX_NUM_ORDERS","maxNumOrders":1000}],"symbolFilters":[{"filterType":"PRICE_FILTER","priceExponent":8,"minPrice":"0.00000100","maxPrice":"100000.00000000","tickSize":"0.00000100"}],"assetFilters":[{"filterType":"MAX_ASSET","qtyExponent":8,"limit":"1000000.00000000","asset":"JPY"}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}}"#).unwrap_or_else(|_| serde_json::json!({}));
2190 resp_json["id"] = id.into();
2191
2192 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2193 let expected_data: models::MyFiltersResponse = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2194 let empty_array = Value::Array(vec![]);
2195 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2196 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2197 match raw_rate_limits.as_array() {
2198 Some(arr) if arr.is_empty() => None,
2199 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2200 None => None,
2201 };
2202
2203 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2204
2205 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2206
2207
2208 let response_rate_limits = response.rate_limits.clone();
2209 let response_data = response.data().expect("deserialize data");
2210
2211 assert_eq!(response_rate_limits, expected_rate_limits);
2212 assert_eq!(response_data, expected_data);
2213 });
2214 }
2215
2216 #[test]
2217 fn my_filters_error_response() {
2218 TOKIO_SHARED_RT.block_on(async {
2219 let (ws_api, conn, mut rx) = setup().await;
2220 let client = AccountApiClient::new(ws_api.clone());
2221
2222 let handle = tokio::spawn(async move {
2223 let params = MyFiltersParams::builder("BNBUSDT".to_string(),).build().unwrap();
2224 client.my_filters(params).await
2225 });
2226
2227 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2228 let Message::Text(text) = sent else { panic!() };
2229 let v: Value = serde_json::from_str(&text).unwrap();
2230 let id = v["id"].as_str().unwrap().to_string();
2231
2232 let resp_json = json!({
2233 "id": id,
2234 "status": 400,
2235 "error": {
2236 "code": -2010,
2237 "msg": "Account has insufficient balance for requested action.",
2238 },
2239 "rateLimits": [
2240 {
2241 "rateLimitType": "ORDERS",
2242 "interval": "SECOND",
2243 "intervalNum": 10,
2244 "limit": 50,
2245 "count": 13
2246 },
2247 ],
2248 });
2249 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2250
2251 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2252 match join {
2253 Ok(Err(e)) => {
2254 let msg = e.to_string();
2255 assert!(
2256 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2257 "Expected error msg to contain server error, got: {msg}"
2258 );
2259 }
2260 Ok(Ok(_)) => panic!("Expected error"),
2261 Err(_) => panic!("Task panicked"),
2262 }
2263 });
2264 }
2265
2266 #[test]
2267 fn my_filters_request_timeout() {
2268 TOKIO_SHARED_RT.block_on(async {
2269 let (ws_api, _conn, mut rx) = setup().await;
2270 let client = AccountApiClient::new(ws_api.clone());
2271
2272 let handle = spawn(async move {
2273 let params = MyFiltersParams::builder("BNBUSDT".to_string())
2274 .build()
2275 .unwrap();
2276 client.my_filters(params).await
2277 });
2278
2279 let sent = timeout(Duration::from_secs(1), rx.recv())
2280 .await
2281 .expect("send should occur")
2282 .expect("channel closed");
2283 let Message::Text(text) = sent else {
2284 panic!("expected Message Text")
2285 };
2286
2287 let _: Value = serde_json::from_str(&text).unwrap();
2288
2289 let result = handle.await.expect("task completed");
2290 match result {
2291 Err(e) => {
2292 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2293 assert!(matches!(inner, WebsocketError::Timeout));
2294 } else {
2295 panic!("Unexpected error type: {:?}", e);
2296 }
2297 }
2298 Ok(_) => panic!("Expected timeout error"),
2299 }
2300 });
2301 }
2302
2303 #[test]
2304 fn my_prevented_matches_success() {
2305 TOKIO_SHARED_RT.block_on(async {
2306 let (ws_api, conn, mut rx) = setup().await;
2307 let client = AccountApiClient::new(ws_api.clone());
2308
2309 let handle = spawn(async move {
2310 let params = MyPreventedMatchesParams::builder("BNBUSDT".to_string(),).build().unwrap();
2311 client.my_prevented_matches(params).await
2312 });
2313
2314 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2315 let Message::Text(text) = sent else { panic!() };
2316 let v: Value = serde_json::from_str(&text).unwrap();
2317 let id = v["id"].as_str().unwrap();
2318 assert_eq!(v["method"], "/myPreventedMatches".trim_start_matches('/'));
2319 let mut resp_json: Value = serde_json::from_str(r#"{"id":"g4ce6a53-a39d-4f71-823b-4ab5r391d6y8","status":200,"result":[{"symbol":"BTCUSDT","preventedMatchId":1,"takerOrderId":5,"makerSymbol":"BTCUSDT","makerOrderId":3,"tradeGroupId":1,"selfTradePreventionMode":"EXPIRE_MAKER","transactTime":1669101687094}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2320 resp_json["id"] = id.into();
2321
2322 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2323 let expected_data: Vec<models::MyPreventedMatchesResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2324 let empty_array = Value::Array(vec![]);
2325 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2326 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2327 match raw_rate_limits.as_array() {
2328 Some(arr) if arr.is_empty() => None,
2329 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2330 None => None,
2331 };
2332
2333 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2334
2335 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2336
2337
2338 let response_rate_limits = response.rate_limits.clone();
2339 let response_data = response.data().expect("deserialize data");
2340
2341 assert_eq!(response_rate_limits, expected_rate_limits);
2342 assert_eq!(response_data, expected_data);
2343 });
2344 }
2345
2346 #[test]
2347 fn my_prevented_matches_error_response() {
2348 TOKIO_SHARED_RT.block_on(async {
2349 let (ws_api, conn, mut rx) = setup().await;
2350 let client = AccountApiClient::new(ws_api.clone());
2351
2352 let handle = tokio::spawn(async move {
2353 let params = MyPreventedMatchesParams::builder("BNBUSDT".to_string(),).build().unwrap();
2354 client.my_prevented_matches(params).await
2355 });
2356
2357 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2358 let Message::Text(text) = sent else { panic!() };
2359 let v: Value = serde_json::from_str(&text).unwrap();
2360 let id = v["id"].as_str().unwrap().to_string();
2361
2362 let resp_json = json!({
2363 "id": id,
2364 "status": 400,
2365 "error": {
2366 "code": -2010,
2367 "msg": "Account has insufficient balance for requested action.",
2368 },
2369 "rateLimits": [
2370 {
2371 "rateLimitType": "ORDERS",
2372 "interval": "SECOND",
2373 "intervalNum": 10,
2374 "limit": 50,
2375 "count": 13
2376 },
2377 ],
2378 });
2379 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2380
2381 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2382 match join {
2383 Ok(Err(e)) => {
2384 let msg = e.to_string();
2385 assert!(
2386 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2387 "Expected error msg to contain server error, got: {msg}"
2388 );
2389 }
2390 Ok(Ok(_)) => panic!("Expected error"),
2391 Err(_) => panic!("Task panicked"),
2392 }
2393 });
2394 }
2395
2396 #[test]
2397 fn my_prevented_matches_request_timeout() {
2398 TOKIO_SHARED_RT.block_on(async {
2399 let (ws_api, _conn, mut rx) = setup().await;
2400 let client = AccountApiClient::new(ws_api.clone());
2401
2402 let handle = spawn(async move {
2403 let params = MyPreventedMatchesParams::builder("BNBUSDT".to_string())
2404 .build()
2405 .unwrap();
2406 client.my_prevented_matches(params).await
2407 });
2408
2409 let sent = timeout(Duration::from_secs(1), rx.recv())
2410 .await
2411 .expect("send should occur")
2412 .expect("channel closed");
2413 let Message::Text(text) = sent else {
2414 panic!("expected Message Text")
2415 };
2416
2417 let _: Value = serde_json::from_str(&text).unwrap();
2418
2419 let result = handle.await.expect("task completed");
2420 match result {
2421 Err(e) => {
2422 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2423 assert!(matches!(inner, WebsocketError::Timeout));
2424 } else {
2425 panic!("Unexpected error type: {:?}", e);
2426 }
2427 }
2428 Ok(_) => panic!("Expected timeout error"),
2429 }
2430 });
2431 }
2432
2433 #[test]
2434 fn my_trades_success() {
2435 TOKIO_SHARED_RT.block_on(async {
2436 let (ws_api, conn, mut rx) = setup().await;
2437 let client = AccountApiClient::new(ws_api.clone());
2438
2439 let handle = spawn(async move {
2440 let params = MyTradesParams::builder("BNBUSDT".to_string(),).build().unwrap();
2441 client.my_trades(params).await
2442 });
2443
2444 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2445 let Message::Text(text) = sent else { panic!() };
2446 let v: Value = serde_json::from_str(&text).unwrap();
2447 let id = v["id"].as_str().unwrap();
2448 assert_eq!(v["method"], "/myTrades".trim_start_matches('/'));
2449 let mut resp_json: Value = serde_json::from_str(r#"{"id":"f4ce6a53-a29d-4f70-823b-4ab59391d6e8","status":200,"result":[{"symbol":"BTCUSDT","id":1650422481,"orderId":12569099453,"orderListId":-1,"commissionAsset":"BNB","time":1660801715793,"isBuyer":false,"isMaker":true,"isBestMatch":true}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2450 resp_json["id"] = id.into();
2451
2452 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2453 let expected_data: Vec<models::MyTradesResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2454 let empty_array = Value::Array(vec![]);
2455 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2456 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2457 match raw_rate_limits.as_array() {
2458 Some(arr) if arr.is_empty() => None,
2459 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2460 None => None,
2461 };
2462
2463 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2464
2465 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2466
2467
2468 let response_rate_limits = response.rate_limits.clone();
2469 let response_data = response.data().expect("deserialize data");
2470
2471 assert_eq!(response_rate_limits, expected_rate_limits);
2472 assert_eq!(response_data, expected_data);
2473 });
2474 }
2475
2476 #[test]
2477 fn my_trades_error_response() {
2478 TOKIO_SHARED_RT.block_on(async {
2479 let (ws_api, conn, mut rx) = setup().await;
2480 let client = AccountApiClient::new(ws_api.clone());
2481
2482 let handle = tokio::spawn(async move {
2483 let params = MyTradesParams::builder("BNBUSDT".to_string(),).build().unwrap();
2484 client.my_trades(params).await
2485 });
2486
2487 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2488 let Message::Text(text) = sent else { panic!() };
2489 let v: Value = serde_json::from_str(&text).unwrap();
2490 let id = v["id"].as_str().unwrap().to_string();
2491
2492 let resp_json = json!({
2493 "id": id,
2494 "status": 400,
2495 "error": {
2496 "code": -2010,
2497 "msg": "Account has insufficient balance for requested action.",
2498 },
2499 "rateLimits": [
2500 {
2501 "rateLimitType": "ORDERS",
2502 "interval": "SECOND",
2503 "intervalNum": 10,
2504 "limit": 50,
2505 "count": 13
2506 },
2507 ],
2508 });
2509 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2510
2511 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2512 match join {
2513 Ok(Err(e)) => {
2514 let msg = e.to_string();
2515 assert!(
2516 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2517 "Expected error msg to contain server error, got: {msg}"
2518 );
2519 }
2520 Ok(Ok(_)) => panic!("Expected error"),
2521 Err(_) => panic!("Task panicked"),
2522 }
2523 });
2524 }
2525
2526 #[test]
2527 fn my_trades_request_timeout() {
2528 TOKIO_SHARED_RT.block_on(async {
2529 let (ws_api, _conn, mut rx) = setup().await;
2530 let client = AccountApiClient::new(ws_api.clone());
2531
2532 let handle = spawn(async move {
2533 let params = MyTradesParams::builder("BNBUSDT".to_string())
2534 .build()
2535 .unwrap();
2536 client.my_trades(params).await
2537 });
2538
2539 let sent = timeout(Duration::from_secs(1), rx.recv())
2540 .await
2541 .expect("send should occur")
2542 .expect("channel closed");
2543 let Message::Text(text) = sent else {
2544 panic!("expected Message Text")
2545 };
2546
2547 let _: Value = serde_json::from_str(&text).unwrap();
2548
2549 let result = handle.await.expect("task completed");
2550 match result {
2551 Err(e) => {
2552 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2553 assert!(matches!(inner, WebsocketError::Timeout));
2554 } else {
2555 panic!("Unexpected error type: {:?}", e);
2556 }
2557 }
2558 Ok(_) => panic!("Expected timeout error"),
2559 }
2560 });
2561 }
2562
2563 #[test]
2564 fn open_order_lists_status_success() {
2565 TOKIO_SHARED_RT.block_on(async {
2566 let (ws_api, conn, mut rx) = setup().await;
2567 let client = AccountApiClient::new(ws_api.clone());
2568
2569 let handle = spawn(async move {
2570 let params = OpenOrderListsStatusParams::builder().build().unwrap();
2571 client.open_order_lists_status(params).await
2572 });
2573
2574 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2575 let Message::Text(text) = sent else { panic!() };
2576 let v: Value = serde_json::from_str(&text).unwrap();
2577 let id = v["id"].as_str().unwrap();
2578 assert_eq!(v["method"], "/openOrderLists.status".trim_start_matches('/'));
2579 let mut resp_json: Value = serde_json::from_str(r#"{"id":"3a4437e2-41a3-4c19-897c-9cadc5dce8b6","status":200,"result":[{"orderListId":0,"contingencyType":"OCO","listStatusType":"EXEC_STARTED","listOrderStatus":"EXECUTING","listClientOrderId":"08985fedd9ea2cf6b28996","transactionTime":1660801713793,"symbol":"BTCUSDT","orders":[{"symbol":"BTCUSDT","orderId":4,"clientOrderId":"CUhLgTXnX5n2c0gWiLpV4d"}]}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2580 resp_json["id"] = id.into();
2581
2582 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2583 let expected_data: Vec<models::OpenOrderListsStatusResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2584 let empty_array = Value::Array(vec![]);
2585 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2586 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2587 match raw_rate_limits.as_array() {
2588 Some(arr) if arr.is_empty() => None,
2589 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2590 None => None,
2591 };
2592
2593 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2594
2595 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2596
2597
2598 let response_rate_limits = response.rate_limits.clone();
2599 let response_data = response.data().expect("deserialize data");
2600
2601 assert_eq!(response_rate_limits, expected_rate_limits);
2602 assert_eq!(response_data, expected_data);
2603 });
2604 }
2605
2606 #[test]
2607 fn open_order_lists_status_error_response() {
2608 TOKIO_SHARED_RT.block_on(async {
2609 let (ws_api, conn, mut rx) = setup().await;
2610 let client = AccountApiClient::new(ws_api.clone());
2611
2612 let handle = tokio::spawn(async move {
2613 let params = OpenOrderListsStatusParams::builder().build().unwrap();
2614 client.open_order_lists_status(params).await
2615 });
2616
2617 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2618 let Message::Text(text) = sent else { panic!() };
2619 let v: Value = serde_json::from_str(&text).unwrap();
2620 let id = v["id"].as_str().unwrap().to_string();
2621
2622 let resp_json = json!({
2623 "id": id,
2624 "status": 400,
2625 "error": {
2626 "code": -2010,
2627 "msg": "Account has insufficient balance for requested action.",
2628 },
2629 "rateLimits": [
2630 {
2631 "rateLimitType": "ORDERS",
2632 "interval": "SECOND",
2633 "intervalNum": 10,
2634 "limit": 50,
2635 "count": 13
2636 },
2637 ],
2638 });
2639 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2640
2641 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2642 match join {
2643 Ok(Err(e)) => {
2644 let msg = e.to_string();
2645 assert!(
2646 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2647 "Expected error msg to contain server error, got: {msg}"
2648 );
2649 }
2650 Ok(Ok(_)) => panic!("Expected error"),
2651 Err(_) => panic!("Task panicked"),
2652 }
2653 });
2654 }
2655
2656 #[test]
2657 fn open_order_lists_status_request_timeout() {
2658 TOKIO_SHARED_RT.block_on(async {
2659 let (ws_api, _conn, mut rx) = setup().await;
2660 let client = AccountApiClient::new(ws_api.clone());
2661
2662 let handle = spawn(async move {
2663 let params = OpenOrderListsStatusParams::builder().build().unwrap();
2664 client.open_order_lists_status(params).await
2665 });
2666
2667 let sent = timeout(Duration::from_secs(1), rx.recv())
2668 .await
2669 .expect("send should occur")
2670 .expect("channel closed");
2671 let Message::Text(text) = sent else {
2672 panic!("expected Message Text")
2673 };
2674
2675 let _: Value = serde_json::from_str(&text).unwrap();
2676
2677 let result = handle.await.expect("task completed");
2678 match result {
2679 Err(e) => {
2680 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2681 assert!(matches!(inner, WebsocketError::Timeout));
2682 } else {
2683 panic!("Unexpected error type: {:?}", e);
2684 }
2685 }
2686 Ok(_) => panic!("Expected timeout error"),
2687 }
2688 });
2689 }
2690
2691 #[test]
2692 fn open_orders_status_success() {
2693 TOKIO_SHARED_RT.block_on(async {
2694 let (ws_api, conn, mut rx) = setup().await;
2695 let client = AccountApiClient::new(ws_api.clone());
2696
2697 let handle = spawn(async move {
2698 let params = OpenOrdersStatusParams::builder().build().unwrap();
2699 client.open_orders_status(params).await
2700 });
2701
2702 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2703 let Message::Text(text) = sent else { panic!() };
2704 let v: Value = serde_json::from_str(&text).unwrap();
2705 let id = v["id"].as_str().unwrap();
2706 assert_eq!(v["method"], "/openOrders.status".trim_start_matches('/'));
2707 let mut resp_json: Value = serde_json::from_str(r#"{"id":"55f07876-4f6f-4c47-87dc-43e5fff3f2e7","status":200,"result":[{"symbol":"BTCUSDT","orderId":12569099453,"orderListId":-1,"clientOrderId":"4d96324ff9d44481926157","status":"PARTIALLY_FILLED","timeInForce":"GTC","type":"LIMIT","side":"SELL","time":1660801715639,"updateTime":1660801717945,"isWorking":true,"workingTime":1660801715639,"selfTradePreventionMode":"NONE","icebergQty":"0.00000000","preventedMatchId":0,"preventedQuantity":"1.200000","stopPrice":"0.00000000","strategyId":1,"strategyType":1000000,"trailingDelta":10,"trailingTime":-1,"usedSor":true,"workingFloor":"SOR","pegPriceType":"PRIMARY_PEG","pegOffsetType":"PRICE_LEVEL","pegOffsetValue":5,"peggedPrice":"87523.83710000","expiryReason":"INSUFFICIENT_LIQUIDITY"}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2708 resp_json["id"] = id.into();
2709
2710 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2711 let expected_data: Vec<models::OpenOrdersStatusResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2712 let empty_array = Value::Array(vec![]);
2713 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2714 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2715 match raw_rate_limits.as_array() {
2716 Some(arr) if arr.is_empty() => None,
2717 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2718 None => None,
2719 };
2720
2721 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2722
2723 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2724
2725
2726 let response_rate_limits = response.rate_limits.clone();
2727 let response_data = response.data().expect("deserialize data");
2728
2729 assert_eq!(response_rate_limits, expected_rate_limits);
2730 assert_eq!(response_data, expected_data);
2731 });
2732 }
2733
2734 #[test]
2735 fn open_orders_status_error_response() {
2736 TOKIO_SHARED_RT.block_on(async {
2737 let (ws_api, conn, mut rx) = setup().await;
2738 let client = AccountApiClient::new(ws_api.clone());
2739
2740 let handle = tokio::spawn(async move {
2741 let params = OpenOrdersStatusParams::builder().build().unwrap();
2742 client.open_orders_status(params).await
2743 });
2744
2745 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2746 let Message::Text(text) = sent else { panic!() };
2747 let v: Value = serde_json::from_str(&text).unwrap();
2748 let id = v["id"].as_str().unwrap().to_string();
2749
2750 let resp_json = json!({
2751 "id": id,
2752 "status": 400,
2753 "error": {
2754 "code": -2010,
2755 "msg": "Account has insufficient balance for requested action.",
2756 },
2757 "rateLimits": [
2758 {
2759 "rateLimitType": "ORDERS",
2760 "interval": "SECOND",
2761 "intervalNum": 10,
2762 "limit": 50,
2763 "count": 13
2764 },
2765 ],
2766 });
2767 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2768
2769 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2770 match join {
2771 Ok(Err(e)) => {
2772 let msg = e.to_string();
2773 assert!(
2774 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2775 "Expected error msg to contain server error, got: {msg}"
2776 );
2777 }
2778 Ok(Ok(_)) => panic!("Expected error"),
2779 Err(_) => panic!("Task panicked"),
2780 }
2781 });
2782 }
2783
2784 #[test]
2785 fn open_orders_status_request_timeout() {
2786 TOKIO_SHARED_RT.block_on(async {
2787 let (ws_api, _conn, mut rx) = setup().await;
2788 let client = AccountApiClient::new(ws_api.clone());
2789
2790 let handle = spawn(async move {
2791 let params = OpenOrdersStatusParams::builder().build().unwrap();
2792 client.open_orders_status(params).await
2793 });
2794
2795 let sent = timeout(Duration::from_secs(1), rx.recv())
2796 .await
2797 .expect("send should occur")
2798 .expect("channel closed");
2799 let Message::Text(text) = sent else {
2800 panic!("expected Message Text")
2801 };
2802
2803 let _: Value = serde_json::from_str(&text).unwrap();
2804
2805 let result = handle.await.expect("task completed");
2806 match result {
2807 Err(e) => {
2808 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2809 assert!(matches!(inner, WebsocketError::Timeout));
2810 } else {
2811 panic!("Unexpected error type: {:?}", e);
2812 }
2813 }
2814 Ok(_) => panic!("Expected timeout error"),
2815 }
2816 });
2817 }
2818
2819 #[test]
2820 fn order_amendments_success() {
2821 TOKIO_SHARED_RT.block_on(async {
2822 let (ws_api, conn, mut rx) = setup().await;
2823 let client = AccountApiClient::new(ws_api.clone());
2824
2825 let handle = spawn(async move {
2826 let params = OrderAmendmentsParams::builder("BNBUSDT".to_string(),1,).build().unwrap();
2827 client.order_amendments(params).await
2828 });
2829
2830 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2831 let Message::Text(text) = sent else { panic!() };
2832 let v: Value = serde_json::from_str(&text).unwrap();
2833 let id = v["id"].as_str().unwrap();
2834 assert_eq!(v["method"], "/order.amendments".trim_start_matches('/'));
2835 let mut resp_json: Value = serde_json::from_str(r#"{"id":"6f5ebe91-01d9-43ac-be99-57cf062e0e30","status":200,"result":[{"symbol":"BTCUSDT","orderId":23,"executionId":60,"origClientOrderId":"my_pending_order","newClientOrderId":"xbxXh5SSwaHS7oUEOCI88B","time":1741924229819}],"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2836 resp_json["id"] = id.into();
2837
2838 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2839 let expected_data: Vec<models::OrderAmendmentsResponseResultInner> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2840 let empty_array = Value::Array(vec![]);
2841 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2842 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2843 match raw_rate_limits.as_array() {
2844 Some(arr) if arr.is_empty() => None,
2845 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2846 None => None,
2847 };
2848
2849 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2850
2851 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2852
2853
2854 let response_rate_limits = response.rate_limits.clone();
2855 let response_data = response.data().expect("deserialize data");
2856
2857 assert_eq!(response_rate_limits, expected_rate_limits);
2858 assert_eq!(response_data, expected_data);
2859 });
2860 }
2861
2862 #[test]
2863 fn order_amendments_error_response() {
2864 TOKIO_SHARED_RT.block_on(async {
2865 let (ws_api, conn, mut rx) = setup().await;
2866 let client = AccountApiClient::new(ws_api.clone());
2867
2868 let handle = tokio::spawn(async move {
2869 let params = OrderAmendmentsParams::builder("BNBUSDT".to_string(),1,).build().unwrap();
2870 client.order_amendments(params).await
2871 });
2872
2873 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
2874 let Message::Text(text) = sent else { panic!() };
2875 let v: Value = serde_json::from_str(&text).unwrap();
2876 let id = v["id"].as_str().unwrap().to_string();
2877
2878 let resp_json = json!({
2879 "id": id,
2880 "status": 400,
2881 "error": {
2882 "code": -2010,
2883 "msg": "Account has insufficient balance for requested action.",
2884 },
2885 "rateLimits": [
2886 {
2887 "rateLimitType": "ORDERS",
2888 "interval": "SECOND",
2889 "intervalNum": 10,
2890 "limit": 50,
2891 "count": 13
2892 },
2893 ],
2894 });
2895 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2896
2897 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
2898 match join {
2899 Ok(Err(e)) => {
2900 let msg = e.to_string();
2901 assert!(
2902 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
2903 "Expected error msg to contain server error, got: {msg}"
2904 );
2905 }
2906 Ok(Ok(_)) => panic!("Expected error"),
2907 Err(_) => panic!("Task panicked"),
2908 }
2909 });
2910 }
2911
2912 #[test]
2913 fn order_amendments_request_timeout() {
2914 TOKIO_SHARED_RT.block_on(async {
2915 let (ws_api, _conn, mut rx) = setup().await;
2916 let client = AccountApiClient::new(ws_api.clone());
2917
2918 let handle = spawn(async move {
2919 let params = OrderAmendmentsParams::builder("BNBUSDT".to_string(), 1)
2920 .build()
2921 .unwrap();
2922 client.order_amendments(params).await
2923 });
2924
2925 let sent = timeout(Duration::from_secs(1), rx.recv())
2926 .await
2927 .expect("send should occur")
2928 .expect("channel closed");
2929 let Message::Text(text) = sent else {
2930 panic!("expected Message Text")
2931 };
2932
2933 let _: Value = serde_json::from_str(&text).unwrap();
2934
2935 let result = handle.await.expect("task completed");
2936 match result {
2937 Err(e) => {
2938 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
2939 assert!(matches!(inner, WebsocketError::Timeout));
2940 } else {
2941 panic!("Unexpected error type: {:?}", e);
2942 }
2943 }
2944 Ok(_) => panic!("Expected timeout error"),
2945 }
2946 });
2947 }
2948
2949 #[test]
2950 fn order_list_status_success() {
2951 TOKIO_SHARED_RT.block_on(async {
2952 let (ws_api, conn, mut rx) = setup().await;
2953 let client = AccountApiClient::new(ws_api.clone());
2954
2955 let handle = spawn(async move {
2956 let params = OrderListStatusParams::builder().build().unwrap();
2957 client.order_list_status(params).await
2958 });
2959
2960 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
2961 let Message::Text(text) = sent else { panic!() };
2962 let v: Value = serde_json::from_str(&text).unwrap();
2963 let id = v["id"].as_str().unwrap();
2964 assert_eq!(v["method"], "/orderList.status".trim_start_matches('/'));
2965 let mut resp_json: Value = serde_json::from_str(r#"{"id":"b53fd5ff-82c7-4a04-bd64-5f9dc42c2100","status":200,"result":{"orderListId":1274512,"contingencyType":"OCO","listStatusType":"EXEC_STARTED","listOrderStatus":"EXECUTING","listClientOrderId":"08985fedd9ea2cf6b28996","transactionTime":1660801713793,"symbol":"BTCUSDT","orders":[{"symbol":"BTCUSDT","orderId":12569138901,"clientOrderId":"BqtFCj5odMoWtSqGk2X9tU"}]},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
2966 resp_json["id"] = id.into();
2967
2968 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
2969 let expected_data: Box<models::OrderListStatusResponseResult> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
2970 let empty_array = Value::Array(vec![]);
2971 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
2972 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
2973 match raw_rate_limits.as_array() {
2974 Some(arr) if arr.is_empty() => None,
2975 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
2976 None => None,
2977 };
2978
2979 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
2980
2981 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
2982
2983
2984 let response_rate_limits = response.rate_limits.clone();
2985 let response_data = response.data().expect("deserialize data");
2986
2987 assert_eq!(response_rate_limits, expected_rate_limits);
2988 assert_eq!(response_data, expected_data);
2989 });
2990 }
2991
2992 #[test]
2993 fn order_list_status_error_response() {
2994 TOKIO_SHARED_RT.block_on(async {
2995 let (ws_api, conn, mut rx) = setup().await;
2996 let client = AccountApiClient::new(ws_api.clone());
2997
2998 let handle = tokio::spawn(async move {
2999 let params = OrderListStatusParams::builder().build().unwrap();
3000 client.order_list_status(params).await
3001 });
3002
3003 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3004 let Message::Text(text) = sent else { panic!() };
3005 let v: Value = serde_json::from_str(&text).unwrap();
3006 let id = v["id"].as_str().unwrap().to_string();
3007
3008 let resp_json = json!({
3009 "id": id,
3010 "status": 400,
3011 "error": {
3012 "code": -2010,
3013 "msg": "Account has insufficient balance for requested action.",
3014 },
3015 "rateLimits": [
3016 {
3017 "rateLimitType": "ORDERS",
3018 "interval": "SECOND",
3019 "intervalNum": 10,
3020 "limit": 50,
3021 "count": 13
3022 },
3023 ],
3024 });
3025 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3026
3027 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3028 match join {
3029 Ok(Err(e)) => {
3030 let msg = e.to_string();
3031 assert!(
3032 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3033 "Expected error msg to contain server error, got: {msg}"
3034 );
3035 }
3036 Ok(Ok(_)) => panic!("Expected error"),
3037 Err(_) => panic!("Task panicked"),
3038 }
3039 });
3040 }
3041
3042 #[test]
3043 fn order_list_status_request_timeout() {
3044 TOKIO_SHARED_RT.block_on(async {
3045 let (ws_api, _conn, mut rx) = setup().await;
3046 let client = AccountApiClient::new(ws_api.clone());
3047
3048 let handle = spawn(async move {
3049 let params = OrderListStatusParams::builder().build().unwrap();
3050 client.order_list_status(params).await
3051 });
3052
3053 let sent = timeout(Duration::from_secs(1), rx.recv())
3054 .await
3055 .expect("send should occur")
3056 .expect("channel closed");
3057 let Message::Text(text) = sent else {
3058 panic!("expected Message Text")
3059 };
3060
3061 let _: Value = serde_json::from_str(&text).unwrap();
3062
3063 let result = handle.await.expect("task completed");
3064 match result {
3065 Err(e) => {
3066 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3067 assert!(matches!(inner, WebsocketError::Timeout));
3068 } else {
3069 panic!("Unexpected error type: {:?}", e);
3070 }
3071 }
3072 Ok(_) => panic!("Expected timeout error"),
3073 }
3074 });
3075 }
3076
3077 #[test]
3078 fn order_status_success() {
3079 TOKIO_SHARED_RT.block_on(async {
3080 let (ws_api, conn, mut rx) = setup().await;
3081 let client = AccountApiClient::new(ws_api.clone());
3082
3083 let handle = spawn(async move {
3084 let params = OrderStatusParams::builder("BNBUSDT".to_string(),).build().unwrap();
3085 client.order_status(params).await
3086 });
3087
3088 let sent = timeout(Duration::from_secs(1), rx.recv()).await.expect("send should occur").expect("channel closed");
3089 let Message::Text(text) = sent else { panic!() };
3090 let v: Value = serde_json::from_str(&text).unwrap();
3091 let id = v["id"].as_str().unwrap();
3092 assert_eq!(v["method"], "/order.status".trim_start_matches('/'));
3093 let mut resp_json: Value = serde_json::from_str(r#"{"id":"aa62318a-5a97-4f3b-bdc7-640bbe33b291","status":200,"result":{"symbol":"BTCUSDT","orderId":12569099453,"orderListId":-1,"clientOrderId":"4d96324ff9d44481926157","status":"FILLED","timeInForce":"GTC","type":"LIMIT","side":"SELL","trailingDelta":10,"trailingTime":-1,"time":1660801715639,"updateTime":1660801717945,"isWorking":true,"workingTime":1660801715639,"strategyId":37463720,"strategyType":1000000,"selfTradePreventionMode":"NONE","preventedMatchId":0,"usedSor":true,"workingFloor":"SOR","pegPriceType":"PRIMARY_PEG","pegOffsetType":"PRICE_LEVEL","pegOffsetValue":5,"peggedPrice":"87523.83710000","expiryReason":"INSUFFICIENT_LIQUIDITY"},"rateLimits":[{"rateLimitType":"REQUEST_WEIGHT","interval":"MINUTE","intervalNum":1,"limit":6000,"count":321}]}"#).unwrap_or_else(|_| serde_json::json!({}));
3094 resp_json["id"] = id.into();
3095
3096 let raw_data = resp_json.get("result").or_else(|| resp_json.get("response")).expect("no response in JSON");
3097 let expected_data: Box<models::OrderStatusResponseResult> = serde_json::from_value(raw_data.clone()).expect("should parse raw response");
3098 let empty_array = Value::Array(vec![]);
3099 let raw_rate_limits = resp_json.get("rateLimits").unwrap_or(&empty_array);
3100 let expected_rate_limits: Option<Vec<WebsocketApiRateLimit>> =
3101 match raw_rate_limits.as_array() {
3102 Some(arr) if arr.is_empty() => None,
3103 Some(_) => Some(serde_json::from_value(raw_rate_limits.clone()).expect("should parse rateLimits array")),
3104 None => None,
3105 };
3106
3107 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3108
3109 let response = timeout(Duration::from_secs(1), handle).await.expect("task done").expect("no panic").expect("no error");
3110
3111
3112 let response_rate_limits = response.rate_limits.clone();
3113 let response_data = response.data().expect("deserialize data");
3114
3115 assert_eq!(response_rate_limits, expected_rate_limits);
3116 assert_eq!(response_data, expected_data);
3117 });
3118 }
3119
3120 #[test]
3121 fn order_status_error_response() {
3122 TOKIO_SHARED_RT.block_on(async {
3123 let (ws_api, conn, mut rx) = setup().await;
3124 let client = AccountApiClient::new(ws_api.clone());
3125
3126 let handle = tokio::spawn(async move {
3127 let params = OrderStatusParams::builder("BNBUSDT".to_string(),).build().unwrap();
3128 client.order_status(params).await
3129 });
3130
3131 let sent = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
3132 let Message::Text(text) = sent else { panic!() };
3133 let v: Value = serde_json::from_str(&text).unwrap();
3134 let id = v["id"].as_str().unwrap().to_string();
3135
3136 let resp_json = json!({
3137 "id": id,
3138 "status": 400,
3139 "error": {
3140 "code": -2010,
3141 "msg": "Account has insufficient balance for requested action.",
3142 },
3143 "rateLimits": [
3144 {
3145 "rateLimitType": "ORDERS",
3146 "interval": "SECOND",
3147 "intervalNum": 10,
3148 "limit": 50,
3149 "count": 13
3150 },
3151 ],
3152 });
3153 WebsocketHandler::on_message(&*ws_api, resp_json.to_string(), conn.clone()).await;
3154
3155 let join = timeout(Duration::from_secs(1), handle).await.unwrap();
3156 match join {
3157 Ok(Err(e)) => {
3158 let msg = e.to_string();
3159 assert!(
3160 msg.contains("Server‐side response error (code -2010): Account has insufficient balance for requested action."),
3161 "Expected error msg to contain server error, got: {msg}"
3162 );
3163 }
3164 Ok(Ok(_)) => panic!("Expected error"),
3165 Err(_) => panic!("Task panicked"),
3166 }
3167 });
3168 }
3169
3170 #[test]
3171 fn order_status_request_timeout() {
3172 TOKIO_SHARED_RT.block_on(async {
3173 let (ws_api, _conn, mut rx) = setup().await;
3174 let client = AccountApiClient::new(ws_api.clone());
3175
3176 let handle = spawn(async move {
3177 let params = OrderStatusParams::builder("BNBUSDT".to_string())
3178 .build()
3179 .unwrap();
3180 client.order_status(params).await
3181 });
3182
3183 let sent = timeout(Duration::from_secs(1), rx.recv())
3184 .await
3185 .expect("send should occur")
3186 .expect("channel closed");
3187 let Message::Text(text) = sent else {
3188 panic!("expected Message Text")
3189 };
3190
3191 let _: Value = serde_json::from_str(&text).unwrap();
3192
3193 let result = handle.await.expect("task completed");
3194 match result {
3195 Err(e) => {
3196 if let Some(inner) = e.downcast_ref::<WebsocketError>() {
3197 assert!(matches!(inner, WebsocketError::Timeout));
3198 } else {
3199 panic!("Unexpected error type: {:?}", e);
3200 }
3201 }
3202 Ok(_) => panic!("Expected timeout error"),
3203 }
3204 });
3205 }
3206}