Skip to main content

nautilus_binance/spot/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 Spot WebSocket 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//! ## Responsibilities
23//!
24//! - Command processing: Receives `BinanceSpotWsTradingCommand` from client, serializes to JSON requests.
25//! - Response decoding: Parses SBE binary responses using schema 3 decoders.
26//! - Request correlation: Matches responses to pending requests by ID.
27//! - Message transformation: Emits `BinanceSpotWsTradingMessage` events to client via channel.
28
29use std::{
30    fmt::Debug,
31    sync::{
32        Arc,
33        atomic::{AtomicBool, AtomicU64, Ordering},
34    },
35};
36
37use ahash::AHashMap;
38use nautilus_network::{RECONNECTED, websocket::WebSocketClient};
39use tokio_tungstenite::tungstenite::Message;
40
41use super::{
42    client::BINANCE_WS_RATE_LIMIT_KEY_ORDER,
43    error::{BinanceWsApiError, BinanceWsApiResult},
44    messages::{
45        BinanceSpotWsTradingCommand, BinanceSpotWsTradingMessage, BinanceSpotWsTradingRequest,
46        BinanceSpotWsTradingRequestMeta, method,
47    },
48};
49use crate::{
50    common::credential::{SigningCredential, canonical_ws_query_string},
51    spot::{
52        enums::BinanceSpotUserDataEventType,
53        http::{models::BinanceCancelOrderResponse, parse},
54        sbe::spot::{
55            ReadBuf,
56            error_response_codec::ErrorResponseDecoder,
57            message_header_codec,
58            web_socket_response_codec::{SBE_TEMPLATE_ID, WebSocketResponseDecoder},
59        },
60    },
61};
62
63/// Binance Spot WebSocket API handler.
64///
65/// Runs in a dedicated Tokio task, processing commands from the client
66/// and transforming raw WebSocket messages into Nautilus domain events.
67/// Messages are sent to the client via the output channel.
68pub struct BinanceSpotWsTradingHandler {
69    signal: Arc<AtomicBool>,
70    inner: Option<WebSocketClient>,
71    cmd_rx: tokio::sync::mpsc::UnboundedReceiver<BinanceSpotWsTradingCommand>,
72    raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
73    out_tx: tokio::sync::mpsc::UnboundedSender<BinanceSpotWsTradingMessage>,
74    credential: Arc<SigningCredential>,
75    pending_requests: AHashMap<String, BinanceSpotWsTradingRequestMeta>,
76    request_id_counter: AtomicU64,
77    recv_window_ms: Option<u64>,
78}
79
80impl Debug for BinanceSpotWsTradingHandler {
81    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
82        f.debug_struct(stringify!(BinanceSpotWsTradingHandler))
83            .field("inner", &self.inner.as_ref().map(|_| "<client>"))
84            .field(
85                "pending_requests",
86                &format!("{} pending", self.pending_requests.len()),
87            )
88            .finish_non_exhaustive()
89    }
90}
91
92impl BinanceSpotWsTradingHandler {
93    /// Creates a new handler instance.
94    #[must_use]
95    pub fn new(
96        signal: Arc<AtomicBool>,
97        cmd_rx: tokio::sync::mpsc::UnboundedReceiver<BinanceSpotWsTradingCommand>,
98        raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
99        out_tx: tokio::sync::mpsc::UnboundedSender<BinanceSpotWsTradingMessage>,
100        credential: Arc<SigningCredential>,
101    ) -> Self {
102        Self {
103            signal,
104            inner: None,
105            cmd_rx,
106            raw_rx,
107            out_tx,
108            credential,
109            pending_requests: AHashMap::new(),
110            request_id_counter: AtomicU64::new(1000),
111            recv_window_ms: None,
112        }
113    }
114
115    /// Configures the receive window added before request signing.
116    #[must_use]
117    pub const fn with_recv_window(mut self, recv_window_ms: Option<u64>) -> Self {
118        self.recv_window_ms = recv_window_ms;
119        self
120    }
121
122    /// Runs the main event loop for commands and raw messages.
123    ///
124    /// Sends output messages via `out_tx` channel. Returns `false` when disconnected
125    /// or the signal is set, indicating the handler should exit.
126    pub async fn run(&mut self) -> bool {
127        loop {
128            if self.signal.load(Ordering::Relaxed) {
129                return false;
130            }
131
132            tokio::select! {
133                Some(cmd) = self.cmd_rx.recv() => {
134                    match cmd {
135                        BinanceSpotWsTradingCommand::SetClient(client) => {
136                            log::debug!("Handler received WebSocket client");
137                            self.inner = Some(client);
138                            self.emit(BinanceSpotWsTradingMessage::Connected);
139                        }
140                        BinanceSpotWsTradingCommand::Disconnect => {
141                            log::debug!("Handler disconnecting WebSocket client");
142                            self.inner = None;
143                            return false;
144                        }
145                        BinanceSpotWsTradingCommand::PlaceOrder { id, params } => {
146                            if let Err(e) = self.handle_place_order(id.clone(), params).await {
147                                log::error!("Failed to handle place order command: {e}");
148                                self.pending_requests.remove(&id);
149                                self.emit(BinanceSpotWsTradingMessage::RequestFailed {
150                                    request_id: id,
151                                    msg: e.to_string(),
152                                });
153                            }
154                        }
155                        BinanceSpotWsTradingCommand::CancelOrder { id, params } => {
156                            if let Err(e) = self.handle_cancel_order(id.clone(), params).await {
157                                log::error!("Failed to handle cancel order command: {e}");
158                                self.pending_requests.remove(&id);
159                                self.emit(BinanceSpotWsTradingMessage::RequestFailed {
160                                    request_id: id,
161                                    msg: e.to_string(),
162                                });
163                            }
164                        }
165                        BinanceSpotWsTradingCommand::CancelReplaceOrder { id, params } => {
166                            if let Err(e) = self.handle_cancel_replace_order(id.clone(), params).await {
167                                log::error!("Failed to handle cancel replace command: {e}");
168                                self.pending_requests.remove(&id);
169                                self.emit(BinanceSpotWsTradingMessage::RequestFailed {
170                                    request_id: id,
171                                    msg: e.to_string(),
172                                });
173                            }
174                        }
175                        BinanceSpotWsTradingCommand::CancelAllOrders { id, symbol } => {
176                            if let Err(e) = self.handle_cancel_all_orders(id.clone(), symbol).await {
177                                log::error!("Failed to handle cancel all command: {e}");
178                                self.pending_requests.remove(&id);
179                                self.emit(BinanceSpotWsTradingMessage::RequestFailed {
180                                    request_id: id,
181                                    msg: e.to_string(),
182                                });
183                            }
184                        }
185                        BinanceSpotWsTradingCommand::SessionLogon => {
186                            if let Err(e) = self.handle_session_logon().await {
187                                log::error!("Session logon failed: {e}");
188                                self.emit(BinanceSpotWsTradingMessage::AuthenticationRejected(
189                                    format!("Session logon failed: {e}"),
190                                ));
191                            }
192                        }
193                        BinanceSpotWsTradingCommand::SubscribeUserData => {
194                            if let Err(e) = self.handle_subscribe_user_data().await {
195                                log::error!("User data subscribe failed: {e}");
196                                self.emit(
197                                    BinanceSpotWsTradingMessage::UserDataSubscriptionRejected(
198                                        format!("User data subscribe failed: {e}"),
199                                    ),
200                                );
201                            }
202                        }
203                    }
204                }
205                Some(msg) = self.raw_rx.recv() => {
206                    if let Message::Text(ref text) = msg
207                        && text.as_str() == RECONNECTED
208                    {
209                        log::debug!("Handler received reconnection signal");
210
211                        // Fail any pending requests - they won't get responses on new connection
212                        self.fail_pending_requests();
213
214                        self.emit(BinanceSpotWsTradingMessage::Reconnected);
215                        continue;
216                    }
217
218                    self.handle_message(msg);
219                }
220                else => {
221                    // Both channels closed
222                    return false;
223                }
224            }
225        }
226    }
227
228    /// Sends a message to the output channel.
229    fn emit(&self, msg: BinanceSpotWsTradingMessage) {
230        if let Err(e) = self.out_tx.send(msg) {
231            log::error!("Failed to send message to output channel: {e}");
232        }
233    }
234
235    /// Fails all pending requests after a reconnection.
236    fn fail_pending_requests(&mut self) {
237        if self.pending_requests.is_empty() {
238            return;
239        }
240
241        let count = self.pending_requests.len();
242        log::warn!("Failing {count} pending requests after reconnection");
243
244        let pending = std::mem::take(&mut self.pending_requests);
245        for (request_id, _meta) in pending {
246            self.emit(BinanceSpotWsTradingMessage::RequestFailed {
247                request_id,
248                msg: "Connection lost before response received".to_string(),
249            });
250        }
251    }
252
253    async fn handle_place_order(
254        &mut self,
255        id: String,
256        params: crate::spot::http::query::NewOrderParams,
257    ) -> BinanceWsApiResult<()> {
258        let params_json = serde_json::to_value(&params)
259            .map_err(|e| BinanceWsApiError::ClientError(e.to_string()))?;
260        let signed_params = self.sign_params(params_json)?;
261
262        let request = BinanceSpotWsTradingRequest::new(&id, method::ORDER_PLACE, signed_params);
263        self.pending_requests
264            .insert(id.clone(), BinanceSpotWsTradingRequestMeta::PlaceOrder);
265        self.send_request(request).await
266    }
267
268    async fn handle_cancel_order(
269        &mut self,
270        id: String,
271        params: crate::spot::http::query::CancelOrderParams,
272    ) -> BinanceWsApiResult<()> {
273        let params_json = serde_json::to_value(&params)
274            .map_err(|e| BinanceWsApiError::ClientError(e.to_string()))?;
275        let signed_params = self.sign_params(params_json)?;
276
277        let request = BinanceSpotWsTradingRequest::new(&id, method::ORDER_CANCEL, signed_params);
278        self.pending_requests
279            .insert(id.clone(), BinanceSpotWsTradingRequestMeta::CancelOrder);
280        self.send_request(request).await
281    }
282
283    async fn handle_cancel_replace_order(
284        &mut self,
285        id: String,
286        params: crate::spot::http::query::CancelReplaceOrderParams,
287    ) -> BinanceWsApiResult<()> {
288        let params_json = serde_json::to_value(&params)
289            .map_err(|e| BinanceWsApiError::ClientError(e.to_string()))?;
290        let signed_params = self.sign_params(params_json)?;
291
292        let request =
293            BinanceSpotWsTradingRequest::new(&id, method::ORDER_CANCEL_REPLACE, signed_params);
294        self.pending_requests.insert(
295            id.clone(),
296            BinanceSpotWsTradingRequestMeta::CancelReplaceOrder,
297        );
298        self.send_request(request).await
299    }
300
301    async fn handle_cancel_all_orders(
302        &mut self,
303        id: String,
304        symbol: String,
305    ) -> BinanceWsApiResult<()> {
306        let params_json = serde_json::json!({ "symbol": symbol });
307        let signed_params = self.sign_params(params_json)?;
308
309        let request =
310            BinanceSpotWsTradingRequest::new(&id, method::OPEN_ORDERS_CANCEL_ALL, signed_params);
311        self.pending_requests
312            .insert(id.clone(), BinanceSpotWsTradingRequestMeta::CancelAllOrders);
313        self.send_request(request).await
314    }
315
316    async fn handle_session_logon(&mut self) -> BinanceWsApiResult<()> {
317        let id = self.next_request_id();
318        let params_json = serde_json::json!({});
319        let signed_params = self.sign_params(params_json)?;
320
321        let request = BinanceSpotWsTradingRequest::new(&id, "session.logon", signed_params);
322        self.pending_requests
323            .insert(id, BinanceSpotWsTradingRequestMeta::SessionLogon);
324        self.send_request(request).await
325    }
326
327    async fn handle_subscribe_user_data(&mut self) -> BinanceWsApiResult<()> {
328        let id = self.next_request_id();
329        let request = BinanceSpotWsTradingRequest::new(
330            &id,
331            "userDataStream.subscribe",
332            serde_json::json!({}),
333        );
334        self.pending_requests
335            .insert(id, BinanceSpotWsTradingRequestMeta::SubscribeUserData);
336        self.send_request(request).await
337    }
338
339    fn next_request_id(&self) -> String {
340        let id = self.request_id_counter.fetch_add(1, Ordering::Relaxed);
341        format!("ws-{id}")
342    }
343
344    fn sign_params(&self, mut params: serde_json::Value) -> BinanceWsApiResult<serde_json::Value> {
345        let timestamp = std::time::SystemTime::now()
346            .duration_since(std::time::UNIX_EPOCH)
347            .map_err(|e| BinanceWsApiError::ClientError(e.to_string()))?
348            .as_millis() as i64;
349
350        if let Some(obj) = params.as_object_mut() {
351            obj.insert("timestamp".to_string(), serde_json::json!(timestamp));
352            obj.insert(
353                "apiKey".to_string(),
354                serde_json::json!(self.credential.api_key()),
355            );
356
357            if let Some(recv_window_ms) = self.recv_window_ms {
358                obj.insert("recvWindow".to_string(), serde_json::json!(recv_window_ms));
359            }
360        }
361
362        // Sign over a key-sorted query string: Binance's WS API verifies the
363        // signature against the parameters sorted by key, which must not depend
364        // on serde_json's Map iteration order (issue #4410).
365        let query_string = canonical_ws_query_string(
366            params
367                .as_object()
368                .into_iter()
369                .flatten()
370                .map(|(k, v)| (k.as_str(), v)),
371        )
372        .map_err(|e| BinanceWsApiError::ClientError(e.to_string()))?;
373        let signature = self.credential.sign(&query_string);
374
375        if let Some(obj) = params.as_object_mut() {
376            obj.insert("signature".to_string(), serde_json::json!(signature));
377        }
378
379        Ok(params)
380    }
381
382    async fn send_request(
383        &mut self,
384        request: BinanceSpotWsTradingRequest,
385    ) -> BinanceWsApiResult<()> {
386        let client = self.inner.as_mut().ok_or_else(|| {
387            BinanceWsApiError::ConnectionError("WebSocket not connected".to_string())
388        })?;
389
390        let json = serde_json::to_string(&request)
391            .map_err(|e| BinanceWsApiError::ClientError(e.to_string()))?;
392
393        log::debug!(
394            "Sending WebSocket API request id={} method={}",
395            request.id,
396            request.method
397        );
398
399        // Apply rate limiting for order operations
400        client
401            .send_text(json, Some(BINANCE_WS_RATE_LIMIT_KEY_ORDER.as_slice()))
402            .await
403            .map_err(|e| {
404                BinanceWsApiError::ConnectionError(format!("Failed to send request: {e}"))
405            })?;
406
407        Ok(())
408    }
409
410    fn handle_message(&mut self, msg: Message) {
411        match msg {
412            Message::Binary(data) => self.handle_binary_response(&data),
413            Message::Text(text) => self.handle_text_response(&text),
414            Message::Ping(_) | Message::Pong(_) => {}
415            Message::Close(frame) => {
416                log::debug!("WebSocket closed: {frame:?}");
417            }
418            Message::Frame(_) => {}
419        }
420    }
421
422    fn handle_binary_response(&mut self, data: &[u8]) {
423        match self.decode_ws_api_response(data) {
424            Ok(response) => self.emit(response),
425            Err(e) => {
426                log::error!("Failed to decode WebSocket API response: {e}");
427                self.emit(BinanceSpotWsTradingMessage::Error(e.to_string()));
428            }
429        }
430    }
431
432    fn handle_text_response(&mut self, text: &str) {
433        let json: serde_json::Value = match serde_json::from_str(text) {
434            Ok(j) => j,
435            Err(e) => {
436                log::warn!("Failed to parse text response as JSON: {e}");
437                return;
438            }
439        };
440
441        // User data events arrive wrapped: {"subscriptionId": N, "event": {...}}
442        if let Some(event) = json.get("event") {
443            self.handle_user_data_event(event);
444            return;
445        }
446
447        // Legacy listen-key streams, including Binance US, send the event directly.
448        if json.get("e").is_some() {
449            self.handle_user_data_event(&json);
450            return;
451        }
452
453        // WS API responses have an "id" field for request correlation
454        if let Some(id) = json.get("id") {
455            let id_str = match id {
456                serde_json::Value::String(s) => s.clone(),
457                serde_json::Value::Number(n) => n.to_string(),
458                _ => return,
459            };
460
461            if let Some(meta) = self.pending_requests.remove(&id_str) {
462                // Check for error: nested {"error": {"code": N, "msg": "..."}}
463                // or top-level {"code": N, "msg": "..."}
464                let error_info = json
465                    .get("error")
466                    .map(|e| {
467                        (
468                            e.get("code").and_then(|v| v.as_i64()).unwrap_or(-1),
469                            e.get("msg")
470                                .and_then(|v| v.as_str())
471                                .unwrap_or("Unknown error")
472                                .to_string(),
473                        )
474                    })
475                    .or_else(|| {
476                        json.get("code").and_then(|c| c.as_i64()).map(|code| {
477                            let msg = json
478                                .get("msg")
479                                .and_then(|v| v.as_str())
480                                .unwrap_or("Unknown error")
481                                .to_string();
482                            (code, msg)
483                        })
484                    });
485
486                if let Some((code, msg)) = error_info {
487                    let rejection = self.create_rejection(id_str, code as i32, msg, meta);
488                    self.emit(rejection);
489                    return;
490                }
491
492                // Success response
493                match meta {
494                    BinanceSpotWsTradingRequestMeta::SessionLogon => {
495                        log::debug!("Session authenticated");
496                        self.emit(BinanceSpotWsTradingMessage::Authenticated);
497                    }
498                    BinanceSpotWsTradingRequestMeta::SubscribeUserData => {
499                        let subscription_id = json
500                            .get("result")
501                            .and_then(|r| r.get("subscriptionId"))
502                            .map(|v| v.to_string())
503                            .unwrap_or_default();
504                        log::debug!("User data stream subscribed: id={subscription_id}");
505                        self.emit(BinanceSpotWsTradingMessage::UserDataSubscribed {
506                            subscription_id,
507                        });
508                    }
509                    _ => {
510                        // Order operation responses come as SBE binary, not JSON text.
511                        // If we get a JSON success for an order operation, log it.
512                        log::debug!("Unexpected JSON success for request {id_str}: {json}");
513                    }
514                }
515                return;
516            }
517
518            // Error response without matching pending request
519            if let Some(code) = json.get("code").and_then(|v| v.as_i64()) {
520                let msg = json
521                    .get("msg")
522                    .and_then(|v| v.as_str())
523                    .unwrap_or("Unknown error");
524                log::warn!(
525                    "Received error response without matching request ID: code={code} msg={msg}"
526                );
527            }
528            return;
529        }
530
531        // Stream termination event
532        if json.get("eventStreamTerminated").is_some() {
533            log::warn!("User data stream terminated, resubscribe needed");
534            return;
535        }
536
537        log::debug!("Unhandled text message: {text}");
538    }
539
540    fn handle_user_data_event(&self, event: &serde_json::Value) {
541        if let Some(msg) = classify_user_data_event(event) {
542            self.emit(msg);
543        }
544    }
545
546    fn decode_ws_api_response(
547        &mut self,
548        data: &[u8],
549    ) -> Result<BinanceSpotWsTradingMessage, BinanceWsApiError> {
550        // Check template ID before parsing
551        if data.len() >= message_header_codec::ENCODED_LENGTH {
552            let buf = ReadBuf::new(data);
553            let template_id = buf.get_u16_at(2);
554
555            // User data stream events arrive as SBE with their own template IDs
556            // (not wrapped in WebSocketResponse template 50).
557            match template_id {
558                601 => {
559                    log::debug!("Received SBE BalanceUpdateEvent ({} bytes)", data.len());
560                    match super::decode_sbe::decode_balance_update(data) {
561                        Ok(msg) => {
562                            log::debug!(
563                                "SBE balance update: asset={}, delta={}",
564                                msg.asset,
565                                msg.delta
566                            );
567                            return Ok(BinanceSpotWsTradingMessage::BalanceUpdate(msg));
568                        }
569                        Err(e) => {
570                            log::error!("Failed to decode SBE BalanceUpdateEvent: {e}");
571                            return Ok(BinanceSpotWsTradingMessage::Error(format!(
572                                "SBE BalanceUpdateEvent decode failed: {e}"
573                            )));
574                        }
575                    }
576                }
577                603 => {
578                    log::debug!("Received SBE ExecutionReportEvent ({} bytes)", data.len());
579                    match super::decode_sbe::decode_execution_report(data) {
580                        Ok(report) => {
581                            log::debug!(
582                                "SBE execution report: symbol={}, order_id={}, exec={:?}, status={:?}",
583                                report.symbol,
584                                report.order_id,
585                                report.execution_type,
586                                report.order_status
587                            );
588                            return Ok(BinanceSpotWsTradingMessage::ExecutionReport(Box::new(
589                                report,
590                            )));
591                        }
592                        Err(e) => {
593                            log::error!("Failed to decode SBE ExecutionReportEvent: {e}");
594                            return Ok(BinanceSpotWsTradingMessage::Error(format!(
595                                "SBE ExecutionReportEvent decode failed: {e}"
596                            )));
597                        }
598                    }
599                }
600                606 => {
601                    log::debug!(
602                        "Received SBE ListStatusEvent ({} bytes), not yet decoded",
603                        data.len()
604                    );
605                    return Ok(BinanceSpotWsTradingMessage::Error(
606                        "SBE ListStatusEvent decoding not yet implemented".to_string(),
607                    ));
608                }
609                607 => {
610                    log::debug!(
611                        "Received SBE OutboundAccountPositionEvent ({} bytes)",
612                        data.len()
613                    );
614
615                    match super::decode_sbe::decode_account_position(data) {
616                        Ok(msg) => {
617                            log::debug!("SBE account position: {} balance(s)", msg.balances.len());
618                            return Ok(BinanceSpotWsTradingMessage::AccountPosition(msg));
619                        }
620                        Err(e) => {
621                            log::error!("Failed to decode SBE OutboundAccountPositionEvent: {e}");
622                            return Ok(BinanceSpotWsTradingMessage::Error(format!(
623                                "SBE OutboundAccountPositionEvent decode failed: {e}"
624                            )));
625                        }
626                    }
627                }
628                610 => {
629                    let event_time = parse_server_shutdown_event_time_ms(data);
630                    log::warn!(
631                        "Binance server shutdown notice (SBE, event_time={event_time}); disconnect expected within ~10 minutes",
632                    );
633                    return Ok(BinanceSpotWsTradingMessage::ServerShutdown { event_time });
634                }
635                _ => {} // Fall through to WebSocketResponse parsing
636            }
637        }
638
639        // Standard WebSocketResponse envelope (template 50)
640        let (request_id, status, result_data) = self.parse_envelope(data)?;
641
642        // Look up the pending request by ID
643        let meta = self.pending_requests.remove(&request_id).ok_or_else(|| {
644            BinanceWsApiError::UnknownRequestId(format!("No pending request for ID: {request_id}"))
645        })?;
646
647        // Check for error status (non-200)
648        if status != 200 {
649            let (code, msg) = Self::try_decode_sbe_error(&result_data).unwrap_or((
650                status as i32,
651                format!("Request failed with status {status}"),
652            ));
653            return Ok(self.create_rejection(request_id, code, msg, meta));
654        }
655
656        // Decode the inner payload based on request type
657        match meta {
658            BinanceSpotWsTradingRequestMeta::PlaceOrder => {
659                let response = parse::decode_new_order_full(&result_data)?;
660                Ok(BinanceSpotWsTradingMessage::OrderAccepted {
661                    request_id,
662                    response,
663                })
664            }
665            BinanceSpotWsTradingRequestMeta::CancelOrder => {
666                let response = parse::decode_cancel_order(&result_data)?;
667                Ok(BinanceSpotWsTradingMessage::OrderCanceled {
668                    request_id,
669                    response,
670                })
671            }
672            BinanceSpotWsTradingRequestMeta::CancelReplaceOrder => {
673                // Cancel-replace returns both cancel and new order info
674                let new_order_response = parse::decode_new_order_full(&result_data)?;
675                let cancel_response = BinanceCancelOrderResponse {
676                    price_exponent: new_order_response.price_exponent,
677                    qty_exponent: new_order_response.qty_exponent,
678                    order_id: 0,
679                    order_list_id: None,
680                    transact_time: new_order_response.transact_time,
681                    price_mantissa: 0,
682                    orig_qty_mantissa: 0,
683                    executed_qty_mantissa: 0,
684                    cummulative_quote_qty_mantissa: 0,
685                    status: crate::spot::sbe::spot::order_status::OrderStatus::Canceled,
686                    time_in_force: new_order_response.time_in_force,
687                    order_type: new_order_response.order_type,
688                    side: new_order_response.side,
689                    self_trade_prevention_mode: new_order_response.self_trade_prevention_mode,
690                    client_order_id: String::new(),
691                    orig_client_order_id: String::new(),
692                    symbol: new_order_response.symbol.clone(),
693                };
694                Ok(BinanceSpotWsTradingMessage::CancelReplaceAccepted {
695                    request_id,
696                    cancel_response,
697                    new_order_response,
698                })
699            }
700            BinanceSpotWsTradingRequestMeta::CancelAllOrders => {
701                let responses = parse::decode_cancel_open_orders(&result_data)?;
702                Ok(BinanceSpotWsTradingMessage::AllOrdersCanceled {
703                    request_id,
704                    responses,
705                })
706            }
707            BinanceSpotWsTradingRequestMeta::SessionLogon => {
708                log::debug!("Session authenticated (SBE response)");
709                Ok(BinanceSpotWsTradingMessage::Authenticated)
710            }
711            BinanceSpotWsTradingRequestMeta::SubscribeUserData => {
712                log::debug!("User data stream subscribed (SBE response)");
713                Ok(BinanceSpotWsTradingMessage::UserDataSubscribed {
714                    subscription_id: request_id,
715                })
716            }
717        }
718    }
719
720    /// Parses the WebSocketResponse SBE envelope.
721    ///
722    /// Returns (request_id, status, result_payload).
723    fn parse_envelope(&self, data: &[u8]) -> Result<(String, u16, Vec<u8>), BinanceWsApiError> {
724        if data.len() < message_header_codec::ENCODED_LENGTH {
725            return Err(BinanceWsApiError::DecodeError(
726                crate::spot::sbe::error::SbeDecodeError::BufferTooShort {
727                    expected: message_header_codec::ENCODED_LENGTH,
728                    actual: data.len(),
729                },
730            ));
731        }
732
733        let buf = ReadBuf::new(data);
734
735        // Parse message header
736        let block_length = buf.get_u16_at(0);
737        let template_id = buf.get_u16_at(2);
738
739        if template_id != SBE_TEMPLATE_ID {
740            return Err(BinanceWsApiError::DecodeError(
741                crate::spot::sbe::error::SbeDecodeError::UnknownTemplateId(template_id),
742            ));
743        }
744
745        let version = buf.get_u16_at(6);
746
747        // Create decoder at offset after message header
748        let decoder = WebSocketResponseDecoder::default().wrap(
749            buf,
750            message_header_codec::ENCODED_LENGTH,
751            block_length,
752            version,
753        );
754
755        // Read status from fixed block (offset 1 within block)
756        let status = decoder.status();
757
758        // Skip rate_limits group
759        let mut rate_limits = decoder.rate_limits_decoder();
760        while rate_limits.advance().unwrap_or(None).is_some() {}
761        let mut decoder = rate_limits.parent().map_err(|e| {
762            BinanceWsApiError::ClientError(format!("Failed to get parent from rate_limits: {e}"))
763        })?;
764
765        // Extract request ID
766        let id_coords = decoder.id_decoder();
767        let id_bytes = decoder.id_slice(id_coords);
768        let request_id = String::from_utf8_lossy(id_bytes).to_string();
769
770        // Extract result payload - copy to owned Vec to avoid lifetime issues
771        let result_coords = decoder.result_decoder();
772        let result_data = decoder.result_slice(result_coords).to_vec();
773
774        Ok((request_id, status, result_data))
775    }
776
777    fn create_rejection(
778        &self,
779        request_id: String,
780        code: i32,
781        msg: String,
782        meta: BinanceSpotWsTradingRequestMeta,
783    ) -> BinanceSpotWsTradingMessage {
784        match meta {
785            BinanceSpotWsTradingRequestMeta::PlaceOrder => {
786                BinanceSpotWsTradingMessage::OrderRejected {
787                    request_id,
788                    code,
789                    msg,
790                }
791            }
792            BinanceSpotWsTradingRequestMeta::CancelOrder => {
793                BinanceSpotWsTradingMessage::CancelRejected {
794                    request_id,
795                    code,
796                    msg,
797                }
798            }
799            BinanceSpotWsTradingRequestMeta::CancelReplaceOrder => {
800                BinanceSpotWsTradingMessage::CancelReplaceRejected {
801                    request_id,
802                    code,
803                    msg,
804                }
805            }
806            BinanceSpotWsTradingRequestMeta::CancelAllOrders => {
807                BinanceSpotWsTradingMessage::CancelRejected {
808                    request_id,
809                    code,
810                    msg,
811                }
812            }
813            BinanceSpotWsTradingRequestMeta::SessionLogon => {
814                BinanceSpotWsTradingMessage::AuthenticationRejected(format!("code={code}: {msg}"))
815            }
816            BinanceSpotWsTradingRequestMeta::SubscribeUserData => {
817                BinanceSpotWsTradingMessage::UserDataSubscriptionRejected(format!(
818                    "code={code}: {msg}"
819                ))
820            }
821        }
822    }
823
824    // Decodes the SBE error response to extract the Binance error code and message
825    fn try_decode_sbe_error(data: &[u8]) -> Option<(i32, String)> {
826        const HEADER_LEN: usize = 8;
827
828        if data.len()
829            < HEADER_LEN + crate::spot::sbe::spot::error_response_codec::SBE_BLOCK_LENGTH as usize
830        {
831            return None;
832        }
833
834        let buf = ReadBuf::new(data);
835        let header = message_header_codec::MessageHeaderDecoder::default().wrap(buf, 0);
836        if header.template_id() != crate::spot::sbe::spot::error_response_codec::SBE_TEMPLATE_ID {
837            return None;
838        }
839
840        let mut decoder = ErrorResponseDecoder::default().header(header, 0);
841        let code = i32::from(decoder.code());
842        let msg_coords = decoder.msg_decoder();
843        let msg_bytes = decoder.msg_slice(msg_coords);
844        let msg = String::from_utf8_lossy(msg_bytes).into_owned();
845
846        Some((code, msg))
847    }
848}
849
850/// Classifies a JSON user-data event into a trading message, if any.
851///
852/// Returns `None` when the event type is unknown or the payload fails to
853/// deserialize; in that case the caller logs and drops the event.
854pub(crate) fn classify_user_data_event(
855    event: &serde_json::Value,
856) -> Option<BinanceSpotWsTradingMessage> {
857    let event_type = event
858        .get("e")
859        .and_then(|v| serde_json::from_value::<BinanceSpotUserDataEventType>(v.clone()).ok())
860        .unwrap_or(BinanceSpotUserDataEventType::Unknown);
861
862    match event_type {
863        BinanceSpotUserDataEventType::ExecutionReport => {
864            match serde_json::from_value::<super::user_data::BinanceSpotExecutionReport>(
865                event.clone(),
866            ) {
867                Ok(report) => {
868                    log::debug!(
869                        "Execution report: symbol={}, order_id={}, exec={:?}, status={:?}",
870                        report.symbol,
871                        report.order_id,
872                        report.execution_type,
873                        report.order_status
874                    );
875                    Some(BinanceSpotWsTradingMessage::ExecutionReport(Box::new(
876                        report,
877                    )))
878                }
879                Err(e) => {
880                    log::warn!("Failed to parse execution report: {e}");
881                    None
882                }
883            }
884        }
885        BinanceSpotUserDataEventType::OutboundAccountPosition => {
886            match serde_json::from_value::<super::user_data::BinanceSpotAccountPositionMsg>(
887                event.clone(),
888            ) {
889                Ok(msg) => {
890                    log::debug!("Account position update: {} balance(s)", msg.balances.len());
891                    Some(BinanceSpotWsTradingMessage::AccountPosition(msg))
892                }
893                Err(e) => {
894                    log::warn!("Failed to parse account position: {e}");
895                    None
896                }
897            }
898        }
899        BinanceSpotUserDataEventType::BalanceUpdate => {
900            match serde_json::from_value::<super::user_data::BinanceSpotBalanceUpdateMsg>(
901                event.clone(),
902            ) {
903                Ok(msg) => {
904                    log::debug!("Balance update: asset={}, delta={}", msg.asset, msg.delta);
905                    Some(BinanceSpotWsTradingMessage::BalanceUpdate(msg))
906                }
907                Err(e) => {
908                    log::warn!("Failed to parse balance update: {e}");
909                    None
910                }
911            }
912        }
913        BinanceSpotUserDataEventType::ServerShutdown => {
914            let event_time = event.get("E").and_then(|v| v.as_i64()).unwrap_or_default();
915            log::warn!(
916                "Binance server shutdown notice (event_time={event_time}); disconnect expected within ~10 minutes",
917            );
918            Some(BinanceSpotWsTradingMessage::ServerShutdown { event_time })
919        }
920        BinanceSpotUserDataEventType::ListenKeyExpired
921        | BinanceSpotUserDataEventType::ExternalLockUpdate
922        | BinanceSpotUserDataEventType::EventStreamTerminated
923        | BinanceSpotUserDataEventType::Unknown => {
924            log::debug!("Unhandled user data event type: {event_type:?}");
925            None
926        }
927    }
928}
929
930/// Parses the `event_time` from an SBE `ServerShutdownEvent` (template 610) frame.
931///
932/// The SBE field is microseconds; the trading message variant documents
933/// milliseconds (matching the JSON dispatch), so this divides by 1_000.
934/// Returns `0` when the buffer is too short to contain the field.
935pub(crate) fn parse_server_shutdown_event_time_ms(data: &[u8]) -> i64 {
936    if data.len() < message_header_codec::ENCODED_LENGTH + 8 {
937        return 0;
938    }
939    let buf = ReadBuf::new(data);
940    buf.get_i64_at(message_header_codec::ENCODED_LENGTH) / 1_000
941}
942
943#[cfg(test)]
944mod tests {
945    use rstest::rstest;
946
947    use super::*;
948
949    #[rstest]
950    #[case::microseconds_converted_to_ms(1_700_000_000_000_000_i64, 1_700_000_000_000_i64)]
951    #[case::zero(0_i64, 0_i64)]
952    #[case::negative(-1_000_i64, -1_i64)]
953    fn test_parse_server_shutdown_event_time_ms(
954        #[case] event_time_us: i64,
955        #[case] expected_ms: i64,
956    ) {
957        let mut buf = vec![0u8; message_header_codec::ENCODED_LENGTH];
958        buf.extend_from_slice(&event_time_us.to_le_bytes());
959        assert_eq!(parse_server_shutdown_event_time_ms(&buf), expected_ms);
960    }
961
962    #[rstest]
963    fn test_parse_server_shutdown_event_time_ms_short_buffer_returns_zero() {
964        let buf = vec![0u8; message_header_codec::ENCODED_LENGTH + 4];
965        assert_eq!(parse_server_shutdown_event_time_ms(&buf), 0);
966    }
967
968    #[rstest]
969    fn test_classify_user_data_event_server_shutdown_emits_variant() {
970        let event = serde_json::json!({"e": "serverShutdown", "E": 1_700_000_000_000_i64});
971        let msg = classify_user_data_event(&event).expect("expected ServerShutdown");
972        match msg {
973            BinanceSpotWsTradingMessage::ServerShutdown { event_time } => {
974                assert_eq!(event_time, 1_700_000_000_000);
975            }
976            other => panic!("expected ServerShutdown variant, was {other:?}"),
977        }
978    }
979
980    #[rstest]
981    fn test_classify_user_data_event_server_shutdown_missing_event_time_defaults_to_zero() {
982        let event = serde_json::json!({"e": "serverShutdown"});
983        let msg = classify_user_data_event(&event).expect("expected ServerShutdown");
984        match msg {
985            BinanceSpotWsTradingMessage::ServerShutdown { event_time } => {
986                assert_eq!(event_time, 0);
987            }
988            other => panic!("expected ServerShutdown variant, was {other:?}"),
989        }
990    }
991
992    #[rstest]
993    fn test_classify_user_data_event_unknown_returns_none() {
994        let event = serde_json::json!({"e": "somethingElse"});
995        assert!(classify_user_data_event(&event).is_none());
996    }
997
998    #[rstest]
999    fn test_sign_params_includes_recv_window_in_signature() {
1000        let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
1001        let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
1002        let (out_tx, _out_rx) = tokio::sync::mpsc::unbounded_channel();
1003        let credential = Arc::new(SigningCredential::new(
1004            "api-key".to_string(),
1005            "hmac-secret".to_string(),
1006        ));
1007        let handler = BinanceSpotWsTradingHandler::new(
1008            Arc::new(AtomicBool::new(false)),
1009            cmd_rx,
1010            raw_rx,
1011            out_tx,
1012            credential.clone(),
1013        )
1014        .with_recv_window(Some(45_000));
1015
1016        let signed = handler
1017            .sign_params(serde_json::json!({"symbol": "BTCUSDT"}))
1018            .unwrap();
1019        let mut unsigned = signed.clone();
1020        let signature = unsigned
1021            .as_object_mut()
1022            .unwrap()
1023            .remove("signature")
1024            .unwrap();
1025        let query = canonical_ws_query_string(
1026            unsigned
1027                .as_object()
1028                .unwrap()
1029                .iter()
1030                .map(|(key, value)| (key.as_str(), value)),
1031        )
1032        .unwrap();
1033
1034        assert_eq!(signed["recvWindow"], 45_000);
1035        assert_eq!(signature, credential.sign(&query));
1036    }
1037}