1use 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
63pub 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 #[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 #[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 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 self.fail_pending_requests();
213
214 self.emit(BinanceSpotWsTradingMessage::Reconnected);
215 continue;
216 }
217
218 self.handle_message(msg);
219 }
220 else => {
221 return false;
223 }
224 }
225 }
226 }
227
228 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 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(¶ms)
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(¶ms)
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(¶ms)
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 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 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 if let Some(event) = json.get("event") {
443 self.handle_user_data_event(event);
444 return;
445 }
446
447 if json.get("e").is_some() {
449 self.handle_user_data_event(&json);
450 return;
451 }
452
453 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 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 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 log::debug!("Unexpected JSON success for request {id_str}");
527 }
528 }
529 return;
530 }
531
532 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 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 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 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 _ => {} }
651 }
652
653 let (request_id, status, result_data) = self.parse_envelope(data)?;
655
656 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 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 None => BinanceSpotWsTradingMessage::RequestFailed {
667 request_id,
668 msg: format!("Request failed with status {status}; error payload undecodable"),
669 },
670 });
671 }
672
673 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 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 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 let decoder = WebSocketResponseDecoder::default().wrap(
747 buf,
748 message_header_codec::ENCODED_LENGTH,
749 block_length,
750 version,
751 );
752
753 let status = decoder.status();
755
756 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 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 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 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
853pub(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
933pub(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}