nautilus_binance/futures/websocket/trading/
handler.rs1use std::{
26 fmt::Debug,
27 sync::{
28 Arc,
29 atomic::{AtomicBool, Ordering},
30 },
31};
32
33use ahash::AHashMap;
34use nautilus_network::{RECONNECTED, websocket::WebSocketClient};
35use tokio_tungstenite::tungstenite::Message;
36
37use super::{
38 client::BINANCE_FUTURES_WS_RATE_LIMIT_KEY_ORDER,
39 error::{BinanceFuturesWsApiError, BinanceFuturesWsApiResult},
40 messages::{
41 BinanceFuturesWsTradingCommand, BinanceFuturesWsTradingMessage,
42 BinanceFuturesWsTradingRequest, BinanceFuturesWsTradingRequestMeta,
43 BinanceFuturesWsTradingResponse, method,
44 },
45};
46use crate::{
47 common::credential::{SigningCredential, canonical_ws_query_string},
48 futures::http::query::{
49 BinanceCancelOrderParams, BinanceModifyOrderParams, BinanceNewOrderParams,
50 },
51};
52
53pub struct BinanceFuturesWsTradingHandler {
58 signal: Arc<AtomicBool>,
59 inner: Option<WebSocketClient>,
60 cmd_rx: tokio::sync::mpsc::UnboundedReceiver<BinanceFuturesWsTradingCommand>,
61 raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
62 out_tx: tokio::sync::mpsc::UnboundedSender<BinanceFuturesWsTradingMessage>,
63 credential: Arc<SigningCredential>,
64 pending_requests: AHashMap<String, BinanceFuturesWsTradingRequestMeta>,
65 recv_window_ms: Option<u64>,
66}
67
68impl Debug for BinanceFuturesWsTradingHandler {
69 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
70 f.debug_struct(stringify!(BinanceFuturesWsTradingHandler))
71 .field("inner", &self.inner.as_ref().map(|_| "<client>"))
72 .field(
73 "pending_requests",
74 &format!("{} pending", self.pending_requests.len()),
75 )
76 .finish_non_exhaustive()
77 }
78}
79
80impl BinanceFuturesWsTradingHandler {
81 #[must_use]
83 pub fn new(
84 signal: Arc<AtomicBool>,
85 cmd_rx: tokio::sync::mpsc::UnboundedReceiver<BinanceFuturesWsTradingCommand>,
86 raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
87 out_tx: tokio::sync::mpsc::UnboundedSender<BinanceFuturesWsTradingMessage>,
88 credential: Arc<SigningCredential>,
89 ) -> Self {
90 Self {
91 signal,
92 inner: None,
93 cmd_rx,
94 raw_rx,
95 out_tx,
96 credential,
97 pending_requests: AHashMap::new(),
98 recv_window_ms: None,
99 }
100 }
101
102 #[must_use]
104 pub const fn with_recv_window(mut self, recv_window_ms: Option<u64>) -> Self {
105 self.recv_window_ms = recv_window_ms;
106 self
107 }
108
109 pub async fn run(&mut self) -> bool {
113 loop {
114 if self.signal.load(Ordering::Relaxed) {
115 return false;
116 }
117
118 tokio::select! {
119 Some(cmd) = self.cmd_rx.recv() => {
120 match cmd {
121 BinanceFuturesWsTradingCommand::SetClient(client) => {
122 log::debug!("Handler received WebSocket client");
123 self.inner = Some(client);
124 self.emit(BinanceFuturesWsTradingMessage::Connected);
125 }
126 BinanceFuturesWsTradingCommand::Disconnect => {
127 log::debug!("Handler disconnecting WebSocket client");
128 self.inner = None;
129 return false;
130 }
131 BinanceFuturesWsTradingCommand::PlaceOrder { id, params } => {
132 if let Err(e) = self.handle_place_order(id.clone(), params).await {
133 log::error!("Failed to handle place order command: {e}");
134 self.pending_requests.remove(&id);
135 self.emit(BinanceFuturesWsTradingMessage::RequestFailed {
136 request_id: id,
137 msg: e.to_string(),
138 });
139 }
140 }
141 BinanceFuturesWsTradingCommand::CancelOrder { id, params } => {
142 if let Err(e) = self.handle_cancel_order(id.clone(), params).await {
143 log::error!("Failed to handle cancel order command: {e}");
144 self.pending_requests.remove(&id);
145 self.emit(BinanceFuturesWsTradingMessage::RequestFailed {
146 request_id: id,
147 msg: e.to_string(),
148 });
149 }
150 }
151 BinanceFuturesWsTradingCommand::ModifyOrder { id, params } => {
152 if let Err(e) = self.handle_modify_order(id.clone(), params).await {
153 log::error!("Failed to handle modify order command: {e}");
154 self.pending_requests.remove(&id);
155 self.emit(BinanceFuturesWsTradingMessage::RequestFailed {
156 request_id: id,
157 msg: e.to_string(),
158 });
159 }
160 }
161 }
162 }
163 Some(msg) = self.raw_rx.recv() => {
164 if let Message::Text(ref text) = msg
165 && text.as_str() == RECONNECTED
166 {
167 log::debug!("Handler received reconnection signal");
168 self.fail_pending_requests();
169 self.emit(BinanceFuturesWsTradingMessage::Reconnected);
170 continue;
171 }
172
173 self.handle_message(msg);
174 }
175 else => {
176 return false;
177 }
178 }
179 }
180 }
181
182 fn emit(&self, msg: BinanceFuturesWsTradingMessage) {
183 if let Err(e) = self.out_tx.send(msg) {
184 log::error!("Failed to send message to output channel: {e}");
185 }
186 }
187
188 fn fail_pending_requests(&mut self) {
189 if self.pending_requests.is_empty() {
190 return;
191 }
192
193 let count = self.pending_requests.len();
194 log::warn!("Failing {count} pending requests after reconnection");
195
196 let pending = std::mem::take(&mut self.pending_requests);
197 for (request_id, _meta) in pending {
198 self.emit(BinanceFuturesWsTradingMessage::RequestFailed {
199 request_id,
200 msg: "Connection lost before response received".to_string(),
201 });
202 }
203 }
204
205 async fn handle_place_order(
206 &mut self,
207 id: String,
208 params: BinanceNewOrderParams,
209 ) -> BinanceFuturesWsApiResult<()> {
210 let params_json = serde_json::to_value(¶ms)
211 .map_err(|e| BinanceFuturesWsApiError::JsonError(e.to_string()))?;
212 let signed_params = self.sign_params(params_json)?;
213
214 let request = BinanceFuturesWsTradingRequest::new(&id, method::ORDER_PLACE, signed_params);
215 self.pending_requests
216 .insert(id.clone(), BinanceFuturesWsTradingRequestMeta::PlaceOrder);
217 self.send_request(request).await
218 }
219
220 async fn handle_cancel_order(
221 &mut self,
222 id: String,
223 params: BinanceCancelOrderParams,
224 ) -> BinanceFuturesWsApiResult<()> {
225 let params_json = serde_json::to_value(¶ms)
226 .map_err(|e| BinanceFuturesWsApiError::JsonError(e.to_string()))?;
227 let signed_params = self.sign_params(params_json)?;
228
229 let request = BinanceFuturesWsTradingRequest::new(&id, method::ORDER_CANCEL, signed_params);
230 self.pending_requests
231 .insert(id.clone(), BinanceFuturesWsTradingRequestMeta::CancelOrder);
232 self.send_request(request).await
233 }
234
235 async fn handle_modify_order(
236 &mut self,
237 id: String,
238 params: BinanceModifyOrderParams,
239 ) -> BinanceFuturesWsApiResult<()> {
240 let params_json = serde_json::to_value(¶ms)
241 .map_err(|e| BinanceFuturesWsApiError::JsonError(e.to_string()))?;
242 let signed_params = self.sign_params(params_json)?;
243
244 let request = BinanceFuturesWsTradingRequest::new(&id, method::ORDER_MODIFY, signed_params);
245 self.pending_requests
246 .insert(id.clone(), BinanceFuturesWsTradingRequestMeta::ModifyOrder);
247 self.send_request(request).await
248 }
249
250 fn sign_params(
251 &self,
252 mut params: serde_json::Value,
253 ) -> BinanceFuturesWsApiResult<serde_json::Value> {
254 let timestamp = std::time::SystemTime::now()
255 .duration_since(std::time::UNIX_EPOCH)
256 .map_err(|e| BinanceFuturesWsApiError::ClientError(e.to_string()))?
257 .as_millis() as i64;
258
259 if let Some(obj) = params.as_object_mut() {
260 obj.insert("timestamp".to_string(), serde_json::json!(timestamp));
261 obj.insert(
262 "apiKey".to_string(),
263 serde_json::json!(self.credential.api_key()),
264 );
265
266 if let Some(recv_window_ms) = self.recv_window_ms {
267 obj.insert("recvWindow".to_string(), serde_json::json!(recv_window_ms));
268 }
269 }
270
271 let query_string = canonical_ws_query_string(
275 params
276 .as_object()
277 .into_iter()
278 .flatten()
279 .map(|(k, v)| (k.as_str(), v)),
280 )
281 .map_err(|e| BinanceFuturesWsApiError::ClientError(e.to_string()))?;
282 let signature = self.credential.sign(&query_string);
283
284 if let Some(obj) = params.as_object_mut() {
285 obj.insert("signature".to_string(), serde_json::json!(signature));
286 }
287
288 Ok(params)
289 }
290
291 async fn send_request(
292 &mut self,
293 request: BinanceFuturesWsTradingRequest,
294 ) -> BinanceFuturesWsApiResult<()> {
295 let client = self.inner.as_mut().ok_or_else(|| {
296 BinanceFuturesWsApiError::ConnectionError("WebSocket not connected".to_string())
297 })?;
298
299 let json = serde_json::to_string(&request)
300 .map_err(|e| BinanceFuturesWsApiError::JsonError(e.to_string()))?;
301
302 log::debug!(
303 "Sending Futures WS Trading API request id={} method={}",
304 request.id,
305 request.method
306 );
307
308 client
309 .send_text(
310 json,
311 Some(BINANCE_FUTURES_WS_RATE_LIMIT_KEY_ORDER.as_slice()),
312 )
313 .await
314 .map_err(|e| {
315 BinanceFuturesWsApiError::ConnectionError(format!("Failed to send request: {e}"))
316 })?;
317
318 Ok(())
319 }
320
321 fn handle_message(&mut self, msg: Message) {
322 match msg {
323 Message::Text(text) => self.handle_text_response(&text),
324 Message::Ping(_) | Message::Pong(_) => {}
325 Message::Close(frame) => {
326 log::debug!("WebSocket closed: {frame:?}");
327 }
328 Message::Binary(_) | Message::Frame(_) => {}
329 }
330 }
331
332 fn handle_text_response(&mut self, text: &str) {
333 let response: BinanceFuturesWsTradingResponse = match serde_json::from_str(text) {
334 Ok(r) => r,
335 Err(e) => {
336 log::warn!("Failed to parse WS Trading API response: {e}");
337 return;
338 }
339 };
340
341 let Some(meta) = self.pending_requests.remove(&response.id) else {
342 log::warn!("Received response for unknown request ID: {}", response.id);
343 return;
344 };
345
346 if response.status != 200 {
347 match response.error {
348 Some(error) => {
349 let rejection = self.create_rejection(
350 response.id,
351 response.status,
352 error.code,
353 error.msg,
354 meta,
355 );
356 self.emit(rejection);
357 }
358 None => {
360 self.emit(BinanceFuturesWsTradingMessage::RequestFailed {
361 request_id: response.id,
362 msg: format!(
363 "Request failed with status {}; error payload missing",
364 response.status
365 ),
366 });
367 }
368 }
369 return;
370 }
371
372 let Some(result) = response.result else {
373 log::warn!(
374 "Missing result in success response for request {}",
375 response.id
376 );
377 return;
378 };
379
380 match meta {
381 BinanceFuturesWsTradingRequestMeta::PlaceOrder => {
382 match serde_json::from_value(result) {
383 Ok(order) => {
384 self.emit(BinanceFuturesWsTradingMessage::OrderAccepted {
385 request_id: response.id,
386 response: Box::new(order),
387 });
388 }
389 Err(e) => {
390 log::error!("Failed to deserialize order response: {e}");
391 self.emit(BinanceFuturesWsTradingMessage::Error(e.to_string()));
392 }
393 }
394 }
395 BinanceFuturesWsTradingRequestMeta::CancelOrder => match serde_json::from_value(result)
396 {
397 Ok(order) => {
398 self.emit(BinanceFuturesWsTradingMessage::OrderCanceled {
399 request_id: response.id,
400 response: Box::new(order),
401 });
402 }
403 Err(e) => {
404 log::error!("Failed to deserialize cancel response: {e}");
405 self.emit(BinanceFuturesWsTradingMessage::Error(e.to_string()));
406 }
407 },
408 BinanceFuturesWsTradingRequestMeta::ModifyOrder => match serde_json::from_value(result)
409 {
410 Ok(order) => {
411 self.emit(BinanceFuturesWsTradingMessage::OrderModified {
412 request_id: response.id,
413 response: Box::new(order),
414 });
415 }
416 Err(e) => {
417 log::error!("Failed to deserialize modify response: {e}");
418 self.emit(BinanceFuturesWsTradingMessage::Error(e.to_string()));
419 }
420 },
421 }
422 }
423
424 fn create_rejection(
425 &self,
426 request_id: String,
427 status: u16,
428 code: i32,
429 msg: String,
430 meta: BinanceFuturesWsTradingRequestMeta,
431 ) -> BinanceFuturesWsTradingMessage {
432 match meta {
433 BinanceFuturesWsTradingRequestMeta::PlaceOrder => {
434 BinanceFuturesWsTradingMessage::OrderRejected {
435 request_id,
436 status,
437 code,
438 msg,
439 }
440 }
441 BinanceFuturesWsTradingRequestMeta::CancelOrder => {
442 BinanceFuturesWsTradingMessage::CancelRejected {
443 request_id,
444 status,
445 code,
446 msg,
447 }
448 }
449 BinanceFuturesWsTradingRequestMeta::ModifyOrder => {
450 BinanceFuturesWsTradingMessage::ModifyRejected {
451 request_id,
452 status,
453 code,
454 msg,
455 }
456 }
457 }
458 }
459}
460
461#[cfg(test)]
462mod tests {
463 use std::sync::atomic::AtomicBool;
464
465 use rstest::rstest;
466
467 use super::*;
468
469 #[rstest]
470 fn test_sign_params_includes_recv_window_in_signature() {
471 let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
472 let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
473 let (out_tx, _out_rx) = tokio::sync::mpsc::unbounded_channel();
474 let credential = Arc::new(SigningCredential::new(
475 "api-key".to_string(),
476 "hmac-secret".to_string(),
477 ));
478 let handler = BinanceFuturesWsTradingHandler::new(
479 Arc::new(AtomicBool::new(false)),
480 cmd_rx,
481 raw_rx,
482 out_tx,
483 credential.clone(),
484 )
485 .with_recv_window(Some(30_000));
486
487 let signed = handler
488 .sign_params(serde_json::json!({"symbol": "BTCUSDT"}))
489 .unwrap();
490 let mut unsigned = signed.clone();
491 let signature = unsigned
492 .as_object_mut()
493 .unwrap()
494 .remove("signature")
495 .unwrap();
496 let query = canonical_ws_query_string(
497 unsigned
498 .as_object()
499 .unwrap()
500 .iter()
501 .map(|(key, value)| (key.as_str(), value)),
502 )
503 .unwrap();
504
505 assert_eq!(signed["recvWindow"], 30_000);
506 assert_eq!(signature, credential.sign(&query));
507 }
508}