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::{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
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").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 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 log::debug!("Unexpected JSON success for request {id_str}: {json}");
513 }
514 }
515 return;
516 }
517
518 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 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 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 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 _ => {} }
637 }
638
639 let (request_id, status, result_data) = self.parse_envelope(data)?;
641
642 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 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 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 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 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 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 let decoder = WebSocketResponseDecoder::default().wrap(
749 buf,
750 message_header_codec::ENCODED_LENGTH,
751 block_length,
752 version,
753 );
754
755 let status = decoder.status();
757
758 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 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 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 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
850pub(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
930pub(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}