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 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}