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::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")
469                                .and_then(|v| v.as_i64())
470                                .and_then(|code| i32::try_from(code).ok()),
471                            e.get("msg")
472                                .and_then(|v| v.as_str())
473                                .unwrap_or("Unknown error")
474                                .to_string(),
475                        )
476                    })
477                    .or_else(|| {
478                        json.get("code").and_then(|c| c.as_i64()).map(|code| {
479                            let msg = json
480                                .get("msg")
481                                .and_then(|v| v.as_str())
482                                .unwrap_or("Unknown error")
483                                .to_string();
484                            (i32::try_from(code).ok(), msg)
485                        })
486                    });
487
488                if let Some((code, msg)) = error_info {
489                    let status = json
490                        .get("status")
491                        .and_then(|v| v.as_u64())
492                        .map(|s| s as u16);
493                    let rejection = match code {
494                        Some(code) => {
495                            self.create_rejection(id_str, status.unwrap_or(0), code, msg, meta)
496                        }
497                        None => BinanceSpotWsTradingMessage::RequestFailed {
498                            request_id: id_str,
499                            msg: format!("Missing or invalid venue error code: {msg}"),
500                        },
501                    };
502                    self.emit(rejection);
503                    return;
504                }
505
506                // Success response
507                match meta {
508                    BinanceSpotWsTradingRequestMeta::SessionLogon => {
509                        log::debug!("Session authenticated");
510                        self.emit(BinanceSpotWsTradingMessage::Authenticated);
511                    }
512                    BinanceSpotWsTradingRequestMeta::SubscribeUserData => {
513                        let subscription_id = json
514                            .get("result")
515                            .and_then(|r| r.get("subscriptionId"))
516                            .map(|v| v.to_string())
517                            .unwrap_or_default();
518                        log::debug!("User data stream subscribed: id={subscription_id}");
519                        self.emit(BinanceSpotWsTradingMessage::UserDataSubscribed {
520                            subscription_id,
521                        });
522                    }
523                    _ => {
524                        // Order operation responses come as SBE binary, not JSON text.
525                        // If we get a JSON success for an order operation, log it.
526                        log::debug!("Unexpected JSON success for request {id_str}");
527                    }
528                }
529                return;
530            }
531
532            // Error response without matching pending request
533            if let Some(code) = json.get("code").and_then(|v| v.as_i64()) {
534                let msg = json
535                    .get("msg")
536                    .and_then(|v| v.as_str())
537                    .unwrap_or("Unknown error");
538                log::warn!(
539                    "Received error response without matching request ID: code={code} msg={msg}"
540                );
541            }
542            return;
543        }
544
545        // Stream termination event
546        if json.get("eventStreamTerminated").is_some() {
547            log::warn!("User data stream terminated, resubscribe needed");
548            return;
549        }
550
551        log::debug!("Unhandled text message: {} bytes", text.len());
552    }
553
554    fn handle_user_data_event(&self, event: &serde_json::Value) {
555        if let Some(msg) = classify_user_data_event(event) {
556            self.emit(msg);
557        }
558    }
559
560    fn decode_ws_api_response(
561        &mut self,
562        data: &[u8],
563    ) -> Result<BinanceSpotWsTradingMessage, BinanceWsApiError> {
564        // Check template ID before parsing
565        if data.len() >= message_header_codec::ENCODED_LENGTH {
566            let buf = ReadBuf::new(data);
567            let template_id = buf.get_u16_at(2);
568
569            // User data stream events arrive as SBE with their own template IDs
570            // (not wrapped in WebSocketResponse template 50).
571            match template_id {
572                601 => {
573                    log::debug!("Received SBE BalanceUpdateEvent ({} bytes)", data.len());
574                    match super::decode_sbe::decode_balance_update(data) {
575                        Ok(msg) => {
576                            log::debug!(
577                                "SBE balance update: asset={}, delta={}",
578                                msg.asset,
579                                msg.delta
580                            );
581                            return Ok(BinanceSpotWsTradingMessage::BalanceUpdate(msg));
582                        }
583                        Err(e) => {
584                            log::error!("Failed to decode SBE BalanceUpdateEvent: {e}");
585                            return Ok(BinanceSpotWsTradingMessage::Error(format!(
586                                "SBE BalanceUpdateEvent decode failed: {e}"
587                            )));
588                        }
589                    }
590                }
591                603 => {
592                    log::debug!("Received SBE ExecutionReportEvent ({} bytes)", data.len());
593                    match super::decode_sbe::decode_execution_report(data) {
594                        Ok(report) => {
595                            log::debug!(
596                                "SBE execution report: symbol={}, order_id={}, exec={:?}, status={:?}",
597                                report.symbol,
598                                report.order_id,
599                                report.execution_type,
600                                report.order_status
601                            );
602                            return Ok(BinanceSpotWsTradingMessage::ExecutionReport(Box::new(
603                                report,
604                            )));
605                        }
606                        Err(e) => {
607                            log::error!("Failed to decode SBE ExecutionReportEvent: {e}");
608                            return Ok(BinanceSpotWsTradingMessage::Error(format!(
609                                "SBE ExecutionReportEvent decode failed: {e}"
610                            )));
611                        }
612                    }
613                }
614                606 => {
615                    log::debug!(
616                        "Received SBE ListStatusEvent ({} bytes), not yet decoded",
617                        data.len()
618                    );
619                    return Ok(BinanceSpotWsTradingMessage::Error(
620                        "SBE ListStatusEvent decoding not yet implemented".to_string(),
621                    ));
622                }
623                607 => {
624                    log::debug!(
625                        "Received SBE OutboundAccountPositionEvent ({} bytes)",
626                        data.len()
627                    );
628
629                    match super::decode_sbe::decode_account_position(data) {
630                        Ok(msg) => {
631                            log::debug!("SBE account position: {} balance(s)", msg.balances.len());
632                            return Ok(BinanceSpotWsTradingMessage::AccountPosition(msg));
633                        }
634                        Err(e) => {
635                            log::error!("Failed to decode SBE OutboundAccountPositionEvent: {e}");
636                            return Ok(BinanceSpotWsTradingMessage::Error(format!(
637                                "SBE OutboundAccountPositionEvent decode failed: {e}"
638                            )));
639                        }
640                    }
641                }
642                610 => {
643                    let event_time = parse_server_shutdown_event_time_ms(data);
644                    log::warn!(
645                        "Binance server shutdown notice (SBE, event_time={event_time}); disconnect expected within ~10 minutes",
646                    );
647                    return Ok(BinanceSpotWsTradingMessage::ServerShutdown { event_time });
648                }
649                _ => {} // Fall through to WebSocketResponse parsing
650            }
651        }
652
653        // Standard WebSocketResponse envelope (template 50)
654        let (request_id, status, result_data) = self.parse_envelope(data)?;
655
656        // Look up the pending request by ID
657        let meta = self.pending_requests.remove(&request_id).ok_or_else(|| {
658            BinanceWsApiError::UnknownRequestId(format!("No pending request for ID: {request_id}"))
659        })?;
660
661        // Check for error status (non-200)
662        if status != 200 {
663            return Ok(match Self::try_decode_sbe_error(&result_data) {
664                Some((code, msg)) => self.create_rejection(request_id, status, code, msg, meta),
665                // An undecodable error payload carries no definitive command evidence
666                None => BinanceSpotWsTradingMessage::RequestFailed {
667                    request_id,
668                    msg: format!("Request failed with status {status}; error payload undecodable"),
669                },
670            });
671        }
672
673        // Decode the inner payload based on request type
674        match meta {
675            BinanceSpotWsTradingRequestMeta::PlaceOrder => {
676                let response = parse::decode_new_order_full(&result_data)?;
677                Ok(BinanceSpotWsTradingMessage::OrderAccepted {
678                    request_id,
679                    response,
680                })
681            }
682            BinanceSpotWsTradingRequestMeta::CancelOrder => {
683                let response = parse::decode_cancel_order(&result_data)?;
684                Ok(BinanceSpotWsTradingMessage::OrderCanceled {
685                    request_id,
686                    response,
687                })
688            }
689            BinanceSpotWsTradingRequestMeta::CancelReplaceOrder => {
690                let (cancel_response, new_order_response) =
691                    parse::decode_cancel_replace_orders(&result_data)?;
692                Ok(BinanceSpotWsTradingMessage::CancelReplaceAccepted {
693                    request_id,
694                    cancel_response,
695                    new_order_response,
696                })
697            }
698            BinanceSpotWsTradingRequestMeta::CancelAllOrders => {
699                let responses = parse::decode_cancel_open_orders(&result_data)?;
700                Ok(BinanceSpotWsTradingMessage::AllOrdersCanceled {
701                    request_id,
702                    responses,
703                })
704            }
705            BinanceSpotWsTradingRequestMeta::SessionLogon => {
706                log::debug!("Session authenticated (SBE response)");
707                Ok(BinanceSpotWsTradingMessage::Authenticated)
708            }
709            BinanceSpotWsTradingRequestMeta::SubscribeUserData => {
710                log::debug!("User data stream subscribed (SBE response)");
711                Ok(BinanceSpotWsTradingMessage::UserDataSubscribed {
712                    subscription_id: request_id,
713                })
714            }
715        }
716    }
717
718    /// Parses the WebSocketResponse SBE envelope.
719    ///
720    /// Returns (request_id, status, result_payload).
721    fn parse_envelope(&self, data: &[u8]) -> Result<(String, u16, Vec<u8>), BinanceWsApiError> {
722        if data.len() < message_header_codec::ENCODED_LENGTH {
723            return Err(BinanceWsApiError::DecodeError(
724                crate::spot::sbe::error::SbeDecodeError::BufferTooShort {
725                    expected: message_header_codec::ENCODED_LENGTH,
726                    actual: data.len(),
727                },
728            ));
729        }
730
731        let buf = ReadBuf::new(data);
732
733        // Parse message header
734        let block_length = buf.get_u16_at(0);
735        let template_id = buf.get_u16_at(2);
736
737        if template_id != SBE_TEMPLATE_ID {
738            return Err(BinanceWsApiError::DecodeError(
739                crate::spot::sbe::error::SbeDecodeError::UnknownTemplateId(template_id),
740            ));
741        }
742
743        let version = buf.get_u16_at(6);
744
745        // Create decoder at offset after message header
746        let decoder = WebSocketResponseDecoder::default().wrap(
747            buf,
748            message_header_codec::ENCODED_LENGTH,
749            block_length,
750            version,
751        );
752
753        // Read status from fixed block (offset 1 within block)
754        let status = decoder.status();
755
756        // Skip rate_limits group
757        let mut rate_limits = decoder.rate_limits_decoder();
758        while rate_limits.advance().unwrap_or(None).is_some() {}
759        let mut decoder = rate_limits.parent().map_err(|e| {
760            BinanceWsApiError::ClientError(format!("Failed to get parent from rate_limits: {e}"))
761        })?;
762
763        // Extract request ID
764        let id_coords = decoder.id_decoder();
765        let id_bytes = decoder.id_slice(id_coords);
766        let request_id = String::from_utf8_lossy(id_bytes).to_string();
767
768        // Extract result payload - copy to owned Vec to avoid lifetime issues
769        let result_coords = decoder.result_decoder();
770        let result_data = decoder.result_slice(result_coords).to_vec();
771
772        Ok((request_id, status, result_data))
773    }
774
775    fn create_rejection(
776        &self,
777        request_id: String,
778        status: u16,
779        code: i32,
780        msg: String,
781        meta: BinanceSpotWsTradingRequestMeta,
782    ) -> BinanceSpotWsTradingMessage {
783        match meta {
784            BinanceSpotWsTradingRequestMeta::PlaceOrder => {
785                BinanceSpotWsTradingMessage::OrderRejected {
786                    request_id,
787                    status,
788                    code,
789                    msg,
790                }
791            }
792            BinanceSpotWsTradingRequestMeta::CancelOrder => {
793                BinanceSpotWsTradingMessage::CancelRejected {
794                    request_id,
795                    status,
796                    code,
797                    msg,
798                }
799            }
800            BinanceSpotWsTradingRequestMeta::CancelReplaceOrder => {
801                BinanceSpotWsTradingMessage::CancelReplaceRejected {
802                    request_id,
803                    status,
804                    code,
805                    msg,
806                }
807            }
808            BinanceSpotWsTradingRequestMeta::CancelAllOrders => {
809                BinanceSpotWsTradingMessage::CancelRejected {
810                    request_id,
811                    status,
812                    code,
813                    msg,
814                }
815            }
816            BinanceSpotWsTradingRequestMeta::SessionLogon => {
817                BinanceSpotWsTradingMessage::AuthenticationRejected(format!("code={code}: {msg}"))
818            }
819            BinanceSpotWsTradingRequestMeta::SubscribeUserData => {
820                BinanceSpotWsTradingMessage::UserDataSubscriptionRejected(format!(
821                    "code={code}: {msg}"
822                ))
823            }
824        }
825    }
826
827    // Decodes the SBE error response to extract the Binance error code and message
828    fn try_decode_sbe_error(data: &[u8]) -> Option<(i32, String)> {
829        const HEADER_LEN: usize = 8;
830
831        if data.len()
832            < HEADER_LEN + crate::spot::sbe::spot::error_response_codec::SBE_BLOCK_LENGTH as usize
833        {
834            return None;
835        }
836
837        let buf = ReadBuf::new(data);
838        let header = message_header_codec::MessageHeaderDecoder::default().wrap(buf, 0);
839        if header.template_id() != crate::spot::sbe::spot::error_response_codec::SBE_TEMPLATE_ID {
840            return None;
841        }
842
843        let mut decoder = ErrorResponseDecoder::default().header(header, 0);
844        let code = i32::from(decoder.code());
845        let msg_coords = decoder.msg_decoder();
846        let msg_bytes = decoder.msg_slice(msg_coords);
847        let msg = String::from_utf8_lossy(msg_bytes).into_owned();
848
849        Some((code, msg))
850    }
851}
852
853/// Classifies a JSON user-data event into a trading message, if any.
854///
855/// Returns `None` when the event type is unknown or the payload fails to
856/// deserialize; in that case the caller logs and drops the event.
857pub(crate) fn classify_user_data_event(
858    event: &serde_json::Value,
859) -> Option<BinanceSpotWsTradingMessage> {
860    let event_type = event
861        .get("e")
862        .and_then(|v| serde_json::from_value::<BinanceSpotUserDataEventType>(v.clone()).ok())
863        .unwrap_or(BinanceSpotUserDataEventType::Unknown);
864
865    match event_type {
866        BinanceSpotUserDataEventType::ExecutionReport => {
867            match serde_json::from_value::<super::user_data::BinanceSpotExecutionReport>(
868                event.clone(),
869            ) {
870                Ok(report) => {
871                    log::debug!(
872                        "Execution report: symbol={}, order_id={}, exec={:?}, status={:?}",
873                        report.symbol,
874                        report.order_id,
875                        report.execution_type,
876                        report.order_status
877                    );
878                    Some(BinanceSpotWsTradingMessage::ExecutionReport(Box::new(
879                        report,
880                    )))
881                }
882                Err(e) => {
883                    log::warn!("Failed to parse execution report: {e}");
884                    None
885                }
886            }
887        }
888        BinanceSpotUserDataEventType::OutboundAccountPosition => {
889            match serde_json::from_value::<super::user_data::BinanceSpotAccountPositionMsg>(
890                event.clone(),
891            ) {
892                Ok(msg) => {
893                    log::debug!("Account position update: {} balance(s)", msg.balances.len());
894                    Some(BinanceSpotWsTradingMessage::AccountPosition(msg))
895                }
896                Err(e) => {
897                    log::warn!("Failed to parse account position: {e}");
898                    None
899                }
900            }
901        }
902        BinanceSpotUserDataEventType::BalanceUpdate => {
903            match serde_json::from_value::<super::user_data::BinanceSpotBalanceUpdateMsg>(
904                event.clone(),
905            ) {
906                Ok(msg) => {
907                    log::debug!("Balance update: asset={}, delta={}", msg.asset, msg.delta);
908                    Some(BinanceSpotWsTradingMessage::BalanceUpdate(msg))
909                }
910                Err(e) => {
911                    log::warn!("Failed to parse balance update: {e}");
912                    None
913                }
914            }
915        }
916        BinanceSpotUserDataEventType::ServerShutdown => {
917            let event_time = event.get("E").and_then(|v| v.as_i64()).unwrap_or_default();
918            log::warn!(
919                "Binance server shutdown notice (event_time={event_time}); disconnect expected within ~10 minutes",
920            );
921            Some(BinanceSpotWsTradingMessage::ServerShutdown { event_time })
922        }
923        BinanceSpotUserDataEventType::ListenKeyExpired
924        | BinanceSpotUserDataEventType::ExternalLockUpdate
925        | BinanceSpotUserDataEventType::EventStreamTerminated
926        | BinanceSpotUserDataEventType::Unknown => {
927            log::debug!("Unhandled user data event type: {event_type:?}");
928            None
929        }
930    }
931}
932
933/// Parses the `event_time` from an SBE `ServerShutdownEvent` (template 610) frame.
934///
935/// The SBE field is microseconds; the trading message variant documents
936/// milliseconds (matching the JSON dispatch), so this divides by 1_000.
937/// Returns `0` when the buffer is too short to contain the field.
938pub(crate) fn parse_server_shutdown_event_time_ms(data: &[u8]) -> i64 {
939    if data.len() < message_header_codec::ENCODED_LENGTH + 8 {
940        return 0;
941    }
942    let buf = ReadBuf::new(data);
943    buf.get_i64_at(message_header_codec::ENCODED_LENGTH) / 1_000
944}
945
946#[cfg(test)]
947mod tests {
948    use rstest::rstest;
949
950    use super::*;
951    use crate::spot::sbe::spot::{
952        cancel_order_response_codec::CancelOrderResponseDecoder,
953        new_order_full_response_codec::NewOrderFullResponseDecoder,
954        self_trade_prevention_mode::SelfTradePreventionMode,
955    };
956
957    #[rstest]
958    fn test_cancel_replace_response_decodes_both_orders() {
959        let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
960        let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
961        let (out_tx, _out_rx) = tokio::sync::mpsc::unbounded_channel();
962        let mut handler = BinanceSpotWsTradingHandler::new(
963            Arc::new(AtomicBool::new(false)),
964            cmd_rx,
965            raw_rx,
966            out_tx,
967            Arc::new(SigningCredential::new(
968                "api-key".to_string(),
969                "secret".to_string(),
970            )),
971        );
972        let placed = include_bytes!(
973            "../../../../test_data/spot/user_data_sbe/mainnet/web_socket_response_1.sbe"
974        );
975        let canceled = include_bytes!(
976            "../../../../test_data/spot/user_data_sbe/mainnet/web_socket_response_2.sbe"
977        );
978        let (request_id, _, replacement) = handler.parse_envelope(placed).unwrap();
979        let (_, _, cancellation) = handler.parse_envelope(canceled).unwrap();
980        let expected_new = parse::decode_new_order_full(&replacement).unwrap();
981        let expected_cancel = parse::decode_cancel_order(&cancellation).unwrap();
982        let mut payload = Vec::new();
983
984        for field in [
985            2_u16,
986            crate::spot::sbe::spot::cancel_replace_order_response_codec::SBE_TEMPLATE_ID,
987            crate::spot::sbe::spot::SBE_SCHEMA_ID,
988            crate::spot::sbe::spot::SBE_SCHEMA_VERSION,
989        ] {
990            payload.extend_from_slice(&field.to_le_bytes());
991        }
992        let success =
993            crate::spot::sbe::spot::cancel_replace_status::CancelReplaceStatus::Success as u8;
994        payload.extend_from_slice(&[success, success]);
995        payload.extend_from_slice(&(cancellation.len() as u16).to_le_bytes());
996        payload.extend_from_slice(&cancellation);
997        payload.extend_from_slice(&(replacement.len() as u32).to_le_bytes());
998        payload.extend_from_slice(&replacement);
999        let mut response = placed[..placed.len() - replacement.len() - 4].to_vec();
1000        response.extend_from_slice(&(payload.len() as u32).to_le_bytes());
1001        response.extend_from_slice(&payload);
1002        handler.pending_requests.insert(
1003            request_id.clone(),
1004            BinanceSpotWsTradingRequestMeta::CancelReplaceOrder,
1005        );
1006
1007        let decoded = handler.decode_ws_api_response(&response).unwrap();
1008
1009        let BinanceSpotWsTradingMessage::CancelReplaceAccepted {
1010            request_id: actual_id,
1011            cancel_response,
1012            new_order_response,
1013        } = decoded
1014        else {
1015            panic!("Expected cancel-replace acceptance");
1016        };
1017        assert_eq!(actual_id, request_id);
1018        assert_eq!(cancel_response, expected_cancel);
1019        assert_eq!(new_order_response, expected_new);
1020        assert!(handler.pending_requests.is_empty());
1021    }
1022
1023    #[rstest]
1024    fn test_mainnet_order_responses_match_generated_decoders() {
1025        let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
1026        let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
1027        let (out_tx, _out_rx) = tokio::sync::mpsc::unbounded_channel();
1028        let handler = BinanceSpotWsTradingHandler::new(
1029            Arc::new(AtomicBool::new(false)),
1030            cmd_rx,
1031            raw_rx,
1032            out_tx,
1033            Arc::new(SigningCredential::new(
1034                "api-key".to_string(),
1035                "secret".to_string(),
1036            )),
1037        );
1038        let placed = include_bytes!(
1039            "../../../../test_data/spot/user_data_sbe/mainnet/web_socket_response_1.sbe"
1040        );
1041        let canceled = include_bytes!(
1042            "../../../../test_data/spot/user_data_sbe/mainnet/web_socket_response_2.sbe"
1043        );
1044        let (_, _, replacement) = handler.parse_envelope(placed).unwrap();
1045        let (_, _, cancellation) = handler.parse_envelope(canceled).unwrap();
1046        let placed_header = message_header_codec::MessageHeaderDecoder::default()
1047            .wrap(ReadBuf::new(&replacement), 0);
1048        let placed_decoder = NewOrderFullResponseDecoder::default().header(placed_header, 0);
1049        let canceled_header = message_header_codec::MessageHeaderDecoder::default()
1050            .wrap(ReadBuf::new(&cancellation), 0);
1051        let canceled_decoder = CancelOrderResponseDecoder::default().header(canceled_header, 0);
1052
1053        let new_order = parse::decode_new_order_full(&replacement).unwrap();
1054        let cancel = parse::decode_cancel_order(&cancellation).unwrap();
1055
1056        assert_eq!(
1057            new_order.self_trade_prevention_mode,
1058            placed_decoder.self_trade_prevention_mode()
1059        );
1060        assert_eq!(new_order.working_time, placed_decoder.working_time());
1061        assert_eq!(new_order.stop_price_mantissa, placed_decoder.stop_price());
1062        assert_eq!(
1063            cancel.self_trade_prevention_mode,
1064            canceled_decoder.self_trade_prevention_mode()
1065        );
1066        assert_eq!(
1067            cancel.self_trade_prevention_mode,
1068            SelfTradePreventionMode::ExpireMaker
1069        );
1070    }
1071
1072    #[rstest]
1073    #[case::missing(None)]
1074    #[case::overflow(Some(i64::MAX))]
1075    fn test_json_error_without_valid_code_is_ambiguous(#[case] code: Option<i64>) {
1076        let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
1077        let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
1078        let (out_tx, mut out_rx) = tokio::sync::mpsc::unbounded_channel();
1079        let mut handler = BinanceSpotWsTradingHandler::new(
1080            Arc::new(AtomicBool::new(false)),
1081            cmd_rx,
1082            raw_rx,
1083            out_tx,
1084            Arc::new(SigningCredential::new(
1085                "test-key".to_string(),
1086                "test-secret".to_string(),
1087            )),
1088        );
1089        handler.pending_requests.insert(
1090            "order-1".to_string(),
1091            BinanceSpotWsTradingRequestMeta::PlaceOrder,
1092        );
1093        let response = serde_json::json!({"id": "order-1", "status": 400, "error": {"code": code, "msg": "invalid request"}});
1094
1095        handler.handle_text_response(&response.to_string());
1096
1097        let BinanceSpotWsTradingMessage::RequestFailed { request_id, msg } =
1098            out_rx.try_recv().unwrap()
1099        else {
1100            panic!("expected ambiguous request failure");
1101        };
1102        assert_eq!(request_id, "order-1");
1103        assert_eq!(msg, "Missing or invalid venue error code: invalid request");
1104        assert!(handler.pending_requests.is_empty());
1105        assert!(out_rx.try_recv().is_err());
1106    }
1107
1108    #[rstest]
1109    #[case::microseconds_converted_to_ms(1_700_000_000_000_000_i64, 1_700_000_000_000_i64)]
1110    #[case::zero(0_i64, 0_i64)]
1111    #[case::negative(-1_000_i64, -1_i64)]
1112    fn test_parse_server_shutdown_event_time_ms(
1113        #[case] event_time_us: i64,
1114        #[case] expected_ms: i64,
1115    ) {
1116        let mut buf = vec![0u8; message_header_codec::ENCODED_LENGTH];
1117        buf.extend_from_slice(&event_time_us.to_le_bytes());
1118        assert_eq!(parse_server_shutdown_event_time_ms(&buf), expected_ms);
1119    }
1120
1121    #[rstest]
1122    fn test_parse_server_shutdown_event_time_ms_short_buffer_returns_zero() {
1123        let buf = vec![0u8; message_header_codec::ENCODED_LENGTH + 4];
1124        assert_eq!(parse_server_shutdown_event_time_ms(&buf), 0);
1125    }
1126
1127    #[rstest]
1128    fn test_classify_user_data_event_server_shutdown_emits_variant() {
1129        let event = serde_json::json!({"e": "serverShutdown", "E": 1_700_000_000_000_i64});
1130        let msg = classify_user_data_event(&event).expect("expected ServerShutdown");
1131        match msg {
1132            BinanceSpotWsTradingMessage::ServerShutdown { event_time } => {
1133                assert_eq!(event_time, 1_700_000_000_000);
1134            }
1135            other => panic!("expected ServerShutdown variant, was {other:?}"),
1136        }
1137    }
1138
1139    #[rstest]
1140    fn test_classify_user_data_event_server_shutdown_missing_event_time_defaults_to_zero() {
1141        let event = serde_json::json!({"e": "serverShutdown"});
1142        let msg = classify_user_data_event(&event).expect("expected ServerShutdown");
1143        match msg {
1144            BinanceSpotWsTradingMessage::ServerShutdown { event_time } => {
1145                assert_eq!(event_time, 0);
1146            }
1147            other => panic!("expected ServerShutdown variant, was {other:?}"),
1148        }
1149    }
1150
1151    #[rstest]
1152    fn test_classify_user_data_event_unknown_returns_none() {
1153        let event = serde_json::json!({"e": "somethingElse"});
1154        assert!(classify_user_data_event(&event).is_none());
1155    }
1156
1157    #[rstest]
1158    fn test_sign_params_includes_recv_window_in_signature() {
1159        let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
1160        let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
1161        let (out_tx, _out_rx) = tokio::sync::mpsc::unbounded_channel();
1162        let credential = Arc::new(SigningCredential::new(
1163            "api-key".to_string(),
1164            "hmac-secret".to_string(),
1165        ));
1166        let handler = BinanceSpotWsTradingHandler::new(
1167            Arc::new(AtomicBool::new(false)),
1168            cmd_rx,
1169            raw_rx,
1170            out_tx,
1171            credential.clone(),
1172        )
1173        .with_recv_window(Some(45_000));
1174
1175        let signed = handler
1176            .sign_params(serde_json::json!({"symbol": "BTCUSDT"}))
1177            .unwrap();
1178        let mut unsigned = signed.clone();
1179        let signature = unsigned
1180            .as_object_mut()
1181            .unwrap()
1182            .remove("signature")
1183            .unwrap();
1184        let query = canonical_ws_query_string(
1185            unsigned
1186                .as_object()
1187                .unwrap()
1188                .iter()
1189                .map(|(key, value)| (key.as_str(), value)),
1190        )
1191        .unwrap();
1192
1193        assert_eq!(signed["recvWindow"], 45_000);
1194        assert_eq!(signature, credential.sign(&query));
1195    }
1196}