Skip to main content

nautilus_binance/futures/websocket/trading/
handler.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Binance Futures WebSocket Trading API message handler.
17//!
18//! The handler runs in a dedicated Tokio task as the I/O boundary between the client
19//! orchestrator and the network layer. It exclusively owns the `WebSocketClient` and
20//! processes commands from the client via an unbounded channel.
21//!
22//! Unlike the Spot handler which decodes SBE binary responses, the Futures handler
23//! works with JSON text responses throughout.
24
25use 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
53/// Binance Futures WebSocket Trading API handler.
54///
55/// Runs in a dedicated Tokio task, processing commands from the client
56/// and transforming raw WebSocket JSON messages into Nautilus domain events.
57pub 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    /// Creates a new handler instance.
82    #[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    /// Configures the receive window added before request signing.
103    #[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    /// Runs the main event loop for commands and raw messages.
110    ///
111    /// Returns `false` when disconnected or the signal is set.
112    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(&params)
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(&params)
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(&params)
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        // Sign over a key-sorted query string: Binance's WS API verifies the
272        // signature against the parameters sorted by key, which must not depend
273        // on serde_json's Map iteration order (issue #4410).
274        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                // A missing error payload carries no definitive command evidence
359                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}