binance_sdk/margin_trading/rest_api/apis/
user_data_stream_api.rs1#![allow(unused_imports)]
15use async_trait::async_trait;
16use derive_builder::Builder;
17use reqwest;
18use rust_decimal::prelude::*;
19use serde::{Deserialize, Serialize};
20use serde_json::{Value, json};
21use std::collections::BTreeMap;
22
23use crate::common::{
24 config::ConfigurationRestApi,
25 models::{ParamBuildError, RestApiResponse},
26 utils::send_request,
27};
28use crate::margin_trading::rest_api::models;
29
30const HAS_TIME_UNIT: bool = false;
31
32#[async_trait]
33pub trait UserDataStreamApi: Send + Sync {
34 async fn close_user_data_stream(&self) -> anyhow::Result<RestApiResponse<Value>>;
35 async fn keepalive_user_data_stream(
36 &self,
37 params: KeepaliveUserDataStreamParams,
38 ) -> anyhow::Result<RestApiResponse<Value>>;
39 async fn start_user_data_stream(
40 &self,
41 ) -> anyhow::Result<RestApiResponse<models::StartUserDataStreamResponse>>;
42}
43
44#[derive(Debug, Clone)]
45pub struct UserDataStreamApiClient {
46 configuration: ConfigurationRestApi,
47}
48
49impl UserDataStreamApiClient {
50 pub fn new(configuration: ConfigurationRestApi) -> Self {
51 Self { configuration }
52 }
53}
54
55#[derive(Clone, Debug, Builder, Deserialize)]
60#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
61pub struct KeepaliveUserDataStreamParams {
62 #[builder(setter(into))]
67 #[serde(rename = "listenKey")]
68 pub listen_key: String,
69}
70
71impl KeepaliveUserDataStreamParams {
72 #[must_use]
79 pub fn builder(listen_key: String) -> KeepaliveUserDataStreamParamsBuilder {
80 KeepaliveUserDataStreamParamsBuilder::default().listen_key(listen_key)
81 }
82}
83
84#[async_trait]
85impl UserDataStreamApi for UserDataStreamApiClient {
86 async fn close_user_data_stream(&self) -> anyhow::Result<RestApiResponse<Value>> {
87 let query_params = BTreeMap::new();
88 let body_params = BTreeMap::new();
89
90 send_request::<Value>(
91 &self.configuration,
92 "/sapi/v1/margin/listen-key",
93 reqwest::Method::DELETE,
94 query_params,
95 body_params,
96 if HAS_TIME_UNIT {
97 self.configuration.time_unit
98 } else {
99 None
100 },
101 false,
102 )
103 .await
104 }
105
106 async fn keepalive_user_data_stream(
107 &self,
108 params: KeepaliveUserDataStreamParams,
109 ) -> anyhow::Result<RestApiResponse<Value>> {
110 let KeepaliveUserDataStreamParams { listen_key } = params;
111
112 let mut query_params = BTreeMap::new();
113 let body_params = BTreeMap::new();
114
115 query_params.insert("listenKey".to_string(), json!(listen_key));
116
117 send_request::<Value>(
118 &self.configuration,
119 "/sapi/v1/margin/listen-key",
120 reqwest::Method::PUT,
121 query_params,
122 body_params,
123 if HAS_TIME_UNIT {
124 self.configuration.time_unit
125 } else {
126 None
127 },
128 false,
129 )
130 .await
131 }
132
133 async fn start_user_data_stream(
134 &self,
135 ) -> anyhow::Result<RestApiResponse<models::StartUserDataStreamResponse>> {
136 let query_params = BTreeMap::new();
137 let body_params = BTreeMap::new();
138
139 send_request::<models::StartUserDataStreamResponse>(
140 &self.configuration,
141 "/sapi/v1/margin/listen-key",
142 reqwest::Method::POST,
143 query_params,
144 body_params,
145 if HAS_TIME_UNIT {
146 self.configuration.time_unit
147 } else {
148 None
149 },
150 false,
151 )
152 .await
153 }
154}
155
156#[cfg(all(test, feature = "margin_trading"))]
157mod tests {
158 use super::*;
159 use crate::TOKIO_SHARED_RT;
160 use crate::{errors::ConnectorError, models::DataFuture, models::RestApiRateLimit};
161 use async_trait::async_trait;
162 use std::collections::HashMap;
163
164 struct DummyRestApiResponse<T> {
165 inner: Box<dyn FnOnce() -> DataFuture<Result<T, ConnectorError>> + Send + Sync>,
166 status: u16,
167 headers: HashMap<String, String>,
168 rate_limits: Option<Vec<RestApiRateLimit>>,
169 }
170
171 impl<T> From<DummyRestApiResponse<T>> for RestApiResponse<T> {
172 fn from(dummy: DummyRestApiResponse<T>) -> Self {
173 Self {
174 data_fn: dummy.inner,
175 status: dummy.status,
176 headers: dummy.headers,
177 rate_limits: dummy.rate_limits,
178 }
179 }
180 }
181
182 struct MockUserDataStreamApiClient {
183 force_error: bool,
184 }
185
186 #[async_trait]
187 impl UserDataStreamApi for MockUserDataStreamApiClient {
188 async fn close_user_data_stream(&self) -> anyhow::Result<RestApiResponse<Value>> {
189 if self.force_error {
190 return Err(ConnectorError::ConnectorClientError {
191 msg: "ResponseError".to_string(),
192 code: None,
193 }
194 .into());
195 }
196
197 let dummy_response = Value::Null;
198
199 let dummy = DummyRestApiResponse {
200 inner: Box::new(move || Box::pin(async move { Ok(dummy_response) })),
201 status: 200,
202 headers: HashMap::new(),
203 rate_limits: None,
204 };
205
206 Ok(dummy.into())
207 }
208
209 async fn keepalive_user_data_stream(
210 &self,
211 _params: KeepaliveUserDataStreamParams,
212 ) -> anyhow::Result<RestApiResponse<Value>> {
213 if self.force_error {
214 return Err(ConnectorError::ConnectorClientError {
215 msg: "ResponseError".to_string(),
216 code: None,
217 }
218 .into());
219 }
220
221 let dummy_response = Value::Null;
222
223 let dummy = DummyRestApiResponse {
224 inner: Box::new(move || Box::pin(async move { Ok(dummy_response) })),
225 status: 200,
226 headers: HashMap::new(),
227 rate_limits: None,
228 };
229
230 Ok(dummy.into())
231 }
232
233 async fn start_user_data_stream(
234 &self,
235 ) -> anyhow::Result<RestApiResponse<models::StartUserDataStreamResponse>> {
236 if self.force_error {
237 return Err(ConnectorError::ConnectorClientError {
238 msg: "ResponseError".to_string(),
239 code: None,
240 }
241 .into());
242 }
243
244 let resp_json: Value = serde_json::from_str(
245 r#"{"listenKey":"T3ee22BIYuWqmvne0HNq2A2WsFlEtLhvWCtItw6ffhhd"}"#,
246 )
247 .unwrap_or_else(|_| serde_json::json!({}));
248 let dummy_response: models::StartUserDataStreamResponse =
249 serde_json::from_value(resp_json.clone())
250 .expect("should parse into models::StartUserDataStreamResponse");
251
252 let dummy = DummyRestApiResponse {
253 inner: Box::new(move || Box::pin(async move { Ok(dummy_response) })),
254 status: 200,
255 headers: HashMap::new(),
256 rate_limits: None,
257 };
258
259 Ok(dummy.into())
260 }
261 }
262
263 #[test]
264 fn close_user_data_stream_required_params_success() {
265 TOKIO_SHARED_RT.block_on(async {
266 let client = MockUserDataStreamApiClient { force_error: false };
267
268 let expected_response = Value::Null;
269
270 let resp = client
271 .close_user_data_stream()
272 .await
273 .expect("Expected a response");
274 let data_future = resp.data();
275 let actual_response = data_future.await.unwrap();
276 assert_eq!(actual_response, expected_response);
277 });
278 }
279
280 #[test]
281 fn close_user_data_stream_optional_params_success() {
282 TOKIO_SHARED_RT.block_on(async {
283 let client = MockUserDataStreamApiClient { force_error: false };
284
285 let expected_response = Value::Null;
286
287 let resp = client
288 .close_user_data_stream()
289 .await
290 .expect("Expected a response");
291 let data_future = resp.data();
292 let actual_response = data_future.await.unwrap();
293 assert_eq!(actual_response, expected_response);
294 });
295 }
296
297 #[test]
298 fn close_user_data_stream_response_error() {
299 TOKIO_SHARED_RT.block_on(async {
300 let client = MockUserDataStreamApiClient { force_error: true };
301
302 match client.close_user_data_stream().await {
303 Ok(_) => panic!("Expected an error"),
304 Err(err) => {
305 assert_eq!(err.to_string(), "Connector client error: ResponseError");
306 }
307 }
308 });
309 }
310
311 #[test]
312 fn keepalive_user_data_stream_required_params_success() {
313 TOKIO_SHARED_RT.block_on(async {
314 let client = MockUserDataStreamApiClient { force_error: false };
315
316 let params = KeepaliveUserDataStreamParams::builder("listen_key_example".to_string())
317 .build()
318 .unwrap();
319
320 let expected_response = Value::Null;
321
322 let resp = client
323 .keepalive_user_data_stream(params)
324 .await
325 .expect("Expected a response");
326 let data_future = resp.data();
327 let actual_response = data_future.await.unwrap();
328 assert_eq!(actual_response, expected_response);
329 });
330 }
331
332 #[test]
333 fn keepalive_user_data_stream_optional_params_success() {
334 TOKIO_SHARED_RT.block_on(async {
335 let client = MockUserDataStreamApiClient { force_error: false };
336
337 let params = KeepaliveUserDataStreamParams::builder("listen_key_example".to_string())
338 .build()
339 .unwrap();
340
341 let expected_response = Value::Null;
342
343 let resp = client
344 .keepalive_user_data_stream(params)
345 .await
346 .expect("Expected a response");
347 let data_future = resp.data();
348 let actual_response = data_future.await.unwrap();
349 assert_eq!(actual_response, expected_response);
350 });
351 }
352
353 #[test]
354 fn keepalive_user_data_stream_response_error() {
355 TOKIO_SHARED_RT.block_on(async {
356 let client = MockUserDataStreamApiClient { force_error: true };
357
358 let params = KeepaliveUserDataStreamParams::builder("listen_key_example".to_string())
359 .build()
360 .unwrap();
361
362 match client.keepalive_user_data_stream(params).await {
363 Ok(_) => panic!("Expected an error"),
364 Err(err) => {
365 assert_eq!(err.to_string(), "Connector client error: ResponseError");
366 }
367 }
368 });
369 }
370
371 #[test]
372 fn start_user_data_stream_required_params_success() {
373 TOKIO_SHARED_RT.block_on(async {
374 let client = MockUserDataStreamApiClient { force_error: false };
375
376 let resp_json: Value = serde_json::from_str(
377 r#"{"listenKey":"T3ee22BIYuWqmvne0HNq2A2WsFlEtLhvWCtItw6ffhhd"}"#,
378 )
379 .unwrap_or_else(|_| serde_json::json!({}));
380 let expected_response: models::StartUserDataStreamResponse =
381 serde_json::from_value(resp_json.clone())
382 .expect("should parse into models::StartUserDataStreamResponse");
383
384 let resp = client
385 .start_user_data_stream()
386 .await
387 .expect("Expected a response");
388 let data_future = resp.data();
389 let actual_response = data_future.await.unwrap();
390 assert_eq!(actual_response, expected_response);
391 });
392 }
393
394 #[test]
395 fn start_user_data_stream_optional_params_success() {
396 TOKIO_SHARED_RT.block_on(async {
397 let client = MockUserDataStreamApiClient { force_error: false };
398
399 let resp_json: Value = serde_json::from_str(
400 r#"{"listenKey":"T3ee22BIYuWqmvne0HNq2A2WsFlEtLhvWCtItw6ffhhd"}"#,
401 )
402 .unwrap_or_else(|_| serde_json::json!({}));
403 let expected_response: models::StartUserDataStreamResponse =
404 serde_json::from_value(resp_json.clone())
405 .expect("should parse into models::StartUserDataStreamResponse");
406
407 let resp = client
408 .start_user_data_stream()
409 .await
410 .expect("Expected a response");
411 let data_future = resp.data();
412 let actual_response = data_future.await.unwrap();
413 assert_eq!(actual_response, expected_response);
414 });
415 }
416
417 #[test]
418 fn start_user_data_stream_response_error() {
419 TOKIO_SHARED_RT.block_on(async {
420 let client = MockUserDataStreamApiClient { force_error: true };
421
422 match client.start_user_data_stream().await {
423 Ok(_) => panic!("Expected an error"),
424 Err(err) => {
425 assert_eq!(err.to_string(), "Connector client error: ResponseError");
426 }
427 }
428 });
429 }
430}