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            let (code, msg) = response.error.map(|e| (e.code, e.msg)).unwrap_or((
348                -1,
349                format!("Request failed with status {}", response.status),
350            ));
351            let rejection = self.create_rejection(response.id, code, msg, meta);
352            self.emit(rejection);
353            return;
354        }
355
356        let Some(result) = response.result else {
357            log::warn!(
358                "Missing result in success response for request {}",
359                response.id
360            );
361            return;
362        };
363
364        match meta {
365            BinanceFuturesWsTradingRequestMeta::PlaceOrder => {
366                match serde_json::from_value(result) {
367                    Ok(order) => {
368                        self.emit(BinanceFuturesWsTradingMessage::OrderAccepted {
369                            request_id: response.id,
370                            response: Box::new(order),
371                        });
372                    }
373                    Err(e) => {
374                        log::error!("Failed to deserialize order response: {e}");
375                        self.emit(BinanceFuturesWsTradingMessage::Error(e.to_string()));
376                    }
377                }
378            }
379            BinanceFuturesWsTradingRequestMeta::CancelOrder => match serde_json::from_value(result)
380            {
381                Ok(order) => {
382                    self.emit(BinanceFuturesWsTradingMessage::OrderCanceled {
383                        request_id: response.id,
384                        response: Box::new(order),
385                    });
386                }
387                Err(e) => {
388                    log::error!("Failed to deserialize cancel response: {e}");
389                    self.emit(BinanceFuturesWsTradingMessage::Error(e.to_string()));
390                }
391            },
392            BinanceFuturesWsTradingRequestMeta::ModifyOrder => match serde_json::from_value(result)
393            {
394                Ok(order) => {
395                    self.emit(BinanceFuturesWsTradingMessage::OrderModified {
396                        request_id: response.id,
397                        response: Box::new(order),
398                    });
399                }
400                Err(e) => {
401                    log::error!("Failed to deserialize modify response: {e}");
402                    self.emit(BinanceFuturesWsTradingMessage::Error(e.to_string()));
403                }
404            },
405        }
406    }
407
408    fn create_rejection(
409        &self,
410        request_id: String,
411        code: i32,
412        msg: String,
413        meta: BinanceFuturesWsTradingRequestMeta,
414    ) -> BinanceFuturesWsTradingMessage {
415        match meta {
416            BinanceFuturesWsTradingRequestMeta::PlaceOrder => {
417                BinanceFuturesWsTradingMessage::OrderRejected {
418                    request_id,
419                    code,
420                    msg,
421                }
422            }
423            BinanceFuturesWsTradingRequestMeta::CancelOrder => {
424                BinanceFuturesWsTradingMessage::CancelRejected {
425                    request_id,
426                    code,
427                    msg,
428                }
429            }
430            BinanceFuturesWsTradingRequestMeta::ModifyOrder => {
431                BinanceFuturesWsTradingMessage::ModifyRejected {
432                    request_id,
433                    code,
434                    msg,
435                }
436            }
437        }
438    }
439}
440
441#[cfg(test)]
442mod tests {
443    use std::sync::atomic::AtomicBool;
444
445    use rstest::rstest;
446
447    use super::*;
448
449    #[rstest]
450    fn test_sign_params_includes_recv_window_in_signature() {
451        let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
452        let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
453        let (out_tx, _out_rx) = tokio::sync::mpsc::unbounded_channel();
454        let credential = Arc::new(SigningCredential::new(
455            "api-key".to_string(),
456            "hmac-secret".to_string(),
457        ));
458        let handler = BinanceFuturesWsTradingHandler::new(
459            Arc::new(AtomicBool::new(false)),
460            cmd_rx,
461            raw_rx,
462            out_tx,
463            credential.clone(),
464        )
465        .with_recv_window(Some(30_000));
466
467        let signed = handler
468            .sign_params(serde_json::json!({"symbol": "BTCUSDT"}))
469            .unwrap();
470        let mut unsigned = signed.clone();
471        let signature = unsigned
472            .as_object_mut()
473            .unwrap()
474            .remove("signature")
475            .unwrap();
476        let query = canonical_ws_query_string(
477            unsigned
478                .as_object()
479                .unwrap()
480                .iter()
481                .map(|(key, value)| (key.as_str(), value)),
482        )
483        .unwrap();
484
485        assert_eq!(signed["recvWindow"], 30_000);
486        assert_eq!(signature, credential.sign(&query));
487    }
488}