1use std::{
19 future::Future,
20 sync::{
21 Arc,
22 atomic::{AtomicBool, Ordering},
23 },
24 time::Duration,
25};
26
27use ahash::{AHashMap, AHashSet};
28use anyhow::Context;
29use async_trait::async_trait;
30use jiff::Timestamp;
31use nautilus_common::{
32 cache::fifo::FifoCache,
33 clients::ExecutionClient,
34 live::runner::get_exec_event_sender,
35 messages::execution::{
36 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
37 GenerateFillReportsBuilder, GenerateOrderStatusReport, GenerateOrderStatusReports,
38 GenerateOrderStatusReportsBuilder, GeneratePositionStatusReports,
39 GeneratePositionStatusReportsBuilder, ModifyOrder, QueryAccount, QueryOrder, SubmitOrder,
40 SubmitOrderList,
41 },
42};
43use nautilus_core::{
44 Params, UUID4, UnixNanos,
45 datetime::{NANOSECONDS_IN_MILLISECOND, checked_mins_to_nanos},
46 time::{AtomicTime, get_atomic_clock_realtime},
47};
48use nautilus_live::{
49 ExecutionClientCore, ExecutionEventEmitter, SocketControlFactory,
50 task::{TaskGroup, TaskGroupGuard, TaskSpawner},
51};
52use nautilus_model::{
53 accounts::AccountAny,
54 enums::{ContingencyType, LiquiditySide, OmsType, OrderStatus, OrderType, TimeInForce},
55 events::{
56 AccountState, OrderAccepted, OrderCancelRejected, OrderCanceled, OrderEventAny,
57 OrderExpired, OrderFilled, OrderModifyRejected, OrderRejected, OrderUpdated,
58 },
59 identifiers::{
60 AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, Venue,
61 VenueOrderId,
62 },
63 instruments::Instrument,
64 orders::{Order, OrderAny},
65 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
66 types::{AccountBalance, Currency, MarginBalance, Money, Price, Quantity},
67};
68use parking_lot::Mutex;
69use rust_decimal::Decimal;
70use ustr::Ustr;
71
72use super::websocket::trading::{
73 client::BinanceSpotWsTradingClient,
74 messages::BinanceSpotWsTradingMessage,
75 parse::{
76 parse_spot_account_position, parse_spot_exec_report_to_fill,
77 parse_spot_exec_report_to_order_status,
78 },
79 user_data::{BinanceSpotExecutionReport, BinanceSpotExecutionType},
80};
81use crate::{
82 common::{
83 consts::{
84 BINANCE_GTX_ORDER_REJECT_CODE, BINANCE_NAUTILUS_SPOT_BROKER_ID,
85 BINANCE_NEW_ORDER_REJECTED_CODE, BINANCE_SPOT_POST_ONLY_REJECT_MSG,
86 BINANCE_SPOT_SBE_WS_API_DEMO_URL, BINANCE_SPOT_SBE_WS_API_TESTNET_URL,
87 BINANCE_SPOT_SBE_WS_API_URL, BINANCE_STATUS_UNKNOWN_CODE,
88 BINANCE_UNEXPECTED_RESPONSE_CODE, BINANCE_VENUE, BINANCE_WS_HEARTBEAT_SECS,
89 },
90 credential::resolve_credentials,
91 dispatch::{
92 OrderIdentity, PendingOperation, PendingRequest, WsDispatchState,
93 ensure_accepted_emitted,
94 },
95 encoder::{decode_client_order_id, encode_broker_id},
96 enums::{BinanceEnvironment, BinanceSide, BinanceTimeInForce},
97 parse::{
98 parse_millis_or_init, parse_required_decimal, parse_required_price_at_precision,
99 parse_required_quantity_at_precision,
100 },
101 urls::{get_http_base_url_with_us, get_spot_user_stream_url},
102 },
103 config::BinanceExecutionClientConfig,
104 spot::{
105 enums::{
106 BinanceCancelReplaceMode, BinanceOrderResponseType, BinanceSpotOrderType,
107 order_type_to_binance_spot, time_in_force_to_binance_spot,
108 },
109 http::{
110 client::BinanceSpotHttpClient,
111 error::BinanceSpotHttpError,
112 models::BatchCancelResult,
113 query::{
114 BatchCancelItem, CancelOrderParams, CancelReplaceOrderParams,
115 NewOcoOrderListParams, NewOrderParams,
116 },
117 },
118 },
119};
120
121const ACCOUNT_TRADES_MAX_INTERVAL_MS: i64 = 24 * 60 * 60 * 1_000;
122
123const ACCOUNT_TRADES_PAGE_LIMIT: u32 = 1_000;
124
125const WS_RECONNECT_SETUP_RETRY_DELAY: Duration = Duration::from_secs(1);
126
127#[derive(Debug)]
133pub struct BinanceSpotExecutionClient {
134 core: ExecutionClientCore,
135 clock: &'static AtomicTime,
136 config: BinanceExecutionClientConfig,
137 emitter: ExecutionEventEmitter,
138 dispatch_state: Arc<WsDispatchState>,
139 http_client: BinanceSpotHttpClient,
140 socket_factory: SocketControlFactory,
141 ws_trading_client: Option<BinanceSpotWsTradingClient>,
142 ws_trading_dispatch_active: Arc<AtomicBool>,
143 ws_user_data_client: Arc<Mutex<Option<BinanceSpotWsTradingClient>>>,
144 ws_user_data_dispatch_active: Arc<AtomicBool>,
145 listen_key: Option<String>,
146 us_credentials: Option<(String, String)>,
147 ws_authenticated: Arc<tokio::sync::Notify>,
148 ws_user_data_subscribed: Arc<tokio::sync::Notify>,
149 session_tasks: TaskGroup,
150 pending_tasks: TaskGroup,
151 shutdown_errors: Vec<String>,
152}
153
154impl BinanceSpotExecutionClient {
155 pub fn new(
161 core: ExecutionClientCore,
162 config: BinanceExecutionClientConfig,
163 ) -> anyhow::Result<Self> {
164 config.validate()?;
165 let (api_key, api_secret) = resolve_credentials(
166 config.api_key.clone(),
167 config.api_secret.clone(),
168 config.environment,
169 config.product_type,
170 )?;
171
172 let clock = get_atomic_clock_realtime();
173 let socket_factory = SocketControlFactory::new(core.client_id, Some(*BINANCE_VENUE));
174 let base_url_http = config.base_url_http.clone().or_else(|| {
175 config.us.then(|| {
176 get_http_base_url_with_us(config.product_type, config.environment, true).to_string()
177 })
178 });
179
180 let http_client = BinanceSpotHttpClient::new_with_json_responses(
181 config.environment,
182 clock,
183 Some(api_key.clone()),
184 Some(api_secret.clone()),
185 base_url_http,
186 Some(config.recv_window_ms),
187 None, config.proxy_url.clone(),
189 config.us,
190 )
191 .context("failed to construct Binance Spot HTTP client")?;
192 let emitter = ExecutionEventEmitter::new(
193 clock,
194 core.trader_id,
195 core.account_id,
196 core.account_type,
197 core.base_currency,
198 );
199
200 let ws_trading_client = if config.us {
201 None
202 } else {
203 let url = Some(Self::resolve_ws_trading_url(
204 config.base_url_ws_trading.clone(),
205 config.environment,
206 ));
207 Some(
208 BinanceSpotWsTradingClient::new(
209 url,
210 api_key.clone(),
211 api_secret.clone(),
212 Some(BINANCE_WS_HEARTBEAT_SECS),
213 config.transport_backend,
214 )
215 .with_proxy(config.proxy_url.clone())
216 .with_recv_window(Some(config.recv_window_ms))
217 .with_socket_control(socket_factory.control("binance-spot-trading")),
218 )
219 };
220 let us_credentials = config.us.then_some((api_key, api_secret));
221
222 let session_tasks = TaskGroup::new();
223 let pending_tasks = TaskGroup::new();
224
225 Ok(Self {
226 core,
227 clock,
228 config,
229 emitter,
230 dispatch_state: Arc::new(WsDispatchState::default()),
231 http_client,
232 socket_factory,
233 ws_trading_client,
234 ws_trading_dispatch_active: Arc::new(AtomicBool::new(false)),
235 ws_user_data_client: Arc::new(Mutex::new(None)),
236 ws_user_data_dispatch_active: Arc::new(AtomicBool::new(false)),
237 listen_key: None,
238 us_credentials,
239 ws_authenticated: Arc::new(tokio::sync::Notify::new()),
240 ws_user_data_subscribed: Arc::new(tokio::sync::Notify::new()),
241 session_tasks,
242 pending_tasks,
243 shutdown_errors: Vec::new(),
244 })
245 }
246
247 fn resolve_ws_trading_url(base_url: Option<String>, environment: BinanceEnvironment) -> String {
248 base_url.unwrap_or_else(|| {
249 match environment {
250 BinanceEnvironment::Live => BINANCE_SPOT_SBE_WS_API_URL,
251 BinanceEnvironment::Testnet => BINANCE_SPOT_SBE_WS_API_TESTNET_URL,
252 BinanceEnvironment::Demo => BINANCE_SPOT_SBE_WS_API_DEMO_URL,
253 }
254 .to_string()
255 })
256 }
257
258 async fn refresh_account_state(&self) -> anyhow::Result<AccountState> {
259 self.http_client
260 .request_account_state(self.core.account_id)
261 .await
262 }
263
264 fn update_account_state(&self) {
265 let http_client = self.http_client.clone();
266 let account_id = self.core.account_id;
267 let emitter = self.emitter.clone();
268 let clock = self.clock;
269
270 self.spawn_task("query_account", async move {
271 let account_state = http_client.request_account_state(account_id).await?;
272 let ts_now = clock.get_time_ns();
273 emitter.emit_account_state(
274 account_state.balances.clone(),
275 account_state.margins.clone(),
276 account_state.is_reported,
277 ts_now,
278 account_state.info,
279 );
280 Ok(())
281 });
282 }
283
284 fn ws_user_data_active(&self) -> bool {
285 let dispatch_running = if self.config.us {
286 self.ws_user_data_dispatch_active.load(Ordering::Acquire)
287 } else {
288 self.ws_trading_dispatch_active.load(Ordering::Acquire)
289 };
290 let user_data_active = if self.config.us {
291 self.ws_user_data_client
292 .lock()
293 .as_ref()
294 .is_some_and(BinanceSpotWsTradingClient::is_user_data_active)
295 } else {
296 self.ws_trading_client
297 .as_ref()
298 .is_some_and(BinanceSpotWsTradingClient::is_user_data_active)
299 };
300
301 user_data_active && dispatch_running
302 }
303
304 fn ensure_ws_user_data_active(&self) -> anyhow::Result<()> {
305 anyhow::ensure!(
306 self.ws_user_data_active(),
307 "Binance Spot user data stream is not active",
308 );
309 Ok(())
310 }
311
312 fn ws_order_transport_active(&self) -> bool {
313 self.config.use_ws_trading && self.ws_trading_client.is_some() && self.ws_user_data_active()
314 }
315
316 fn submit_order_internal(&self, cmd: &SubmitOrder) -> anyhow::Result<()> {
317 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
318
319 let event_emitter = self.emitter.clone();
320 let trader_id = self.core.trader_id;
321 let account_id = self.core.account_id;
322 let client_order_id = order.client_order_id();
323 let strategy_id = order.strategy_id();
324 let instrument_id = order.instrument_id();
325 let order_side = order.order_side();
326 let order_type = order.order_type();
327 let quantity = order.quantity();
328 let time_in_force = order.time_in_force();
329 let price = order.price();
330 let trigger_price = order.trigger_price();
331 let is_post_only = order.is_post_only();
332 let is_quote_quantity = order.is_quote_quantity();
333 let display_qty = order.display_qty();
334 let use_gtd = self.config.use_gtd;
335 let clock = self.clock;
336 let ts_init = self.clock.get_time_ns();
337
338 self.dispatch_state.order_identities.insert(
340 client_order_id,
341 OrderIdentity {
342 instrument_id,
343 strategy_id,
344 order_side,
345 order_type,
346 price,
347 quantity,
348 venue_position_id: None,
349 },
350 );
351
352 if self.ws_order_transport_active() {
353 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
354 let dispatch_state = self.dispatch_state.clone();
355 let params = build_new_order_params(
356 &order,
357 client_order_id,
358 is_post_only,
359 is_quote_quantity,
360 use_gtd,
361 )?;
362
363 let request_id = ws_client.next_request_id();
365 dispatch_state.pending_requests.insert(
366 request_id.clone(),
367 PendingRequest {
368 client_order_id,
369 venue_order_id: None,
370 operation: PendingOperation::Place,
371 },
372 );
373
374 self.spawn_task("submit_order_ws", async move {
375 if let Err(e) = ws_client
376 .place_order_with_id(request_id.clone(), params)
377 .await
378 {
379 dispatch_state.pending_requests.remove(&request_id);
380 log::warn!(
381 "WS submit request failed for {client_order_id}, awaiting reconciliation: {e}"
382 );
383 anyhow::bail!("WS submit order failed: {e}");
384 }
385 Ok(())
386 });
387 } else {
388 let http_client = self.http_client.clone();
389 let dispatch_state = self.dispatch_state.clone();
390 log::debug!("WS trading not active, falling back to HTTP for submit_order");
391
392 self.spawn_task("submit_order_http", async move {
393 let result = http_client
394 .submit_order(
395 account_id,
396 instrument_id,
397 client_order_id,
398 order_side,
399 order_type,
400 quantity,
401 time_in_force,
402 price,
403 trigger_price,
404 is_post_only,
405 is_quote_quantity,
406 display_qty,
407 use_gtd,
408 )
409 .await;
410
411 match result {
412 Ok(report) => handle_spot_order_submit_success(
413 client_order_id,
414 report.venue_order_id,
415 ),
416 Err(e) => {
417 if is_ambiguous_submit_error(&e) {
418 log::warn!(
419 "Ambiguous submit failure for {client_order_id}, awaiting reconciliation: {e}"
420 );
421 } else if is_structured_venue_rejection(&e)
422 || is_local_command_failure(&e)
423 {
424 let due_post_only = e
425 .downcast_ref::<BinanceSpotHttpError>()
426 .is_some_and(is_spot_post_only_rejection);
427 dispatch_state.cleanup_terminal(client_order_id);
428 let rejected = OrderRejected::new(
429 trader_id,
430 strategy_id,
431 instrument_id,
432 client_order_id,
433 account_id,
434 format!("submit-order-error: {e}").into(),
435 UUID4::new(),
436 ts_init,
437 clock.get_time_ns(),
438 false,
439 due_post_only,
440 );
441 event_emitter.send_order_event(OrderEventAny::Rejected(rejected));
442 } else {
443 log::warn!(
444 "Ambiguous submit failure for {client_order_id}, awaiting reconciliation: {e}"
445 );
446 }
447 return Err(e);
448 }
449 }
450 Ok(())
451 });
452 }
453
454 Ok(())
455 }
456
457 fn cancel_order_internal(&self, cmd: &CancelOrder) {
458 let event_emitter = self.emitter.clone();
459 let trader_id = self.core.trader_id;
460 let account_id = self.core.account_id;
461 let clock = self.clock;
462 let command = cmd.clone();
463
464 if self.ws_order_transport_active() {
465 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
466 let dispatch_state = self.dispatch_state.clone();
467 let params = build_cancel_order_params(&command);
468
469 let request_id = ws_client.next_request_id();
471 dispatch_state.pending_requests.insert(
472 request_id.clone(),
473 PendingRequest {
474 client_order_id: command.client_order_id,
475 venue_order_id: command.venue_order_id,
476 operation: PendingOperation::Cancel,
477 },
478 );
479
480 self.spawn_task("cancel_order_ws", async move {
481 if let Err(e) = ws_client
482 .cancel_order_with_id(request_id.clone(), params)
483 .await
484 {
485 dispatch_state.pending_requests.remove(&request_id);
486 log::warn!(
487 "WS cancel request failed for {}, awaiting reconciliation: {e}",
488 command.client_order_id
489 );
490 anyhow::bail!("WS cancel order failed: {e}");
491 }
492 Ok(())
493 });
494 } else {
495 let http_client = self.http_client.clone();
496 let dispatch_state = self.dispatch_state.clone();
497 log::debug!("WS trading not active, falling back to HTTP for cancel_order");
498
499 self.spawn_task("cancel_order_http", async move {
500 let result = http_client
501 .cancel_order(
502 command.instrument_id,
503 command.venue_order_id,
504 Some(command.client_order_id),
505 )
506 .await;
507
508 match result {
509 Ok(venue_order_id) => {
510 dispatch_state.cleanup_terminal(command.client_order_id);
511 let ts_now = clock.get_time_ns();
512 let canceled_event = OrderCanceled::new(
513 trader_id,
514 command.strategy_id,
515 command.instrument_id,
516 command.client_order_id,
517 UUID4::new(),
518 ts_now,
519 ts_now,
520 false,
521 Some(venue_order_id),
522 Some(account_id),
523 );
524 event_emitter.send_order_event(OrderEventAny::Canceled(canceled_event));
525 }
526 Err(e) => {
527 if is_structured_venue_rejection(&e) {
528 let ts_now = clock.get_time_ns();
529 let rejected_event = OrderCancelRejected::new(
530 trader_id,
531 command.strategy_id,
532 command.instrument_id,
533 command.client_order_id,
534 format!("cancel-order-error: {e}").into(),
535 UUID4::new(),
536 ts_now,
537 ts_now,
538 false,
539 command.venue_order_id,
540 Some(account_id),
541 );
542 event_emitter
543 .send_order_event(OrderEventAny::CancelRejected(rejected_event));
544 } else if is_local_command_failure(&e) {
545 log::warn!(
546 "Cancel command failed local validation for {}: {e}",
547 command.client_order_id
548 );
549 } else {
550 log::warn!(
551 "Ambiguous cancel failure for {}, awaiting reconciliation: {e}",
552 command.client_order_id
553 );
554 }
555 return Err(e);
556 }
557 }
558 Ok(())
559 });
560 }
561 }
562
563 fn spawn_task<F>(&self, description: &'static str, fut: F)
564 where
565 F: Future<Output = anyhow::Result<()>> + Send + 'static,
566 {
567 crate::common::execution::spawn_task(&self.pending_tasks, description, fut);
568 }
569
570 fn abort_pending_tasks(&self) {
571 crate::common::execution::abort_pending_tasks(&self.pending_tasks);
572 }
573
574 fn abort_session_tasks(&self) {
575 self.session_tasks.begin_shutdown();
576 }
577
578 async fn await_pending_tasks(&self) -> anyhow::Result<()> {
579 crate::common::execution::await_pending_tasks(&self.pending_tasks).await
580 }
581
582 async fn await_session_tasks(&self) -> anyhow::Result<()> {
583 self.session_tasks.begin_shutdown();
584 self.session_tasks
585 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
586 .await
587 .map_err(|e| anyhow::anyhow!("Failed to terminate Binance Spot session tasks: {e}"))?;
588 Ok(())
589 }
590
591 async fn ws_setup_failure(
592 &mut self,
593 mut ws_trading: BinanceSpotWsTradingClient,
594 reason: String,
595 ) -> anyhow::Error {
596 ws_trading.mark_user_data_inactive();
597 log::error!("{reason}; Binance Spot private user data is required for execution");
598
599 self.abort_session_tasks();
600
601 if let Err(e) = ws_trading.disconnect().await {
602 log::warn!("Failed to stop Binance Spot trading WebSocket after setup failure: {e}");
603 }
604 self.disconnect_us_user_data().await;
605
606 if let Err(e) = self.await_session_tasks().await {
607 log::warn!("Failed to drain Binance Spot session tasks after setup failure: {e}");
608 }
609 self.ws_trading_client = Some(ws_trading);
610 anyhow::anyhow!(reason)
611 }
612
613 async fn connect_us_user_data(&mut self) -> anyhow::Result<()> {
614 let (api_key, api_secret) = self
615 .us_credentials
616 .clone()
617 .context("Binance US user data credentials are unavailable")?;
618 let listen_key = self
619 .http_client
620 .inner()
621 .create_listen_key()
622 .await
623 .context("failed to create Binance US listen key")?
624 .listen_key;
625 self.listen_key = Some(listen_key.clone());
626 let url = get_spot_user_stream_url(self.config.base_url_ws.as_deref(), &listen_key);
627 let mut ws_user_data = BinanceSpotWsTradingClient::new(
628 Some(url),
629 api_key,
630 api_secret,
631 Some(BINANCE_WS_HEARTBEAT_SECS),
632 self.config.transport_backend,
633 )
634 .with_proxy(self.config.proxy_url.clone())
635 .with_socket_control(self.socket_factory.control("binance-spot-user-streams"));
636 *self.ws_user_data_client.lock() = Some(ws_user_data.clone());
637 ws_user_data
638 .connect()
639 .await
640 .context("failed to connect Binance US user data stream")?;
641
642 let ws_clone = ws_user_data.clone();
643 let emitter = self.emitter.clone();
644 let account_id = self.core.account_id;
645 let clock = self.clock;
646 let http_client = self.http_client.clone();
647 let dispatch_state = self.dispatch_state.clone();
648 let treat_expired_as_canceled = self.config.treat_expired_as_canceled;
649 let ws_authenticated = self.ws_authenticated.clone();
650 let ws_user_data_subscribed = self.ws_user_data_subscribed.clone();
651 let (setup_error_tx, _setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
652 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
653 let dispatch_active = Arc::clone(&self.ws_user_data_dispatch_active);
654 let task_spawner = self
655 .session_tasks
656 .spawner()
657 .context("Binance Spot session task admission is closed")?;
658
659 if let Err(e) = self.session_tasks.spawn(async move {
660 dispatch_active.store(true, Ordering::Release);
661 let _active = DispatchActiveGuard(dispatch_active);
662
663 while let Some(message) = ws_clone.recv().await {
664 if matches!(&message, BinanceSpotWsTradingMessage::Reconnected) {
665 ws_clone.mark_user_data_active();
666 }
667 dispatch_ws_trading_message(
668 message,
669 &emitter,
670 &http_client,
671 account_id,
672 treat_expired_as_canceled,
673 clock,
674 &dispatch_state,
675 &ws_authenticated,
676 &ws_user_data_subscribed,
677 &setup_error_tx,
678 &seen_trade_ids,
679 &task_spawner,
680 );
681 }
682 log::warn!("Binance US user data dispatch loop ended");
683 }) {
684 return Err(e.into());
685 }
686
687 let keepalive_http = self.http_client.clone();
688 let keepalive_key = listen_key.clone();
689
690 if let Err(e) = self.session_tasks.spawn(async move {
691 let mut interval = tokio::time::interval(Duration::from_secs(30 * 60));
692 interval.tick().await;
693
694 loop {
695 interval.tick().await;
696
697 if let Err(e) = keepalive_http
698 .inner()
699 .extend_listen_key(&keepalive_key)
700 .await
701 {
702 log::warn!("Binance US listen key keepalive failed: {e}");
703 }
704 }
705 }) {
706 return Err(e.into());
707 }
708
709 ws_user_data.mark_user_data_active();
710 *self.ws_user_data_client.lock() = Some(ws_user_data);
711 Ok(())
712 }
713
714 async fn disconnect_us_user_data(&mut self) {
715 let mut client_drained = true;
716 let client = self.ws_user_data_client.lock().clone();
717
718 if let Some(mut client) = client {
719 client.mark_user_data_inactive();
720 if let Err(e) = client.disconnect().await {
721 client_drained = false;
722 self.shutdown_errors.push(format!(
723 "failed to stop Binance US user data WebSocket: {e}"
724 ));
725 }
726 }
727
728 if client_drained {
729 *self.ws_user_data_client.lock() = None;
730 }
731
732 if let Some(listen_key) = self.listen_key.clone() {
733 match self.http_client.inner().close_listen_key(&listen_key).await {
734 Ok(()) if self.listen_key.as_deref() == Some(listen_key.as_str()) => {
735 self.listen_key = None;
736 }
737 Ok(()) => {}
738 Err(e) => self
739 .shutdown_errors
740 .push(format!("failed to close Binance US listen key: {e}")),
741 }
742 }
743 }
744
745 async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
746 if let Some(client) = self.ws_trading_client.as_ref() {
747 client.mark_user_data_inactive();
748 client.begin_shutdown();
749 }
750
751 if let Some(client) = self.ws_user_data_client.lock().as_ref() {
752 client.mark_user_data_inactive();
753 client.begin_shutdown();
754 }
755
756 self.abort_session_tasks();
757 self.abort_pending_tasks();
758
759 if let Some(ref mut ws_trading) = self.ws_trading_client
760 && let Err(e) = ws_trading.disconnect().await
761 {
762 self.shutdown_errors
763 .push(format!("trading WebSocket shutdown failed: {e}"));
764 }
765 self.disconnect_us_user_data().await;
766
767 let (session_result, pending_result) =
768 tokio::join!(self.await_session_tasks(), self.await_pending_tasks());
769 self.core.set_disconnected();
770
771 if let Err(e) = session_result {
772 self.shutdown_errors.push(e.to_string());
773 }
774
775 if let Err(e) = pending_result {
776 self.shutdown_errors.push(e.to_string());
777 }
778
779 if !self.shutdown_errors.is_empty() {
780 let errors = std::mem::take(&mut self.shutdown_errors);
781 anyhow::bail!("Binance Spot shutdown failed: {}", errors.join("; "));
782 }
783 Ok(())
784 }
785}
786
787#[async_trait(?Send)]
788impl ExecutionClient for BinanceSpotExecutionClient {
789 fn is_connected(&self) -> bool {
790 self.core.is_connected()
791 }
792
793 fn client_id(&self) -> ClientId {
794 self.core.client_id
795 }
796
797 fn account_id(&self) -> AccountId {
798 self.core.account_id
799 }
800
801 fn venue(&self) -> Venue {
802 *BINANCE_VENUE
803 }
804
805 fn oms_type(&self) -> OmsType {
806 self.core.oms_type
807 }
808
809 fn get_account(&self) -> Option<AccountAny> {
810 self.core.cache().account_owned(&self.core.account_id)
811 }
812
813 async fn connect(&mut self) -> anyhow::Result<()> {
814 if self.core.is_connected() && self.session_tasks.is_open() && self.pending_tasks.is_open()
815 {
816 return Ok(());
817 }
818
819 if !self.pending_tasks.is_open() || !self.session_tasks.is_open() {
820 self.teardown_partial_connect().await?;
821 }
822
823 if !self.pending_tasks.is_open() {
824 self.await_pending_tasks().await?;
825 self.pending_tasks.start_generation().map_err(|e| {
826 anyhow::anyhow!("Failed to start Binance Spot task generation: {e}")
827 })?;
828 }
829
830 if !self.session_tasks.is_open() {
831 self.await_session_tasks().await?;
832 self.session_tasks.start_generation().map_err(|e| {
833 anyhow::anyhow!("Failed to start Binance Spot session generation: {e}")
834 })?;
835 }
836 let ws_trading_client = self.ws_trading_client.clone();
837 let ws_user_data_client = Arc::clone(&self.ws_user_data_client);
838 let setup_guard =
839 TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
840 if let Some(client) = ws_trading_client {
841 client.begin_shutdown();
842 }
843
844 if let Some(client) = ws_user_data_client.lock().as_ref() {
845 client.begin_shutdown();
846 }
847 });
848
849 let ws_setup_timeout = Duration::from_millis(self.config.ws_trading_setup_timeout_ms);
850
851 if !self.core.instruments_initialized() {
853 let instruments = self
854 .http_client
855 .request_instruments_with_config(&self.config.instrument_provider, self.config.us)
856 .await
857 .context("failed to request Binance Spot instruments")?;
858
859 if instruments.is_empty() {
860 log::warn!("No instruments returned for Binance Spot");
861 } else {
862 log::debug!("Loaded {} Spot instruments", instruments.len());
863 self.http_client.cache_instruments(instruments);
864 }
865
866 self.core.set_instruments_initialized();
867 }
868
869 let account_state = self
871 .refresh_account_state()
872 .await
873 .context("failed to request Binance account state")?;
874
875 if !account_state.balances.is_empty() {
876 log::debug!(
877 "Received account state with {} balance(s)",
878 account_state.balances.len()
879 );
880 }
881
882 self.emitter.send_account_state(account_state);
883
884 crate::common::execution::await_account_registered(&self.core, self.core.account_id, 30.0)
886 .await?;
887
888 let session_result = async {
889 if self.config.us {
890 self.connect_us_user_data().await?;
891 }
892
893 if let Some(mut ws_trading) = self.ws_trading_client.clone() {
894 match ws_trading.connect().await {
895 Ok(()) => {
896 log::debug!("Connected to Binance Spot WS trading API");
897
898 let ws_trading_clone = ws_trading.clone();
899 let emitter = self.emitter.clone();
900 let account_id = self.core.account_id;
901 let clock = self.clock;
902 let http_client = self.http_client.clone();
903 let dispatch_state = self.dispatch_state.clone();
904 let treat_expired_as_canceled = self.config.treat_expired_as_canceled;
905 let ws_authenticated = self.ws_authenticated.clone();
906 let ws_user_data_subscribed = self.ws_user_data_subscribed.clone();
907 let (ws_setup_error_tx, mut ws_setup_error_rx) =
908 tokio::sync::mpsc::unbounded_channel();
909 let seen_trade_ids = std::sync::Arc::new(Mutex::new(FifoCache::new()));
910 let dispatch_active = Arc::clone(&self.ws_trading_dispatch_active);
911 let task_spawner = self
912 .session_tasks
913 .spawner()
914 .context("Binance Spot session task admission is closed")?;
915
916 self.session_tasks.spawn(async move {
917 dispatch_active.store(true, Ordering::Release);
918 let _active = DispatchActiveGuard(dispatch_active);
919 let mut resubscribing = false;
920
921 loop {
922 match ws_trading_clone.recv().await {
923 Some(msg) => {
924 match &msg {
925 BinanceSpotWsTradingMessage::Reconnected => {
926 ws_trading_clone.mark_user_data_inactive();
927 resubscribing = true;
928 if let Err(e) = ws_trading_clone.session_logon().await {
929 resubscribing = false;
930 log::error!(
931 "Failed to re-authenticate Binance Spot user data stream: {e}"
932 );
933 }
934 }
935 BinanceSpotWsTradingMessage::Authenticated if resubscribing => {
936 if let Err(e) =
937 ws_trading_clone.subscribe_user_data().await
938 {
939 resubscribing = false;
940 log::error!(
941 "Failed to resubscribe Binance Spot user data stream: {e}"
942 );
943 }
944 continue;
945 }
946 BinanceSpotWsTradingMessage::UserDataSubscribed { .. } => {
947 let was_resubscribing = resubscribing;
948 resubscribing = false;
949 ws_trading_clone.mark_user_data_active();
950
951 if was_resubscribing {
952 continue;
953 }
954 }
955 BinanceSpotWsTradingMessage::AuthenticationRejected(reason)
956 if resubscribing =>
957 {
958 log::warn!(
959 "Binance Spot reconnect authentication failed; retrying: {reason}"
960 );
961 tokio::time::sleep(WS_RECONNECT_SETUP_RETRY_DELAY).await;
962
963 if let Err(e) = ws_trading_clone.session_logon().await {
964 resubscribing = false;
965 log::error!(
966 "Failed to retry Binance Spot reconnect authentication: {e}"
967 );
968 }
969 continue;
970 }
971 BinanceSpotWsTradingMessage::UserDataSubscriptionRejected(reason)
972 if resubscribing =>
973 {
974 log::warn!(
975 "Binance Spot reconnect user data subscription failed; retrying: {reason}"
976 );
977 tokio::time::sleep(WS_RECONNECT_SETUP_RETRY_DELAY).await;
978
979 if let Err(e) =
980 ws_trading_clone.subscribe_user_data().await
981 {
982 resubscribing = false;
983 log::error!(
984 "Failed to retry Binance Spot user data subscription: {e}"
985 );
986 }
987 continue;
988 }
989 _ => {}
990 }
991
992 dispatch_ws_trading_message(
993 msg,
994 &emitter,
995 &http_client,
996 account_id,
997 treat_expired_as_canceled,
998 clock,
999 &dispatch_state,
1000 &ws_authenticated,
1001 &ws_user_data_subscribed,
1002 &ws_setup_error_tx,
1003 &seen_trade_ids,
1004 &task_spawner,
1005 );
1006 }
1007 None => {
1008 log::warn!("WS trading dispatch loop ended");
1009 break;
1010 }
1011 }
1012 }
1013 })?;
1014
1015 if let Err(e) = ws_trading.session_logon().await {
1016 let reason = format!("WS session logon failed: {e}");
1017 return Err(self.ws_setup_failure(ws_trading, reason).await);
1018 } else {
1019 let auth_result = wait_for_ws_setup_response(
1020 ws_setup_timeout,
1021 self.ws_authenticated.notified(),
1022 &mut ws_setup_error_rx,
1023 "WS session authentication timed out",
1024 )
1025 .await;
1026
1027 if let Err(e) = auth_result {
1028 return Err(self.ws_setup_failure(ws_trading, e.to_string()).await);
1029 } else if let Err(e) = ws_trading.subscribe_user_data().await {
1030 let reason = format!("WS user data subscribe failed: {e}");
1031 return Err(self.ws_setup_failure(ws_trading, reason).await);
1032 } else {
1033 let subscribe_result = wait_for_ws_setup_response(
1034 ws_setup_timeout,
1035 self.ws_user_data_subscribed.notified(),
1036 &mut ws_setup_error_rx,
1037 "WS user data subscription timed out",
1038 )
1039 .await;
1040
1041 if let Err(e) = subscribe_result {
1042 return Err(self.ws_setup_failure(ws_trading, e.to_string()).await);
1043 } else {
1044 self.ws_trading_client = Some(ws_trading);
1045 }
1046 }
1047 }
1048 }
1049 Err(e) => {
1050 let reason = format!("Failed to connect WS trading API: {e}");
1051 return Err(self.ws_setup_failure(ws_trading, reason).await);
1052 }
1053 }
1054 }
1055
1056 let refresh_secs = self.config.instrument_refresh_interval_secs;
1057 if refresh_secs > 0 {
1058 let http_client = self.http_client.clone();
1059 let provider = self.config.instrument_provider.clone();
1060 let us = self.config.us;
1061
1062 self.session_tasks.spawn(async move {
1063 let mut interval = tokio::time::interval(Duration::from_secs(refresh_secs));
1064 interval.tick().await;
1065
1066 loop {
1067 interval.tick().await;
1068
1069 match http_client
1070 .request_instruments_with_config(&provider, us)
1071 .await
1072 {
1073 Ok(instruments) => log::debug!(
1074 "Refreshed Binance Spot execution instruments: count={}",
1075 instruments.len()
1076 ),
1077 Err(e) => {
1078 log::warn!("Binance Spot execution instrument refresh failed: {e}");
1079 }
1080 }
1081 }
1082 })?;
1083 }
1084
1085 Ok::<(), anyhow::Error>(())
1086 }
1087 .await;
1088
1089 if let Err(e) = session_result {
1090 if let Err(teardown_error) = self.teardown_partial_connect().await {
1091 return Err(e.context(format!(
1092 "Binance Spot startup teardown failed: {teardown_error}"
1093 )));
1094 }
1095 return Err(e);
1096 }
1097
1098 setup_guard.disarm();
1099 self.core.set_connected();
1100 log::info!("Connected: client_id={}", self.core.client_id);
1101 Ok(())
1102 }
1103
1104 async fn disconnect(&mut self) -> anyhow::Result<()> {
1105 self.teardown_partial_connect().await?;
1106 log::info!("Disconnected: client_id={}", self.core.client_id);
1107 Ok(())
1108 }
1109
1110 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
1111 self.update_account_state();
1112 Ok(())
1113 }
1114
1115 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
1116 log::debug!("query_order: client_order_id={}", cmd.client_order_id);
1117
1118 let http_client = self.http_client.clone();
1119 let command = cmd;
1120 let event_emitter = self.emitter.clone();
1121 let account_id = self.core.account_id;
1122 let treat_expired_as_canceled = self.config.treat_expired_as_canceled;
1123
1124 self.spawn_task("query_order", async move {
1125 let result = http_client
1126 .request_order_status_report(
1127 account_id,
1128 command.instrument_id,
1129 command.venue_order_id,
1130 Some(command.client_order_id),
1131 )
1132 .await;
1133
1134 match result {
1135 Ok(Some(mut report)) => {
1136 normalize_spot_order_status_report(&mut report, treat_expired_as_canceled);
1137 event_emitter.send_order_status_report(report);
1138 }
1139 Ok(None) => log::debug!(
1140 "No order status report returned: client_order_id={}",
1141 command.client_order_id
1142 ),
1143 Err(e) => log::warn!("Failed to query order status: {e}"),
1144 }
1145
1146 Ok(())
1147 });
1148
1149 Ok(())
1150 }
1151
1152 fn generate_account_state(
1153 &self,
1154 balances: Vec<AccountBalance>,
1155 margins: Vec<MarginBalance>,
1156 reported: bool,
1157 ts_event: UnixNanos,
1158 info: Option<Params>,
1159 ) -> anyhow::Result<()> {
1160 self.emitter
1161 .emit_account_state(balances, margins, reported, ts_event, info);
1162 Ok(())
1163 }
1164
1165 fn start(&mut self) -> anyhow::Result<()> {
1166 if self.core.is_started() {
1167 return Ok(());
1168 }
1169
1170 self.emitter.set_sender(get_exec_event_sender());
1171 self.core.set_started();
1172
1173 let http_client = self.http_client.clone();
1175 let provider = self.config.instrument_provider.clone();
1176 let us = self.config.us;
1177
1178 self.session_tasks.spawn(async move {
1179 match http_client
1180 .request_instruments_with_config(&provider, us)
1181 .await
1182 {
1183 Ok(instruments) => {
1184 if instruments.is_empty() {
1185 log::warn!("No instruments returned for Binance Spot");
1186 } else {
1187 http_client.cache_instruments(instruments);
1188 log::debug!("Instruments initialized");
1189 }
1190 }
1191 Err(e) => {
1192 log::error!("Failed to request Binance Spot instruments: {e}");
1193 }
1194 }
1195 })?;
1196
1197 log::info!(
1198 "Started: client_id={}, account_id={}, account_type={:?}, environment={:?}, product_type={:?}",
1199 self.core.client_id,
1200 self.core.account_id,
1201 self.core.account_type,
1202 self.config.environment,
1203 self.config.product_type,
1204 );
1205 Ok(())
1206 }
1207
1208 fn stop(&mut self) -> anyhow::Result<()> {
1209 if self.core.is_stopped() {
1210 return Ok(());
1211 }
1212
1213 if let Some(client) = self.ws_trading_client.as_ref() {
1214 client.mark_user_data_inactive();
1215 client.begin_shutdown();
1216 }
1217
1218 if let Some(client) = self.ws_user_data_client.lock().as_ref() {
1219 client.mark_user_data_inactive();
1220 client.begin_shutdown();
1221 }
1222
1223 self.core.set_stopped();
1224 self.core.set_disconnected();
1225 self.abort_session_tasks();
1226 self.abort_pending_tasks();
1227 log::info!("Stopped: client_id={}", self.core.client_id);
1228 Ok(())
1229 }
1230
1231 async fn generate_order_status_report(
1232 &self,
1233 cmd: &GenerateOrderStatusReport,
1234 ) -> anyhow::Result<Option<OrderStatusReport>> {
1235 let Some(instrument_id) = cmd.instrument_id else {
1236 log::warn!("generate_order_status_report requires instrument_id: {cmd:?}");
1237 return Ok(None);
1238 };
1239
1240 let venue_order_id = cmd
1242 .venue_order_id
1243 .as_ref()
1244 .map(|id| VenueOrderId::new(id.inner()));
1245
1246 let report = self
1247 .http_client
1248 .request_order_status_report(
1249 self.core.account_id,
1250 instrument_id,
1251 venue_order_id,
1252 cmd.client_order_id,
1253 )
1254 .await?;
1255
1256 Ok(report.map(|mut report| {
1257 normalize_spot_order_status_report(&mut report, self.config.treat_expired_as_canceled);
1258 report
1259 }))
1260 }
1261
1262 async fn generate_order_status_reports(
1263 &self,
1264 cmd: &GenerateOrderStatusReports,
1265 ) -> anyhow::Result<Vec<OrderStatusReport>> {
1266 let start_dt = cmd.start.map(|nanos| nanos.to_datetime_utc());
1267 let end_dt = cmd.end.map(|nanos| nanos.to_datetime_utc());
1268
1269 let mut reports = self
1270 .http_client
1271 .request_order_status_reports(
1272 self.core.account_id,
1273 cmd.instrument_id,
1274 start_dt,
1275 end_dt,
1276 cmd.open_only,
1277 None, )
1279 .await?;
1280
1281 normalize_spot_order_status_reports(&mut reports, self.config.treat_expired_as_canceled);
1282
1283 Ok(reports)
1284 }
1285
1286 async fn generate_fill_reports(
1287 &self,
1288 cmd: GenerateFillReports,
1289 ) -> anyhow::Result<Vec<FillReport>> {
1290 let Some(instrument_id) = cmd.instrument_id else {
1291 log::warn!("generate_fill_reports requires instrument_id for Binance Spot");
1292 return Ok(Vec::new());
1293 };
1294
1295 let venue_order_id = cmd
1297 .venue_order_id
1298 .as_ref()
1299 .map(|id| VenueOrderId::new(id.inner()));
1300 let requested_start_time = cmd
1301 .start
1302 .map(|start| start.as_i64() / NANOSECONDS_IN_MILLISECOND as i64);
1303 let requested_end_time = cmd
1304 .end
1305 .map(|end| end.as_i64() / NANOSECONDS_IN_MILLISECOND as i64);
1306 if let (Some(start), Some(end)) = (requested_start_time, requested_end_time) {
1307 anyhow::ensure!(
1308 start <= end,
1309 "fill report start time must not exceed end time"
1310 );
1311 }
1312
1313 let mut reports = Vec::new();
1314 let mut seen_trade_ids = AHashSet::new();
1315
1316 if venue_order_id.is_some() {
1317 let mut from_id = 0;
1318
1319 loop {
1320 let page = self
1321 .http_client
1322 .request_fill_reports_with_cursor(
1323 self.core.account_id,
1324 instrument_id,
1325 venue_order_id,
1326 None,
1327 None,
1328 Some(from_id),
1329 Some(ACCOUNT_TRADES_PAGE_LIMIT),
1330 )
1331 .await?;
1332
1333 if page.is_empty() {
1334 break;
1335 }
1336
1337 let page_len = page.len();
1338 let max_trade_id = max_trade_id(&page)?;
1339 let passed_end = requested_end_time.is_some_and(|end_time| {
1340 page.iter().any(|report| report_time_ms(report) > end_time)
1341 });
1342
1343 reports.extend(page.into_iter().filter(|report| {
1344 requested_start_time
1345 .is_none_or(|start_time| report_time_ms(report) >= start_time)
1346 && requested_end_time
1347 .is_none_or(|end_time| report_time_ms(report) <= end_time)
1348 && seen_trade_ids.insert(report.trade_id)
1349 }));
1350
1351 if page_len < ACCOUNT_TRADES_PAGE_LIMIT as usize || passed_end {
1352 break;
1353 }
1354
1355 let next_from_id = max_trade_id
1356 .checked_add(1)
1357 .context("Binance Spot trade ID overflow during pagination")?;
1358 anyhow::ensure!(
1359 next_from_id > from_id,
1360 "Binance Spot account-trades pagination made no progress"
1361 );
1362 from_id = next_from_id;
1363 }
1364 } else if let Some(query_start_time) = requested_start_time {
1365 let query_end_time = requested_end_time.unwrap_or_else(|| {
1366 self.clock.get_time_ns().as_i64() / NANOSECONDS_IN_MILLISECOND as i64
1367 });
1368 anyhow::ensure!(
1369 query_start_time <= query_end_time,
1370 "fill report start time must not exceed end time"
1371 );
1372 let mut window_start = query_start_time;
1373
1374 loop {
1375 let window_end = window_start
1376 .saturating_add(ACCOUNT_TRADES_MAX_INTERVAL_MS)
1377 .min(query_end_time);
1378 let mut from_id = None;
1379
1380 loop {
1381 let start = if from_id.is_none() {
1382 Some(
1383 Timestamp::from_millisecond(window_start)
1384 .context("invalid Binance Spot account-trades start time")?,
1385 )
1386 } else {
1387 None
1388 };
1389 let end = if from_id.is_none() {
1390 Some(
1391 Timestamp::from_millisecond(window_end)
1392 .context("invalid Binance Spot account-trades end time")?,
1393 )
1394 } else {
1395 None
1396 };
1397 let page = self
1398 .http_client
1399 .request_fill_reports_with_cursor(
1400 self.core.account_id,
1401 instrument_id,
1402 None,
1403 start,
1404 end,
1405 from_id,
1406 Some(ACCOUNT_TRADES_PAGE_LIMIT),
1407 )
1408 .await?;
1409
1410 if page.is_empty() {
1411 break;
1412 }
1413
1414 let page_len = page.len();
1415 let max_trade_id = max_trade_id(&page)?;
1416 let passed_window_end = page
1417 .iter()
1418 .any(|report| report_time_ms(report) > window_end);
1419
1420 reports.extend(page.into_iter().filter(|report| {
1421 let report_time = report_time_ms(report);
1422 report_time >= window_start
1423 && report_time <= window_end
1424 && seen_trade_ids.insert(report.trade_id)
1425 }));
1426
1427 if page_len < ACCOUNT_TRADES_PAGE_LIMIT as usize || passed_window_end {
1428 break;
1429 }
1430
1431 let next_from_id = max_trade_id
1432 .checked_add(1)
1433 .context("Binance Spot trade ID overflow during pagination")?;
1434 anyhow::ensure!(
1435 from_id.is_none_or(|cursor| next_from_id > cursor),
1436 "Binance Spot account-trades pagination made no progress"
1437 );
1438 from_id = Some(next_from_id);
1439 }
1440
1441 if window_end >= query_end_time {
1442 break;
1443 }
1444 window_start = window_end.saturating_add(1);
1445 }
1446 } else {
1447 let mut from_id = 0;
1448
1449 loop {
1450 let page = self
1451 .http_client
1452 .request_fill_reports_with_cursor(
1453 self.core.account_id,
1454 instrument_id,
1455 None,
1456 None,
1457 None,
1458 Some(from_id),
1459 Some(ACCOUNT_TRADES_PAGE_LIMIT),
1460 )
1461 .await?;
1462
1463 if page.is_empty() {
1464 break;
1465 }
1466
1467 let page_len = page.len();
1468 let max_trade_id = max_trade_id(&page)?;
1469 let passed_end = requested_end_time.is_some_and(|end_time| {
1470 page.iter().any(|report| report_time_ms(report) > end_time)
1471 });
1472
1473 reports.extend(page.into_iter().filter(|report| {
1474 requested_end_time.is_none_or(|end_time| report_time_ms(report) <= end_time)
1475 && seen_trade_ids.insert(report.trade_id)
1476 }));
1477
1478 if page_len < ACCOUNT_TRADES_PAGE_LIMIT as usize || passed_end {
1479 break;
1480 }
1481
1482 let next_from_id = max_trade_id
1483 .checked_add(1)
1484 .context("Binance Spot trade ID overflow during pagination")?;
1485 anyhow::ensure!(
1486 next_from_id > from_id,
1487 "Binance Spot account-trades pagination made no progress"
1488 );
1489 from_id = next_from_id;
1490 }
1491 }
1492
1493 let mut reports_with_trade_ids = reports
1494 .into_iter()
1495 .map(|report| parse_trade_id(&report).map(|trade_id| (report, trade_id)))
1496 .collect::<anyhow::Result<Vec<_>>>()?;
1497 reports_with_trade_ids
1498 .sort_unstable_by_key(|(report, trade_id)| (report.ts_event, *trade_id));
1499 Ok(reports_with_trade_ids
1500 .into_iter()
1501 .map(|(report, _)| report)
1502 .collect())
1503 }
1504
1505 async fn generate_position_status_reports(
1506 &self,
1507 _cmd: &GeneratePositionStatusReports,
1508 ) -> anyhow::Result<Vec<PositionStatusReport>> {
1509 Ok(Vec::new())
1512 }
1513
1514 async fn generate_mass_status(
1515 &self,
1516 lookback_mins: Option<u64>,
1517 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
1518 log::info!("Generating ExecutionMassStatus (lookback_mins={lookback_mins:?})");
1519
1520 let ts_now = self.clock.get_time_ns();
1521
1522 let start = if let Some(mins) = lookback_mins {
1523 let lookback_ns = checked_mins_to_nanos(mins)
1524 .context("lookback minutes exceed the nanosecond range")?;
1525 Some(UnixNanos::from(ts_now.as_u64().saturating_sub(lookback_ns)))
1526 } else {
1527 None
1528 };
1529
1530 let order_cmd = GenerateOrderStatusReportsBuilder::default()
1533 .ts_init(ts_now)
1534 .open_only(true)
1535 .start(start)
1536 .build()
1537 .map_err(|e| anyhow::anyhow!("{e}"))?;
1538
1539 let position_cmd = GeneratePositionStatusReportsBuilder::default()
1540 .ts_init(ts_now)
1541 .start(start)
1542 .build()
1543 .map_err(|e| anyhow::anyhow!("{e}"))?;
1544
1545 let (order_reports, position_reports) = tokio::try_join!(
1546 self.generate_order_status_reports(&order_cmd),
1547 self.generate_position_status_reports(&position_cmd),
1548 )?;
1549
1550 let mut instrument_ids: Vec<_> = order_reports
1551 .iter()
1552 .map(|report| report.instrument_id)
1553 .collect();
1554 {
1555 let cache = self.core.cache();
1556 instrument_ids.extend(
1557 cache
1558 .orders_open(
1559 Some(&BINANCE_VENUE),
1560 None,
1561 None,
1562 Some(&self.core.account_id),
1563 None,
1564 )
1565 .into_iter()
1566 .chain(cache.orders_inflight(
1567 Some(&BINANCE_VENUE),
1568 None,
1569 None,
1570 Some(&self.core.account_id),
1571 None,
1572 ))
1573 .map(|order| order.instrument_id())
1574 .filter(|instrument_id| {
1575 self.http_client
1576 .get_instrument(&instrument_id.symbol.inner())
1577 .is_some_and(|instrument| instrument.id() == *instrument_id)
1578 }),
1579 );
1580 }
1581 instrument_ids.sort_unstable();
1582 instrument_ids.dedup();
1583
1584 let mut fill_reports = Vec::new();
1585
1586 for instrument_id in instrument_ids {
1587 let fill_cmd = GenerateFillReportsBuilder::default()
1588 .ts_init(ts_now)
1589 .instrument_id(Some(instrument_id))
1590 .start(start)
1591 .build()
1592 .map_err(|e| anyhow::anyhow!("{e}"))?;
1593 fill_reports.extend(self.generate_fill_reports(fill_cmd).await?);
1594 }
1595
1596 log::info!("Received {} OrderStatusReports", order_reports.len());
1597 log::info!("Received {} FillReports", fill_reports.len());
1598 log::info!("Received {} PositionReports", position_reports.len());
1599
1600 let mut mass_status = ExecutionMassStatus::new(
1601 self.core.client_id,
1602 self.core.account_id,
1603 *BINANCE_VENUE,
1604 ts_now,
1605 None,
1606 );
1607
1608 mass_status.add_order_reports(order_reports);
1609 mass_status.add_fill_reports(fill_reports);
1610 mass_status.add_position_reports(position_reports);
1611
1612 Ok(Some(mass_status))
1613 }
1614
1615 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
1616 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
1617
1618 if order.is_closed() {
1619 let client_order_id = order.client_order_id();
1620 log::warn!("Cannot submit closed order {client_order_id}");
1621 return Ok(());
1622 }
1623
1624 if order.time_in_force() == TimeInForce::Gtd && self.config.use_gtd {
1625 time_in_force_to_binance_spot(order.time_in_force(), self.config.use_gtd)?;
1626 }
1627
1628 self.ensure_ws_user_data_active()?;
1629
1630 log::debug!("OrderSubmitted client_order_id={}", order.client_order_id());
1631 self.emitter.emit_order_submitted(&order);
1632
1633 self.submit_order_internal(&cmd)
1634 }
1635
1636 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
1637 if cmd.order_list.client_order_ids.is_empty() {
1638 log::debug!("submit_order_list called with empty order list");
1639 return Ok(());
1640 }
1641
1642 self.ensure_ws_user_data_active()?;
1643
1644 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
1645
1646 if let Some(order) = orders.iter().find(|order| order.is_closed()) {
1647 let reason = format!("Cannot submit closed order {}", order.client_order_id());
1648 for order in &orders {
1649 self.emitter.emit_order_denied(order, &reason);
1650 }
1651 return Ok(());
1652 }
1653
1654 let params = match build_spot_order_list_params(
1655 cmd.order_list.id.as_ref(),
1656 &orders,
1657 self.config.use_gtd,
1658 ) {
1659 Ok(request) => request,
1660 Err(reason) => {
1661 for order in &orders {
1662 self.emitter.emit_order_denied(order, &reason);
1663 }
1664 return Ok(());
1665 }
1666 };
1667
1668 for order in &orders {
1669 self.dispatch_state.order_identities.insert(
1670 order.client_order_id(),
1671 OrderIdentity {
1672 instrument_id: order.instrument_id(),
1673 strategy_id: order.strategy_id(),
1674 order_side: order.order_side(),
1675 order_type: order.order_type(),
1676 price: order.price(),
1677 quantity: order.quantity(),
1678 venue_position_id: None,
1679 },
1680 );
1681 self.emitter.emit_order_submitted(order);
1682 }
1683
1684 let event_emitter = self.emitter.clone();
1685 let trader_id = self.core.trader_id;
1686 let account_id = self.core.account_id;
1687 let clock = self.clock;
1688 let http_client = self.http_client.clone();
1689 let dispatch_state = self.dispatch_state.clone();
1690
1691 self.spawn_task("submit_order_list_http", async move {
1692 if let Err(e) = submit_spot_order_list(&http_client, ¶ms).await {
1693 handle_spot_order_list_submit_error(
1694 &event_emitter,
1695 &dispatch_state,
1696 trader_id,
1697 account_id,
1698 clock,
1699 &orders,
1700 e,
1701 )?;
1702 }
1703
1704 Ok(())
1705 });
1706
1707 Ok(())
1708 }
1709
1710 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1711 let order = self
1715 .core
1716 .cache()
1717 .order(&cmd.client_order_id)
1718 .map(|o| o.clone());
1719
1720 let Some(order) = order else {
1721 log::warn!(
1722 "Cannot modify order {}: not found in cache",
1723 cmd.client_order_id
1724 );
1725 let ts_init = self.clock.get_time_ns();
1726 let rejected_event = OrderModifyRejected::new(
1727 self.core.trader_id,
1728 cmd.strategy_id,
1729 cmd.instrument_id,
1730 cmd.client_order_id,
1731 "Order not found in cache for modify".into(),
1732 UUID4::new(),
1733 ts_init, ts_init,
1735 false,
1736 cmd.venue_order_id,
1737 Some(self.core.account_id),
1738 );
1739
1740 self.emitter
1741 .send_order_event(OrderEventAny::ModifyRejected(rejected_event));
1742 return Ok(());
1743 };
1744
1745 let event_emitter = self.emitter.clone();
1746 let trader_id = self.core.trader_id;
1747 let account_id = self.core.account_id;
1748 let clock = self.clock;
1749
1750 let order_side = order.order_side();
1751 let order_type = order.order_type();
1752 let time_in_force = order.time_in_force();
1753 let quantity = cmd.quantity.unwrap_or_else(|| order.quantity());
1754 let use_gtd = self.config.use_gtd;
1755
1756 if self.ws_order_transport_active() {
1757 let command = cmd;
1758 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
1759 let dispatch_state = self.dispatch_state.clone();
1760 let params = build_cancel_replace_params(&command, &order, quantity, use_gtd)?;
1761
1762 let request_id = ws_client.next_request_id();
1764 dispatch_state.pending_requests.insert(
1765 request_id.clone(),
1766 PendingRequest {
1767 client_order_id: command.client_order_id,
1768 venue_order_id: command.venue_order_id,
1769 operation: PendingOperation::Modify,
1770 },
1771 );
1772
1773 self.spawn_task("modify_order_ws", async move {
1774 if let Err(e) = ws_client
1775 .cancel_replace_order_with_id(request_id.clone(), params)
1776 .await
1777 {
1778 dispatch_state.pending_requests.remove(&request_id);
1779 log::warn!(
1780 "WS modify request failed for {}, awaiting reconciliation: {e}",
1781 command.client_order_id
1782 );
1783 anyhow::bail!("WS modify order failed: {e}");
1784 }
1785 Ok(())
1786 });
1787 } else {
1788 let command = cmd;
1789 let http_client = self.http_client.clone();
1790 log::debug!("WS trading not active, falling back to HTTP for modify_order");
1791
1792 self.spawn_task("modify_order_http", async move {
1793 let result = match command.venue_order_id {
1794 Some(venue_order_id) => {
1795 http_client
1796 .modify_order(
1797 account_id,
1798 command.instrument_id,
1799 venue_order_id,
1800 command.client_order_id,
1801 order_side,
1802 order_type,
1803 quantity,
1804 time_in_force,
1805 command.price,
1806 use_gtd,
1807 )
1808 .await
1809 }
1810 None => Err(anyhow::anyhow!(BinanceSpotHttpError::ValidationError(
1811 "venue_order_id required for modify".to_string()
1812 ))),
1813 };
1814
1815 match result {
1816 Ok(report) => {
1817 let ts_now = clock.get_time_ns();
1818 let updated_event = OrderUpdated::new(
1819 trader_id,
1820 command.strategy_id,
1821 command.instrument_id,
1822 command.client_order_id,
1823 report.quantity,
1824 UUID4::new(),
1825 ts_now,
1826 ts_now,
1827 false,
1828 Some(report.venue_order_id),
1829 Some(account_id),
1830 report.price,
1831 None, None, false, );
1835 event_emitter.send_order_event(OrderEventAny::Updated(updated_event));
1836 }
1837 Err(e) => {
1838 if is_structured_venue_rejection(&e) || is_local_command_failure(&e) {
1839 let ts_now = clock.get_time_ns();
1840 let rejected_event = OrderModifyRejected::new(
1841 trader_id,
1842 command.strategy_id,
1843 command.instrument_id,
1844 command.client_order_id,
1845 format!("modify-order-error: {e}").into(),
1846 UUID4::new(),
1847 ts_now,
1848 ts_now,
1849 false,
1850 command.venue_order_id,
1851 Some(account_id),
1852 );
1853 event_emitter
1854 .send_order_event(OrderEventAny::ModifyRejected(rejected_event));
1855 } else {
1856 log::warn!(
1857 "Ambiguous modify failure for {}, awaiting reconciliation: {e}",
1858 command.client_order_id
1859 );
1860 }
1861 return Err(e);
1862 }
1863 }
1864 Ok(())
1865 });
1866 }
1867
1868 Ok(())
1869 }
1870
1871 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
1872 self.cancel_order_internal(&cmd);
1873 Ok(())
1874 }
1875
1876 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
1877 let event_emitter = self.emitter.clone();
1878 let trader_id = self.core.trader_id;
1879 let account_id = self.core.account_id;
1880 let clock = self.clock;
1881
1882 if self.ws_order_transport_active() {
1883 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
1884 let symbol = cmd.instrument_id.symbol.to_string();
1885
1886 self.spawn_task("cancel_all_orders_ws", async move {
1887 if let Err(e) = ws_client.cancel_all_orders(symbol).await {
1888 log::error!("WS cancel_all_orders failed: {e}");
1889 }
1890 Ok(())
1892 });
1893
1894 return Ok(());
1895 }
1896
1897 log::debug!("WS trading not active, falling back to HTTP for cancel_all_orders");
1898 let http_client = self.http_client.clone();
1899
1900 let strategy_lookup: AHashMap<ClientOrderId, StrategyId> = {
1902 let cache = self.core.cache();
1903 cache
1904 .orders_open(None, Some(&cmd.instrument_id), None, None, None)
1905 .into_iter()
1906 .map(|order| (order.client_order_id(), order.strategy_id()))
1907 .collect()
1908 };
1909
1910 let command = cmd;
1911 self.spawn_task("cancel_all_orders_http", async move {
1912 let canceled_orders = http_client.cancel_all_orders(command.instrument_id).await?;
1913
1914 for (venue_order_id, client_order_id) in canceled_orders {
1915 let strategy_id = strategy_lookup
1916 .get(&client_order_id)
1917 .copied()
1918 .unwrap_or(command.strategy_id);
1919
1920 let canceled_event = OrderCanceled::new(
1921 trader_id,
1922 strategy_id,
1923 command.instrument_id,
1924 client_order_id,
1925 UUID4::new(),
1926 command.ts_init,
1927 clock.get_time_ns(),
1928 false,
1929 Some(venue_order_id),
1930 Some(account_id),
1931 );
1932
1933 event_emitter.send_order_event(OrderEventAny::Canceled(canceled_event));
1934 }
1935
1936 Ok(())
1937 });
1938
1939 Ok(())
1940 }
1941
1942 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1943 const BATCH_SIZE: usize = 5;
1944
1945 if cmd.cancels.is_empty() {
1946 return Ok(());
1947 }
1948
1949 let http_client = self.http_client.clone();
1950 let command = cmd;
1951
1952 let event_emitter = self.emitter.clone();
1953 let trader_id = self.core.trader_id;
1954 let account_id = self.core.account_id;
1955 let clock = self.clock;
1956
1957 self.spawn_task("batch_cancel_orders", async move {
1958 for chunk in command.cancels.chunks(BATCH_SIZE) {
1959 let batch_items: Vec<BatchCancelItem> = chunk
1960 .iter()
1961 .map(|cancel| {
1962 if let Some(venue_order_id) = cancel.venue_order_id {
1963 let order_id = venue_order_id.inner().parse::<i64>().unwrap_or(0);
1964 if order_id != 0 {
1965 BatchCancelItem::by_order_id(
1966 command.instrument_id.symbol.to_string(),
1967 order_id,
1968 )
1969 } else {
1970 BatchCancelItem::by_client_order_id(
1971 command.instrument_id.symbol.to_string(),
1972 encode_broker_id(
1973 &cancel.client_order_id,
1974 BINANCE_NAUTILUS_SPOT_BROKER_ID,
1975 ),
1976 )
1977 }
1978 } else {
1979 BatchCancelItem::by_client_order_id(
1980 command.instrument_id.symbol.to_string(),
1981 encode_broker_id(
1982 &cancel.client_order_id,
1983 BINANCE_NAUTILUS_SPOT_BROKER_ID,
1984 ),
1985 )
1986 }
1987 })
1988 .collect();
1989
1990 match http_client.batch_cancel_orders(&batch_items).await {
1991 Ok(results) => {
1992 for (i, result) in results.iter().enumerate() {
1993 let cancel = &chunk[i];
1994
1995 match result {
1996 BatchCancelResult::Success(success) => {
1997 let venue_order_id =
1998 VenueOrderId::new(success.order_id.to_string());
1999 let canceled_event = OrderCanceled::new(
2000 trader_id,
2001 cancel.strategy_id,
2002 cancel.instrument_id,
2003 cancel.client_order_id,
2004 UUID4::new(),
2005 cancel.ts_init,
2006 clock.get_time_ns(),
2007 false,
2008 Some(venue_order_id),
2009 Some(account_id),
2010 );
2011
2012 event_emitter
2013 .send_order_event(OrderEventAny::Canceled(canceled_event));
2014 }
2015 BatchCancelResult::Error(error) => {
2016 let rejected_event = OrderCancelRejected::new(
2017 trader_id,
2018 cancel.strategy_id,
2019 cancel.instrument_id,
2020 cancel.client_order_id,
2021 format!(
2022 "batch-cancel-error: code={}, msg={}",
2023 error.code, error.msg
2024 )
2025 .into(),
2026 UUID4::new(),
2027 clock.get_time_ns(),
2028 cancel.ts_init,
2029 false,
2030 cancel.venue_order_id,
2031 Some(account_id),
2032 );
2033
2034 event_emitter.send_order_event(OrderEventAny::CancelRejected(
2035 rejected_event,
2036 ));
2037 }
2038 }
2039 }
2040 }
2041 Err(e) => {
2042 if is_local_http_command_failure(&e) {
2043 log::warn!(
2044 "Batch cancel command failed local validation for {} orders: {e}",
2045 chunk.len()
2046 );
2047 } else {
2048 log::warn!(
2049 "Ambiguous batch cancel failure for {} orders, awaiting reconciliation: {e}",
2050 chunk.len()
2051 );
2052 }
2053 }
2054 }
2055 }
2056
2057 Ok(())
2058 });
2059
2060 Ok(())
2061 }
2062}
2063
2064struct DispatchActiveGuard(Arc<AtomicBool>);
2065
2066impl Drop for DispatchActiveGuard {
2067 fn drop(&mut self) {
2068 self.0.store(false, Ordering::Release);
2069 }
2070}
2071
2072fn max_trade_id(reports: &[FillReport]) -> anyhow::Result<i64> {
2073 let mut max_trade_id = None;
2074
2075 for report in reports {
2076 let trade_id = parse_trade_id(report)?;
2077 max_trade_id = Some(max_trade_id.map_or(trade_id, |current: i64| current.max(trade_id)));
2078 }
2079
2080 max_trade_id.context("Binance Spot account-trades page was empty")
2081}
2082
2083fn parse_trade_id(report: &FillReport) -> anyhow::Result<i64> {
2084 report
2085 .trade_id
2086 .to_string()
2087 .parse::<i64>()
2088 .with_context(|| format!("invalid Binance Spot trade ID {}", report.trade_id))
2089}
2090
2091fn report_time_ms(report: &FillReport) -> i64 {
2092 report.ts_event.as_i64() / NANOSECONDS_IN_MILLISECOND as i64
2093}
2094
2095fn normalize_spot_order_status_report(
2096 report: &mut OrderStatusReport,
2097 treat_expired_as_canceled: bool,
2098) {
2099 if treat_expired_as_canceled && report.order_status == OrderStatus::Expired {
2100 report.order_status = OrderStatus::Canceled;
2101 }
2102}
2103
2104fn normalize_spot_order_status_reports(
2105 reports: &mut [OrderStatusReport],
2106 treat_expired_as_canceled: bool,
2107) {
2108 for report in reports {
2109 normalize_spot_order_status_report(report, treat_expired_as_canceled);
2110 }
2111}
2112
2113async fn wait_for_ws_setup_response(
2114 timeout: Duration,
2115 success: impl Future<Output = ()>,
2116 setup_errors: &mut tokio::sync::mpsc::UnboundedReceiver<String>,
2117 timeout_message: &'static str,
2118) -> anyhow::Result<()> {
2119 tokio::pin!(success);
2120
2121 let result = tokio::time::timeout(timeout, async {
2122 tokio::select! {
2123 () = &mut success => Ok(()),
2124 err = setup_errors.recv() => {
2125 anyhow::bail!(
2126 "{}",
2127 err.unwrap_or_else(|| "WS setup error channel closed".to_string()),
2128 )
2129 }
2130 }
2131 })
2132 .await;
2133
2134 result.map_err(|_| anyhow::anyhow!(timeout_message))?
2135}
2136
2137#[expect(clippy::too_many_arguments)]
2138fn dispatch_ws_trading_message(
2139 msg: BinanceSpotWsTradingMessage,
2140 emitter: &ExecutionEventEmitter,
2141 http_client: &BinanceSpotHttpClient,
2142 account_id: AccountId,
2143 treat_expired_as_canceled: bool,
2144 clock: &'static AtomicTime,
2145 dispatch_state: &WsDispatchState,
2146 ws_authenticated: &tokio::sync::Notify,
2147 ws_user_data_subscribed: &tokio::sync::Notify,
2148 ws_setup_error_tx: &tokio::sync::mpsc::UnboundedSender<String>,
2149 seen_trade_ids: &std::sync::Arc<Mutex<FifoCache<(Ustr, i64), 10_000>>>,
2150 task_spawner: &TaskSpawner,
2151) {
2152 match msg {
2153 BinanceSpotWsTradingMessage::OrderAccepted {
2154 request_id,
2155 response,
2156 } => {
2157 dispatch_state.pending_requests.remove(&request_id);
2158 log::debug!(
2159 "WS order accepted: request_id={request_id}, order_id={}",
2160 response.order_id
2161 );
2162 }
2164 BinanceSpotWsTradingMessage::OrderRejected {
2165 request_id,
2166 code,
2167 msg,
2168 } => {
2169 log::debug!("WS order rejected: request_id={request_id}, code={code}, msg={msg}");
2170 if let Some((_, pending)) = dispatch_state.pending_requests.remove(&request_id) {
2171 let code_i64 = i64::from(code);
2172 if matches!(
2173 code_i64,
2174 BINANCE_UNEXPECTED_RESPONSE_CODE | BINANCE_STATUS_UNKNOWN_CODE
2175 ) {
2176 log::warn!(
2177 "Ambiguous WS submit failure for {}, awaiting reconciliation: code={code}, msg={msg}",
2178 pending.client_order_id,
2179 );
2180 return;
2181 }
2182
2183 let identity = dispatch_state
2185 .order_identities
2186 .get(&pending.client_order_id)
2187 .map(|r| r.clone());
2188
2189 if let Some(identity) = identity {
2190 let due_post_only = code_i64 == BINANCE_GTX_ORDER_REJECT_CODE
2191 || (code_i64 == BINANCE_NEW_ORDER_REJECTED_CODE
2192 && msg == BINANCE_SPOT_POST_ONLY_REJECT_MSG);
2193 let ts_now = clock.get_time_ns();
2194 let rejected = OrderRejected::new(
2195 emitter.trader_id(),
2196 identity.strategy_id,
2197 identity.instrument_id,
2198 pending.client_order_id,
2199 account_id,
2200 Ustr::from(&format!("code={code}: {msg}")),
2201 UUID4::new(),
2202 ts_now,
2203 ts_now,
2204 false,
2205 due_post_only,
2206 );
2207 dispatch_state.cleanup_terminal(pending.client_order_id);
2208 emitter.send_order_event(OrderEventAny::Rejected(rejected));
2209 } else {
2210 log::warn!(
2211 "No order identity for {}, cannot emit OrderRejected",
2212 pending.client_order_id
2213 );
2214 }
2215 } else {
2216 log::warn!("No pending request for {request_id}, cannot emit OrderRejected");
2217 }
2218 }
2219 BinanceSpotWsTradingMessage::OrderCanceled {
2220 request_id,
2221 response,
2222 } => {
2223 dispatch_state.pending_requests.remove(&request_id);
2224 log::debug!(
2225 "WS order canceled: request_id={request_id}, order_id={}",
2226 response.order_id
2227 );
2228 }
2230 BinanceSpotWsTradingMessage::CancelRejected {
2231 request_id,
2232 code,
2233 msg,
2234 } => {
2235 log::warn!("WS cancel rejected: request_id={request_id}, code={code}, msg={msg}");
2236 if let Some((_, pending)) = dispatch_state.pending_requests.remove(&request_id)
2237 && let Some(identity) = dispatch_state
2238 .order_identities
2239 .get(&pending.client_order_id)
2240 {
2241 let ts_now = clock.get_time_ns();
2242 let rejected = OrderCancelRejected::new(
2243 emitter.trader_id(),
2244 identity.strategy_id,
2245 identity.instrument_id,
2246 pending.client_order_id,
2247 Ustr::from(&format!("code={code}: {msg}")),
2248 UUID4::new(),
2249 ts_now,
2250 ts_now,
2251 false,
2252 pending.venue_order_id,
2253 Some(account_id),
2254 );
2255 emitter.send_order_event(OrderEventAny::CancelRejected(rejected));
2256 }
2257 }
2258 BinanceSpotWsTradingMessage::CancelReplaceAccepted {
2259 request_id,
2260 cancel_response,
2261 new_order_response,
2262 } => {
2263 dispatch_state.pending_requests.remove(&request_id);
2264 log::debug!(
2265 "WS cancel-replace accepted: request_id={request_id}, \
2266 canceled_id={}, new_id={}",
2267 cancel_response.order_id,
2268 new_order_response.order_id,
2269 );
2270 }
2272 BinanceSpotWsTradingMessage::CancelReplaceRejected {
2273 request_id,
2274 code,
2275 msg,
2276 } => {
2277 log::warn!(
2278 "WS cancel-replace rejected: request_id={request_id}, code={code}, msg={msg}"
2279 );
2280
2281 if let Some((_, pending)) = dispatch_state.pending_requests.remove(&request_id)
2282 && let Some(identity) = dispatch_state
2283 .order_identities
2284 .get(&pending.client_order_id)
2285 {
2286 let ts_now = clock.get_time_ns();
2287 let rejected = OrderModifyRejected::new(
2288 emitter.trader_id(),
2289 identity.strategy_id,
2290 identity.instrument_id,
2291 pending.client_order_id,
2292 Ustr::from(&format!("code={code}: {msg}")),
2293 UUID4::new(),
2294 ts_now,
2295 ts_now,
2296 false,
2297 pending.venue_order_id,
2298 Some(account_id),
2299 );
2300 emitter.send_order_event(OrderEventAny::ModifyRejected(rejected));
2301 }
2302 }
2303 BinanceSpotWsTradingMessage::RequestFailed { request_id, msg } => {
2304 dispatch_state.pending_requests.remove(&request_id);
2305 log::error!(
2306 "WS trading request failed without structured venue response: request_id={request_id}, {msg}"
2307 );
2308 }
2309 BinanceSpotWsTradingMessage::AllOrdersCanceled {
2310 request_id,
2311 responses,
2312 } => {
2313 dispatch_state.pending_requests.remove(&request_id);
2314 log::debug!(
2315 "WS all orders canceled: request_id={request_id}, count={}",
2316 responses.len()
2317 );
2318 }
2320 BinanceSpotWsTradingMessage::UserDataSubscribed { subscription_id } => {
2321 log::debug!("User data stream subscribed: id={subscription_id}");
2322 ws_user_data_subscribed.notify_one();
2323 }
2324 BinanceSpotWsTradingMessage::ExecutionReport(report) => {
2325 let ts_init = clock.get_time_ns();
2326 dispatch_execution_report(
2327 &report,
2328 emitter,
2329 http_client,
2330 account_id,
2331 treat_expired_as_canceled,
2332 dispatch_state,
2333 seen_trade_ids,
2334 ts_init,
2335 );
2336 }
2337 BinanceSpotWsTradingMessage::AccountPosition(position) => {
2338 let ts_init = clock.get_time_ns();
2339 let state = parse_spot_account_position(&position, account_id, ts_init);
2340 emitter.send_account_state(state);
2341 }
2342 BinanceSpotWsTradingMessage::BalanceUpdate(update) => {
2343 log::debug!(
2344 "Balance update: asset={}, delta={}",
2345 update.asset,
2346 update.delta,
2347 );
2348 let http_client = http_client.clone();
2349 let emitter = emitter.clone();
2350
2351 if let Err(e) = task_spawner.spawn(async move {
2352 match http_client.request_account_state(account_id).await {
2353 Ok(state) => emitter.send_account_state(state),
2354 Err(e) => {
2355 log::error!("Failed to refresh account state after balance update: {e}");
2356 }
2357 }
2358 }) {
2359 log::warn!("Skipping Binance Spot balance refresh after shutdown began: {e}");
2360 }
2361 }
2362 BinanceSpotWsTradingMessage::Connected => {
2363 log::debug!("WS trading API connected");
2364 }
2365 BinanceSpotWsTradingMessage::Authenticated => {
2366 log::debug!("WS trading API authenticated");
2367 ws_authenticated.notify_one();
2368 }
2369 BinanceSpotWsTradingMessage::AuthenticationRejected(reason) => {
2370 log::error!("WS trading API authentication failed: {reason}");
2371 let _ = ws_setup_error_tx.send(reason);
2372 }
2373 BinanceSpotWsTradingMessage::Reconnected => {
2374 log::info!("WS trading API reconnected");
2375 }
2376 BinanceSpotWsTradingMessage::ServerShutdown { event_time } => {
2377 log::warn!(
2378 "WS trading API server shutdown notice (event_time={event_time}); reconnect expected within ~10 minutes"
2379 );
2380 }
2381 BinanceSpotWsTradingMessage::Error(err) => {
2382 log::error!("WS trading API error: {err}");
2383 let _ = ws_setup_error_tx.send(err);
2384 }
2385 BinanceSpotWsTradingMessage::UserDataSubscriptionRejected(reason) => {
2386 log::error!("WS trading API user data subscription failed: {reason}");
2387 let _ = ws_setup_error_tx.send(reason);
2388 }
2389 }
2390}
2391
2392fn build_new_order_params(
2393 order: &impl Order,
2394 client_order_id: ClientOrderId,
2395 is_post_only: bool,
2396 is_quote_quantity: bool,
2397 use_gtd: bool,
2398) -> anyhow::Result<NewOrderParams> {
2399 let binance_side = BinanceSide::try_from(order.order_side())?;
2400 let binance_order_type = order_type_to_binance_spot(order.order_type(), is_post_only)?;
2401
2402 let requires_trigger = matches!(
2403 order.order_type(),
2404 OrderType::StopMarket
2405 | OrderType::StopLimit
2406 | OrderType::MarketIfTouched
2407 | OrderType::LimitIfTouched
2408 );
2409
2410 if requires_trigger && order.trigger_price().is_none() {
2411 anyhow::bail!("Conditional orders require a trigger price");
2412 }
2413
2414 let supports_tif = matches!(
2415 binance_order_type,
2416 BinanceSpotOrderType::Limit
2417 | BinanceSpotOrderType::StopLossLimit
2418 | BinanceSpotOrderType::TakeProfitLimit
2419 );
2420 let binance_tif = time_in_force_to_binance_spot(order.time_in_force(), use_gtd)?;
2421 let binance_tif = if supports_tif {
2422 Some(binance_tif)
2423 } else {
2424 None
2425 };
2426
2427 let qty_str = order.quantity().to_string();
2428 let (base_qty, quote_qty) = if is_quote_quantity {
2429 (None, Some(qty_str))
2430 } else {
2431 (Some(qty_str), None)
2432 };
2433
2434 let client_id_str = encode_broker_id(&client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
2435
2436 Ok(NewOrderParams {
2437 symbol: order.instrument_id().symbol.to_string(),
2438 side: binance_side,
2439 order_type: binance_order_type,
2440 time_in_force: binance_tif,
2441 quantity: base_qty,
2442 quote_order_qty: quote_qty,
2443 price: order.price().map(|p| p.to_string()),
2444 new_client_order_id: Some(client_id_str),
2445 stop_price: order.trigger_price().map(|p| p.to_string()),
2446 trailing_delta: None,
2447 iceberg_qty: order.display_qty().map(|q| q.to_string()),
2448 new_order_resp_type: Some(BinanceOrderResponseType::Full),
2449 self_trade_prevention_mode: None,
2450 strategy_id: None,
2451 strategy_type: None,
2452 })
2453}
2454
2455fn build_spot_order_list_params(
2456 order_list_id: &str,
2457 orders: &[OrderAny],
2458 use_gtd: bool,
2459) -> Result<NewOcoOrderListParams, String> {
2460 let has_grouped_order = orders.iter().any(is_grouped_order);
2461
2462 if has_grouped_order {
2463 return build_spot_oco_order_list_params(order_list_id, orders, use_gtd);
2464 }
2465
2466 Err("Binance Spot order-list submission currently supports only OCO lists".to_string())
2467}
2468
2469fn build_spot_oco_order_list_params(
2470 order_list_id: &str,
2471 orders: &[OrderAny],
2472 use_gtd: bool,
2473) -> Result<NewOcoOrderListParams, String> {
2474 if orders.len() != 2 {
2475 return Err(format!(
2476 "Binance Spot OCO order-list submission requires exactly 2 orders, was {}",
2477 orders.len()
2478 ));
2479 }
2480
2481 if orders
2482 .iter()
2483 .any(|order| order.contingency_type() != Some(ContingencyType::Oco))
2484 {
2485 return Err(
2486 "Binance Spot grouped order-list submission currently supports only OCO lists"
2487 .to_string(),
2488 );
2489 }
2490
2491 let first = &orders[0];
2492 let second = &orders[1];
2493 if first.instrument_id() != second.instrument_id() {
2494 return Err("Binance Spot OCO order-list legs must use the same instrument".to_string());
2495 }
2496
2497 if first.order_side() != second.order_side() {
2498 return Err("Binance Spot OCO order-list legs must use the same side".to_string());
2499 }
2500
2501 if first.quantity() != second.quantity() {
2502 return Err("Binance Spot OCO order-list legs must use the same quantity".to_string());
2503 }
2504
2505 if first.is_quote_quantity() || second.is_quote_quantity() {
2506 return Err("Binance Spot OCO order-list legs do not support quote quantity".to_string());
2507 }
2508
2509 let mut above = None;
2510 let mut below = None;
2511
2512 for order in orders {
2513 let params = build_new_order_params(
2514 order,
2515 order.client_order_id(),
2516 order.is_post_only(),
2517 false,
2518 use_gtd,
2519 )
2520 .map_err(|e| e.to_string())?;
2521
2522 match spot_oco_leg_position(params.side, params.order_type)? {
2523 SpotOcoLegPosition::Above => {
2524 if above.replace(params).is_some() {
2525 return Err(
2526 "Binance Spot OCO order-list resolved more than one above leg".to_string(),
2527 );
2528 }
2529 }
2530 SpotOcoLegPosition::Below => {
2531 if below.replace(params).is_some() {
2532 return Err(
2533 "Binance Spot OCO order-list resolved more than one below leg".to_string(),
2534 );
2535 }
2536 }
2537 }
2538 }
2539
2540 let above = above.ok_or_else(|| "Binance Spot OCO order-list missing above leg".to_string())?;
2541 let below = below.ok_or_else(|| "Binance Spot OCO order-list missing below leg".to_string())?;
2542 let quantity = above
2543 .quantity
2544 .clone()
2545 .ok_or_else(|| "Binance Spot OCO order-list requires base quantity".to_string())?;
2546
2547 Ok(NewOcoOrderListParams {
2548 symbol: first.instrument_id().symbol.to_string(),
2549 list_client_order_id: Some(order_list_id.to_string()),
2550 side: above.side,
2551 quantity,
2552 above_type: above.order_type,
2553 above_client_order_id: above.new_client_order_id,
2554 above_iceberg_qty: above.iceberg_qty,
2555 above_price: above.price,
2556 above_stop_price: above.stop_price,
2557 above_time_in_force: above.time_in_force,
2558 below_type: below.order_type,
2559 below_client_order_id: below.new_client_order_id,
2560 below_iceberg_qty: below.iceberg_qty,
2561 below_price: below.price,
2562 below_stop_price: below.stop_price,
2563 below_time_in_force: below.time_in_force,
2564 new_order_resp_type: Some(BinanceOrderResponseType::Full),
2565 self_trade_prevention_mode: None,
2566 })
2567}
2568
2569enum SpotOcoLegPosition {
2570 Above,
2571 Below,
2572}
2573
2574fn spot_oco_leg_position(
2575 side: BinanceSide,
2576 order_type: BinanceSpotOrderType,
2577) -> Result<SpotOcoLegPosition, String> {
2578 match (side, order_type) {
2579 (
2580 BinanceSide::Sell,
2581 BinanceSpotOrderType::LimitMaker
2582 | BinanceSpotOrderType::TakeProfit
2583 | BinanceSpotOrderType::TakeProfitLimit,
2584 )
2585 | (
2586 BinanceSide::Buy,
2587 BinanceSpotOrderType::StopLoss | BinanceSpotOrderType::StopLossLimit,
2588 ) => Ok(SpotOcoLegPosition::Above),
2589 (
2590 BinanceSide::Sell,
2591 BinanceSpotOrderType::StopLoss | BinanceSpotOrderType::StopLossLimit,
2592 )
2593 | (
2594 BinanceSide::Buy,
2595 BinanceSpotOrderType::LimitMaker
2596 | BinanceSpotOrderType::TakeProfit
2597 | BinanceSpotOrderType::TakeProfitLimit,
2598 ) => Ok(SpotOcoLegPosition::Below),
2599 (_, unsupported) => Err(format!(
2600 "Unsupported Binance Spot OCO leg order type: {unsupported:?}"
2601 )),
2602 }
2603}
2604
2605fn is_grouped_order(order: &OrderAny) -> bool {
2606 order.contingency_type().is_some()
2607 || order
2608 .linked_order_ids()
2609 .is_some_and(|linked_order_ids| !linked_order_ids.is_empty())
2610}
2611
2612fn handle_spot_order_submit_success(client_order_id: ClientOrderId, venue_order_id: VenueOrderId) {
2613 log::debug!(
2614 "Order submit succeeded: client_order_id={client_order_id}, venue_order_id={venue_order_id}",
2615 );
2616}
2617
2618async fn submit_spot_order_list(
2619 http_client: &BinanceSpotHttpClient,
2620 params: &NewOcoOrderListParams,
2621) -> Result<(), BinanceSpotHttpError> {
2622 let response = http_client.submit_oco_order_list(params).await?;
2623 log::debug!(
2624 "Order list submit succeeded: order_list_id={}, order_count={}",
2625 response.order_list_id,
2626 response.orders.len(),
2627 );
2628 Ok(())
2629}
2630
2631fn handle_spot_order_list_submit_error(
2632 event_emitter: &ExecutionEventEmitter,
2633 dispatch_state: &WsDispatchState,
2634 trader_id: TraderId,
2635 account_id: AccountId,
2636 clock: &'static AtomicTime,
2637 orders: &[OrderAny],
2638 error: BinanceSpotHttpError,
2639) -> anyhow::Result<()> {
2640 let ambiguous = matches!(
2641 error,
2642 BinanceSpotHttpError::BinanceError {
2643 code: BINANCE_UNEXPECTED_RESPONSE_CODE | BINANCE_STATUS_UNKNOWN_CODE,
2644 ..
2645 }
2646 );
2647
2648 if ambiguous {
2649 log::error!("Ambiguous order-list submit failure, awaiting reconciliation: {error}");
2650 return Err(error.into());
2651 }
2652
2653 let reject_orders = matches!(
2654 error,
2655 BinanceSpotHttpError::BinanceError { .. }
2656 | BinanceSpotHttpError::MissingCredentials
2657 | BinanceSpotHttpError::ValidationError(_)
2658 );
2659
2660 if reject_orders {
2661 let ts_now = clock.get_time_ns();
2662 let reason = format!("submit-order-list-error: {error}");
2663 for order in orders {
2664 let client_order_id = order.client_order_id();
2665 dispatch_state.cleanup_terminal(client_order_id);
2666 let rejected = OrderRejected::new(
2667 trader_id,
2668 order.strategy_id(),
2669 order.instrument_id(),
2670 client_order_id,
2671 account_id,
2672 reason.clone().into(),
2673 UUID4::new(),
2674 ts_now,
2675 ts_now,
2676 false,
2677 false,
2678 );
2679 event_emitter.send_order_event(OrderEventAny::Rejected(rejected));
2680 }
2681 } else {
2682 log::error!("Order-list submit failed, awaiting reconciliation: {error}");
2683 }
2684
2685 Err(error.into())
2686}
2687
2688fn build_cancel_order_params(cmd: &CancelOrder) -> CancelOrderParams {
2689 let order_id = cmd
2690 .venue_order_id
2691 .and_then(|id| id.inner().parse::<i64>().ok());
2692
2693 if let Some(order_id) = order_id {
2694 CancelOrderParams::by_order_id(cmd.instrument_id.symbol.to_string(), order_id)
2695 } else {
2696 let client_id_str = encode_broker_id(&cmd.client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
2697 CancelOrderParams::by_client_order_id(cmd.instrument_id.symbol.to_string(), client_id_str)
2698 }
2699}
2700
2701fn build_cancel_replace_params(
2702 cmd: &ModifyOrder,
2703 order: &impl Order,
2704 quantity: Quantity,
2705 use_gtd: bool,
2706) -> anyhow::Result<CancelReplaceOrderParams> {
2707 let binance_side = BinanceSide::try_from(order.order_side())?;
2708 let binance_order_type = order_type_to_binance_spot(order.order_type(), false)?;
2709 let binance_tif = time_in_force_to_binance_spot(order.time_in_force(), use_gtd)?;
2710
2711 let cancel_order_id: Option<i64> = cmd
2712 .venue_order_id
2713 .map(|id| {
2714 id.inner()
2715 .parse::<i64>()
2716 .map_err(|_| anyhow::anyhow!("Invalid venue order ID: {id}"))
2717 })
2718 .transpose()?;
2719
2720 let client_id_str = encode_broker_id(&cmd.client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
2721
2722 Ok(CancelReplaceOrderParams {
2723 symbol: cmd.instrument_id.symbol.to_string(),
2724 side: binance_side,
2725 order_type: binance_order_type,
2726 cancel_replace_mode: BinanceCancelReplaceMode::StopOnFailure,
2727 time_in_force: Some(binance_tif),
2728 quantity: Some(quantity.to_string()),
2729 quote_order_qty: None,
2730 price: cmd.price.map(|p| p.to_string()),
2731 cancel_order_id,
2732 cancel_orig_client_order_id: if cancel_order_id.is_none() {
2733 Some(client_id_str.clone())
2734 } else {
2735 None
2736 },
2737 new_client_order_id: Some(client_id_str),
2738 stop_price: None,
2739 trailing_delta: None,
2740 iceberg_qty: None,
2741 new_order_resp_type: Some(BinanceOrderResponseType::Full),
2742 self_trade_prevention_mode: None,
2743 })
2744}
2745
2746#[expect(clippy::too_many_arguments)]
2751fn dispatch_execution_report(
2752 report: &BinanceSpotExecutionReport,
2753 emitter: &ExecutionEventEmitter,
2754 http_client: &BinanceSpotHttpClient,
2755 account_id: AccountId,
2756 treat_expired_as_canceled: bool,
2757 dispatch_state: &WsDispatchState,
2758 seen_trade_ids: &std::sync::Arc<Mutex<FifoCache<(Ustr, i64), 10_000>>>,
2759 ts_init: UnixNanos,
2760) {
2761 let symbol = report.symbol;
2762 let instrument_id = InstrumentId::new(symbol.into(), *BINANCE_VENUE);
2763 let (price_precision, size_precision) = http_client
2764 .get_instrument(&symbol)
2765 .map_or((8, 8), |i| (i.price_precision(), i.size_precision()));
2766
2767 let client_order_id =
2768 match decode_client_order_id(&report.client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID) {
2769 Ok(client_order_id) => client_order_id,
2770 Err(e) => {
2771 log::warn!("Skipping Spot execution report with invalid client order ID: {e}");
2772 return;
2773 }
2774 };
2775
2776 let identity = dispatch_state
2777 .order_identities
2778 .get(&client_order_id)
2779 .map(|r| r.clone());
2780
2781 if let Some(identity) = identity {
2782 dispatch_tracked_execution_report(
2783 report,
2784 emitter,
2785 account_id,
2786 treat_expired_as_canceled,
2787 dispatch_state,
2788 seen_trade_ids,
2789 client_order_id,
2790 &identity,
2791 instrument_id,
2792 price_precision,
2793 size_precision,
2794 ts_init,
2795 );
2796 } else {
2797 dispatch_untracked_execution_report(
2798 report,
2799 emitter,
2800 http_client,
2801 account_id,
2802 treat_expired_as_canceled,
2803 seen_trade_ids,
2804 instrument_id,
2805 price_precision,
2806 size_precision,
2807 ts_init,
2808 );
2809 }
2810}
2811
2812#[expect(clippy::too_many_arguments)]
2814fn dispatch_tracked_execution_report(
2815 report: &BinanceSpotExecutionReport,
2816 emitter: &ExecutionEventEmitter,
2817 account_id: AccountId,
2818 treat_expired_as_canceled: bool,
2819 state: &WsDispatchState,
2820 seen_trade_ids: &std::sync::Arc<Mutex<FifoCache<(Ustr, i64), 10_000>>>,
2821 client_order_id: ClientOrderId,
2822 identity: &OrderIdentity,
2823 instrument_id: InstrumentId,
2824 price_precision: u8,
2825 size_precision: u8,
2826 ts_init: UnixNanos,
2827) {
2828 let venue_order_id = VenueOrderId::new(report.order_id.to_string());
2829 let ts_event = parse_millis_or_init(report.event_time, "Spot execution event time", ts_init);
2830
2831 match report.execution_type {
2832 BinanceSpotExecutionType::New => {
2833 if state.has_filled(&client_order_id) {
2834 log::debug!("Skipping New for already-filled {client_order_id}");
2835 return;
2836 }
2837
2838 if state.has_emitted_accepted(&client_order_id) {
2839 let Some(price) = parse_spot_execution_report_price(
2841 report,
2842 &report.price,
2843 price_precision,
2844 "price",
2845 ) else {
2846 return;
2847 };
2848 let Some(quantity) = parse_spot_execution_report_quantity(
2849 report,
2850 &report.original_qty,
2851 size_precision,
2852 "original_qty",
2853 ) else {
2854 return;
2855 };
2856 let Some(stop_price) =
2857 parse_spot_execution_report_decimal(report, &report.stop_price, "stop_price")
2858 else {
2859 return;
2860 };
2861 let trigger = if stop_price > Decimal::ZERO {
2862 let Some(trigger_price) = parse_spot_execution_report_price(
2863 report,
2864 &report.stop_price,
2865 price_precision,
2866 "stop_price",
2867 ) else {
2868 return;
2869 };
2870 Some(trigger_price)
2871 } else {
2872 None
2873 };
2874 let updated = OrderUpdated::new(
2875 emitter.trader_id(),
2876 identity.strategy_id,
2877 identity.instrument_id,
2878 client_order_id,
2879 quantity,
2880 UUID4::new(),
2881 ts_event,
2882 ts_init,
2883 false,
2884 Some(venue_order_id),
2885 Some(account_id),
2886 Some(price),
2887 trigger,
2888 None, false, );
2891 emitter.send_order_event(OrderEventAny::Updated(updated));
2892 return;
2893 }
2894 state.insert_accepted(client_order_id);
2895 let accepted = OrderAccepted::new(
2896 emitter.trader_id(),
2897 identity.strategy_id,
2898 identity.instrument_id,
2899 client_order_id,
2900 venue_order_id,
2901 account_id,
2902 UUID4::new(),
2903 ts_event,
2904 ts_init,
2905 false,
2906 );
2907 emitter.send_order_event(OrderEventAny::Accepted(accepted));
2908 }
2909 BinanceSpotExecutionType::Trade => {
2910 let dedup_key = (report.symbol, report.trade_id);
2911 let mut guard = seen_trade_ids.lock();
2912 let is_duplicate = guard.contains(&dedup_key);
2913 guard.add(dedup_key);
2914 drop(guard);
2915
2916 if is_duplicate {
2917 log::debug!(
2918 "Duplicate trade_id={} for {}, skipping",
2919 report.trade_id,
2920 report.symbol
2921 );
2922 return;
2923 }
2924
2925 ensure_accepted_emitted(
2926 client_order_id,
2927 account_id,
2928 venue_order_id,
2929 identity,
2930 emitter,
2931 state,
2932 ts_init,
2933 );
2934
2935 let Some(last_qty) = parse_spot_execution_report_quantity(
2936 report,
2937 &report.last_filled_qty,
2938 size_precision,
2939 "last_filled_qty",
2940 ) else {
2941 return;
2942 };
2943 let Some(last_px) = parse_spot_execution_report_price(
2944 report,
2945 &report.last_filled_price,
2946 price_precision,
2947 "last_filled_price",
2948 ) else {
2949 return;
2950 };
2951 let Some(commission) =
2952 parse_spot_execution_report_decimal(report, &report.commission, "commission")
2953 else {
2954 return;
2955 };
2956 let commission_currency = report
2957 .commission_asset
2958 .as_ref()
2959 .map_or_else(Currency::USDT, |a| {
2960 Currency::get_or_create_crypto(a.as_str())
2961 });
2962 let commission_money = match Money::from_decimal(commission, commission_currency) {
2963 Ok(money) => money,
2964 Err(e) => {
2965 log::warn!(
2966 "Failed to build Spot commission money for symbol={}, order_id={}, \
2967 trade_id={}: {e}",
2968 report.symbol,
2969 report.order_id,
2970 report.trade_id,
2971 );
2972 return;
2973 }
2974 };
2975
2976 let liquidity_side = if report.is_maker {
2977 LiquiditySide::Maker
2978 } else {
2979 LiquiditySide::Taker
2980 };
2981
2982 let filled = OrderFilled::new(
2983 emitter.trader_id(),
2984 identity.strategy_id,
2985 instrument_id,
2986 client_order_id,
2987 venue_order_id,
2988 account_id,
2989 TradeId::new(report.trade_id.to_string()),
2990 identity.order_side,
2991 identity.order_type,
2992 last_qty,
2993 last_px,
2994 commission_currency,
2995 liquidity_side,
2996 UUID4::new(),
2997 ts_event,
2998 ts_init,
2999 false,
3000 None,
3001 Some(commission_money),
3002 None,
3003 );
3004
3005 state.insert_filled(client_order_id);
3006 emitter.send_order_event(OrderEventAny::Filled(filled));
3007
3008 let cumulative_qty = parse_spot_execution_report_decimal(
3009 report,
3010 &report.cumulative_filled_qty,
3011 "cumulative_filled_qty",
3012 );
3013 let original_qty =
3014 parse_spot_execution_report_decimal(report, &report.original_qty, "original_qty");
3015 if let (Some(original_qty), Some(cumulative_qty)) = (original_qty, cumulative_qty)
3016 && original_qty <= cumulative_qty
3017 {
3018 state.cleanup_terminal(client_order_id);
3019 }
3020 }
3021 BinanceSpotExecutionType::Replaced => {
3022 log::debug!(
3025 "Order replaced: client_order_id={client_order_id}, venue_order_id={venue_order_id}"
3026 );
3027 }
3028 BinanceSpotExecutionType::Canceled | BinanceSpotExecutionType::TradePrevention => {
3029 ensure_accepted_emitted(
3030 client_order_id,
3031 account_id,
3032 venue_order_id,
3033 identity,
3034 emitter,
3035 state,
3036 ts_init,
3037 );
3038 let canceled = OrderCanceled::new(
3039 emitter.trader_id(),
3040 identity.strategy_id,
3041 identity.instrument_id,
3042 client_order_id,
3043 UUID4::new(),
3044 ts_event,
3045 ts_init,
3046 false,
3047 Some(venue_order_id),
3048 Some(account_id),
3049 );
3050 state.cleanup_terminal(client_order_id);
3051 emitter.send_order_event(OrderEventAny::Canceled(canceled));
3052 }
3053 BinanceSpotExecutionType::Expired => {
3054 ensure_accepted_emitted(
3055 client_order_id,
3056 account_id,
3057 venue_order_id,
3058 identity,
3059 emitter,
3060 state,
3061 ts_init,
3062 );
3063 state.cleanup_terminal(client_order_id);
3064
3065 if treat_expired_as_canceled {
3066 let canceled = OrderCanceled::new(
3067 emitter.trader_id(),
3068 identity.strategy_id,
3069 identity.instrument_id,
3070 client_order_id,
3071 UUID4::new(),
3072 ts_event,
3073 ts_init,
3074 false,
3075 Some(venue_order_id),
3076 Some(account_id),
3077 );
3078 emitter.send_order_event(OrderEventAny::Canceled(canceled));
3079 } else {
3080 let expired = OrderExpired::new(
3081 emitter.trader_id(),
3082 identity.strategy_id,
3083 identity.instrument_id,
3084 client_order_id,
3085 UUID4::new(),
3086 ts_event,
3087 ts_init,
3088 false,
3089 Some(venue_order_id),
3090 Some(account_id),
3091 );
3092 emitter.send_order_event(OrderEventAny::Expired(expired));
3093 }
3094 }
3095 BinanceSpotExecutionType::Rejected => {
3096 let reason = if report.reject_reason.is_empty() {
3097 Ustr::from("Order rejected by venue")
3098 } else {
3099 Ustr::from(&report.reject_reason)
3100 };
3101 let due_post_only = report.time_in_force == BinanceTimeInForce::Gtx
3102 || (report.order_type == "LIMIT_MAKER"
3103 && (report.reject_reason.is_empty() || report.reject_reason == "NONE"));
3104 state.cleanup_terminal(client_order_id);
3105 emitter.emit_order_rejected_event(
3106 identity.strategy_id,
3107 identity.instrument_id,
3108 client_order_id,
3109 reason.as_str(),
3110 ts_init,
3111 due_post_only,
3112 );
3113 }
3114 }
3115}
3116
3117fn parse_spot_execution_report_quantity(
3118 report: &BinanceSpotExecutionReport,
3119 raw: &str,
3120 precision: u8,
3121 field: &str,
3122) -> Option<Quantity> {
3123 match parse_required_quantity_at_precision(raw, precision, field) {
3124 Ok(value) => Some(value),
3125 Err(e) => {
3126 warn_invalid_spot_execution_report_field(report, field, &e);
3127 None
3128 }
3129 }
3130}
3131
3132fn parse_spot_execution_report_price(
3133 report: &BinanceSpotExecutionReport,
3134 raw: &str,
3135 precision: u8,
3136 field: &str,
3137) -> Option<Price> {
3138 match parse_required_price_at_precision(raw, precision, field) {
3139 Ok(value) => Some(value),
3140 Err(e) => {
3141 warn_invalid_spot_execution_report_field(report, field, &e);
3142 None
3143 }
3144 }
3145}
3146
3147fn parse_spot_execution_report_decimal(
3148 report: &BinanceSpotExecutionReport,
3149 raw: &str,
3150 field: &str,
3151) -> Option<Decimal> {
3152 match parse_required_decimal(raw, field) {
3153 Ok(value) => Some(value),
3154 Err(e) => {
3155 warn_invalid_spot_execution_report_field(report, field, &e);
3156 None
3157 }
3158 }
3159}
3160
3161fn warn_invalid_spot_execution_report_field(
3162 report: &BinanceSpotExecutionReport,
3163 field: &str,
3164 error: &anyhow::Error,
3165) {
3166 log::warn!(
3167 "Failed to parse Spot execution report {field} for symbol={}, order_id={}, \
3168 trade_id={}, client_order_id={}: {error}",
3169 report.symbol,
3170 report.order_id,
3171 report.trade_id,
3172 report.client_order_id,
3173 );
3174}
3175
3176#[expect(clippy::too_many_arguments)]
3178fn dispatch_untracked_execution_report(
3179 report: &BinanceSpotExecutionReport,
3180 emitter: &ExecutionEventEmitter,
3181 _http_client: &BinanceSpotHttpClient,
3182 account_id: AccountId,
3183 treat_expired_as_canceled: bool,
3184 seen_trade_ids: &std::sync::Arc<Mutex<FifoCache<(Ustr, i64), 10_000>>>,
3185 instrument_id: InstrumentId,
3186 price_precision: u8,
3187 size_precision: u8,
3188 ts_init: UnixNanos,
3189) {
3190 match report.execution_type {
3191 BinanceSpotExecutionType::Trade => {
3192 let dedup_key = (report.symbol, report.trade_id);
3193 let mut guard = seen_trade_ids.lock();
3194 let is_duplicate = guard.contains(&dedup_key);
3195 guard.add(dedup_key);
3196 drop(guard);
3197
3198 if is_duplicate {
3199 log::debug!(
3200 "Duplicate trade_id={} for {}, skipping",
3201 report.trade_id,
3202 report.symbol
3203 );
3204 return;
3205 }
3206
3207 match parse_spot_exec_report_to_order_status(
3208 report,
3209 instrument_id,
3210 price_precision,
3211 size_precision,
3212 account_id,
3213 treat_expired_as_canceled,
3214 ts_init,
3215 ) {
3216 Ok(status) => emitter.send_order_status_report(status),
3217 Err(e) => log::error!("Failed to parse order status report: {e}"),
3218 }
3219
3220 match parse_spot_exec_report_to_fill(
3221 report,
3222 instrument_id,
3223 price_precision,
3224 size_precision,
3225 account_id,
3226 ts_init,
3227 ) {
3228 Ok(fill) => emitter.send_fill_report(fill),
3229 Err(e) => log::error!("Failed to parse fill report: {e}"),
3230 }
3231 }
3232 BinanceSpotExecutionType::New
3233 | BinanceSpotExecutionType::Canceled
3234 | BinanceSpotExecutionType::Replaced
3235 | BinanceSpotExecutionType::Rejected
3236 | BinanceSpotExecutionType::Expired
3237 | BinanceSpotExecutionType::TradePrevention => {
3238 match parse_spot_exec_report_to_order_status(
3239 report,
3240 instrument_id,
3241 price_precision,
3242 size_precision,
3243 account_id,
3244 treat_expired_as_canceled,
3245 ts_init,
3246 ) {
3247 Ok(status) => emitter.send_order_status_report(status),
3248 Err(e) => log::error!("Failed to parse order status report: {e}"),
3249 }
3250 }
3251 }
3252}
3253
3254fn is_spot_post_only_rejection(error: &BinanceSpotHttpError) -> bool {
3256 match error {
3257 BinanceSpotHttpError::BinanceError { code, message } => {
3258 *code == BINANCE_GTX_ORDER_REJECT_CODE
3259 || (*code == BINANCE_NEW_ORDER_REJECTED_CODE
3260 && message == BINANCE_SPOT_POST_ONLY_REJECT_MSG)
3261 }
3262 _ => false,
3263 }
3264}
3265
3266fn is_structured_venue_rejection(err: &anyhow::Error) -> bool {
3267 err.downcast_ref::<BinanceSpotHttpError>()
3268 .is_some_and(|be| matches!(be, BinanceSpotHttpError::BinanceError { .. }))
3269}
3270
3271fn is_ambiguous_submit_error(err: &anyhow::Error) -> bool {
3272 err.downcast_ref::<BinanceSpotHttpError>()
3273 .is_some_and(|be| {
3274 matches!(
3275 be,
3276 BinanceSpotHttpError::BinanceError {
3277 code: BINANCE_UNEXPECTED_RESPONSE_CODE | BINANCE_STATUS_UNKNOWN_CODE,
3278 ..
3279 }
3280 )
3281 })
3282}
3283
3284fn is_local_command_failure(err: &anyhow::Error) -> bool {
3285 err.downcast_ref::<BinanceSpotHttpError>()
3286 .is_some_and(is_local_http_command_failure)
3287}
3288
3289fn is_local_http_command_failure(err: &BinanceSpotHttpError) -> bool {
3290 matches!(
3291 err,
3292 BinanceSpotHttpError::MissingCredentials | BinanceSpotHttpError::ValidationError(_)
3293 )
3294}
3295
3296#[cfg(test)]
3297mod tests {
3298 use nautilus_common::messages::ExecutionEvent;
3299 use nautilus_core::time::get_atomic_clock_realtime;
3300 use nautilus_model::{
3301 enums::{AccountType, LiquiditySide, OrderSide},
3302 identifiers::{StrategyId, TraderId},
3303 };
3304 use rstest::rstest;
3305
3306 use super::*;
3307 use crate::common::enums::BinanceEnvironment;
3308
3309 #[rstest]
3310 #[case::live(BinanceEnvironment::Live, BINANCE_SPOT_SBE_WS_API_URL)]
3311 #[case::testnet(BinanceEnvironment::Testnet, BINANCE_SPOT_SBE_WS_API_TESTNET_URL)]
3312 #[case::demo(BinanceEnvironment::Demo, BINANCE_SPOT_SBE_WS_API_DEMO_URL)]
3313 fn test_resolve_ws_trading_url_uses_environment_default(
3314 #[case] environment: BinanceEnvironment,
3315 #[case] expected: &str,
3316 ) {
3317 assert_eq!(
3318 BinanceSpotExecutionClient::resolve_ws_trading_url(None, environment),
3319 expected
3320 );
3321 }
3322
3323 #[rstest]
3324 fn test_resolve_ws_trading_url_preserves_override() {
3325 let expected = "wss://example.com/ws-api/v3";
3326
3327 assert_eq!(
3328 BinanceSpotExecutionClient::resolve_ws_trading_url(
3329 Some(expected.to_string()),
3330 BinanceEnvironment::Testnet,
3331 ),
3332 expected
3333 );
3334 }
3335
3336 #[rstest]
3337 fn test_dispatch_ws_trading_message_emits_cancel_rejected_and_clears_pending_request() {
3338 let tasks = TaskGroup::new();
3339 let task_spawner = tasks.spawner().expect("task spawner");
3340 let clock = get_atomic_clock_realtime();
3341 let (emitter, mut rx) = create_test_emitter(clock);
3342 let http_client = create_test_http_client(clock);
3343 let dispatch_state = create_tracked_dispatch_state(
3344 ClientOrderId::from("TEST"),
3345 InstrumentId::from("BTCUSDT.BINANCE"),
3346 );
3347 let ws_authenticated = tokio::sync::Notify::new();
3348 let ws_user_data_subscribed = tokio::sync::Notify::new();
3349 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
3350 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3351
3352 dispatch_state.pending_requests.insert(
3353 "req-cancel".to_string(),
3354 PendingRequest {
3355 client_order_id: ClientOrderId::from("TEST"),
3356 venue_order_id: Some(VenueOrderId::from("12345")),
3357 operation: PendingOperation::Cancel,
3358 },
3359 );
3360
3361 dispatch_ws_trading_message(
3362 BinanceSpotWsTradingMessage::CancelRejected {
3363 request_id: "req-cancel".to_string(),
3364 code: -2011,
3365 msg: "Unknown order sent".to_string(),
3366 },
3367 &emitter,
3368 &http_client,
3369 AccountId::from("BINANCE-001"),
3370 false,
3371 clock,
3372 &dispatch_state,
3373 &ws_authenticated,
3374 &ws_user_data_subscribed,
3375 &ws_setup_error_tx,
3376 &seen_trade_ids,
3377 &task_spawner,
3378 );
3379
3380 assert!(dispatch_state.pending_requests.get("req-cancel").is_none());
3381
3382 match rx
3383 .try_recv()
3384 .expect("Cancel rejection event should be emitted")
3385 {
3386 ExecutionEvent::Order(OrderEventAny::CancelRejected(event)) => {
3387 assert_eq!(event.client_order_id, ClientOrderId::from("TEST"));
3388 assert_eq!(event.account_id, Some(AccountId::from("BINANCE-001")));
3389 assert!(event.reason.as_str().contains("code=-2011"));
3390 }
3391 other => panic!("Expected CancelRejected event, was {other:?}"),
3392 }
3393 }
3394
3395 #[rstest]
3396 #[case(
3397 BINANCE_UNEXPECTED_RESPONSE_CODE,
3398 "An unexpected response was received from the message bus"
3399 )]
3400 #[case(
3401 BINANCE_STATUS_UNKNOWN_CODE,
3402 "Timeout waiting for response from backend server"
3403 )]
3404 fn test_dispatch_ws_trading_message_unknown_status_keeps_order_registered(
3405 #[case] code: i64,
3406 #[case] msg: &str,
3407 ) {
3408 let tasks = TaskGroup::new();
3409 let task_spawner = tasks.spawner().expect("task spawner");
3410 let clock = get_atomic_clock_realtime();
3411 let (emitter, mut rx) = create_test_emitter(clock);
3412 let http_client = create_test_http_client(clock);
3413 let client_order_id = ClientOrderId::from("TEST");
3414 let dispatch_state =
3415 create_tracked_dispatch_state(client_order_id, InstrumentId::from("BTCUSDT.BINANCE"));
3416 let ws_authenticated = tokio::sync::Notify::new();
3417 let ws_user_data_subscribed = tokio::sync::Notify::new();
3418 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
3419 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3420
3421 dispatch_state.pending_requests.insert(
3422 "req-submit".to_string(),
3423 PendingRequest {
3424 client_order_id,
3425 venue_order_id: None,
3426 operation: PendingOperation::Place,
3427 },
3428 );
3429
3430 dispatch_ws_trading_message(
3431 BinanceSpotWsTradingMessage::OrderRejected {
3432 request_id: "req-submit".to_string(),
3433 code: code as i32,
3434 msg: msg.to_string(),
3435 },
3436 &emitter,
3437 &http_client,
3438 AccountId::from("BINANCE-001"),
3439 false,
3440 clock,
3441 &dispatch_state,
3442 &ws_authenticated,
3443 &ws_user_data_subscribed,
3444 &ws_setup_error_tx,
3445 &seen_trade_ids,
3446 &task_spawner,
3447 );
3448
3449 assert!(dispatch_state.pending_requests.get("req-submit").is_none());
3450 assert!(
3451 dispatch_state
3452 .order_identities
3453 .get(&client_order_id)
3454 .is_some()
3455 );
3456 assert!(rx.try_recv().is_err());
3457 }
3458
3459 #[rstest]
3460 fn test_dispatch_ws_trading_message_definite_submit_rejection_emits_order_rejected() {
3461 let tasks = TaskGroup::new();
3462 let task_spawner = tasks.spawner().expect("task spawner");
3463 let clock = get_atomic_clock_realtime();
3464 let (emitter, mut rx) = create_test_emitter(clock);
3465 let http_client = create_test_http_client(clock);
3466 let client_order_id = ClientOrderId::from("TEST");
3467 let dispatch_state =
3468 create_tracked_dispatch_state(client_order_id, InstrumentId::from("BTCUSDT.BINANCE"));
3469 let ws_authenticated = tokio::sync::Notify::new();
3470 let ws_user_data_subscribed = tokio::sync::Notify::new();
3471 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
3472 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3473
3474 dispatch_state.pending_requests.insert(
3475 "req-submit".to_string(),
3476 PendingRequest {
3477 client_order_id,
3478 venue_order_id: None,
3479 operation: PendingOperation::Place,
3480 },
3481 );
3482
3483 dispatch_ws_trading_message(
3484 BinanceSpotWsTradingMessage::OrderRejected {
3485 request_id: "req-submit".to_string(),
3486 code: BINANCE_NEW_ORDER_REJECTED_CODE as i32,
3487 msg: BINANCE_SPOT_POST_ONLY_REJECT_MSG.to_string(),
3488 },
3489 &emitter,
3490 &http_client,
3491 AccountId::from("BINANCE-001"),
3492 false,
3493 clock,
3494 &dispatch_state,
3495 &ws_authenticated,
3496 &ws_user_data_subscribed,
3497 &ws_setup_error_tx,
3498 &seen_trade_ids,
3499 &task_spawner,
3500 );
3501
3502 assert!(dispatch_state.pending_requests.get("req-submit").is_none());
3503 assert!(
3504 dispatch_state
3505 .order_identities
3506 .get(&client_order_id)
3507 .is_none()
3508 );
3509
3510 match rx
3511 .try_recv()
3512 .expect("OrderRejected event should be emitted")
3513 {
3514 ExecutionEvent::Order(OrderEventAny::Rejected(event)) => {
3515 assert_eq!(event.client_order_id, client_order_id);
3516 assert_eq!(event.account_id, AccountId::from("BINANCE-001"));
3517 assert!(event.reason.as_str().contains("code=-2010"));
3518 assert!(event.due_post_only);
3519 }
3520 other => panic!("Expected OrderRejected event, was {other:?}"),
3521 }
3522 }
3523
3524 #[rstest]
3525 fn test_dispatch_ws_trading_message_emits_modify_rejected_and_clears_pending_request() {
3526 let tasks = TaskGroup::new();
3527 let task_spawner = tasks.spawner().expect("task spawner");
3528 let clock = get_atomic_clock_realtime();
3529 let (emitter, mut rx) = create_test_emitter(clock);
3530 let http_client = create_test_http_client(clock);
3531 let dispatch_state = create_tracked_dispatch_state(
3532 ClientOrderId::from("TEST"),
3533 InstrumentId::from("BTCUSDT.BINANCE"),
3534 );
3535 let ws_authenticated = tokio::sync::Notify::new();
3536 let ws_user_data_subscribed = tokio::sync::Notify::new();
3537 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
3538 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3539
3540 dispatch_state.pending_requests.insert(
3541 "req-modify".to_string(),
3542 PendingRequest {
3543 client_order_id: ClientOrderId::from("TEST"),
3544 venue_order_id: Some(VenueOrderId::from("12345")),
3545 operation: PendingOperation::Modify,
3546 },
3547 );
3548
3549 dispatch_ws_trading_message(
3550 BinanceSpotWsTradingMessage::CancelReplaceRejected {
3551 request_id: "req-modify".to_string(),
3552 code: -2021,
3553 msg: "Order cancel-replace partially failed".to_string(),
3554 },
3555 &emitter,
3556 &http_client,
3557 AccountId::from("BINANCE-001"),
3558 false,
3559 clock,
3560 &dispatch_state,
3561 &ws_authenticated,
3562 &ws_user_data_subscribed,
3563 &ws_setup_error_tx,
3564 &seen_trade_ids,
3565 &task_spawner,
3566 );
3567
3568 assert!(dispatch_state.pending_requests.get("req-modify").is_none());
3569
3570 match rx
3571 .try_recv()
3572 .expect("Modify rejection event should be emitted")
3573 {
3574 ExecutionEvent::Order(OrderEventAny::ModifyRejected(event)) => {
3575 assert_eq!(event.client_order_id, ClientOrderId::from("TEST"));
3576 assert_eq!(event.account_id, Some(AccountId::from("BINANCE-001")));
3577 assert!(event.reason.as_str().contains("code=-2021"));
3578 }
3579 other => panic!("Expected ModifyRejected event, was {other:?}"),
3580 }
3581 }
3582
3583 fn create_test_emitter(
3584 clock: &'static AtomicTime,
3585 ) -> (
3586 ExecutionEventEmitter,
3587 tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3588 ) {
3589 let mut emitter = ExecutionEventEmitter::new(
3590 clock,
3591 TraderId::from("TESTER-001"),
3592 AccountId::from("BINANCE-001"),
3593 AccountType::Cash,
3594 None,
3595 );
3596 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
3597 emitter.set_sender(tx);
3598 (emitter, rx)
3599 }
3600
3601 fn create_test_http_client(clock: &'static AtomicTime) -> BinanceSpotHttpClient {
3602 BinanceSpotHttpClient::new(
3603 BinanceEnvironment::Live,
3604 clock,
3605 None,
3606 None,
3607 None,
3608 None,
3609 None,
3610 None,
3611 )
3612 .expect("Test HTTP client should be created")
3613 }
3614
3615 fn create_tracked_dispatch_state(
3616 client_order_id: ClientOrderId,
3617 instrument_id: InstrumentId,
3618 ) -> WsDispatchState {
3619 let dispatch_state = WsDispatchState::default();
3620 dispatch_state.order_identities.insert(
3621 client_order_id,
3622 OrderIdentity {
3623 instrument_id,
3624 strategy_id: StrategyId::from("TEST-STRATEGY"),
3625 order_side: OrderSide::Buy,
3626 order_type: OrderType::Limit,
3627 price: None,
3628 quantity: Quantity::from("1"),
3629 venue_position_id: None,
3630 },
3631 );
3632 dispatch_state
3633 }
3634
3635 #[rstest]
3636 fn test_http_submit_success_defers_acceptance_to_user_stream() {
3637 let clock = get_atomic_clock_realtime();
3638 let (emitter, mut rx) = create_test_emitter(clock);
3639 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
3640 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
3641 let dispatch_state = Arc::new(create_tracked_dispatch_state(
3642 client_order_id,
3643 instrument_id,
3644 ));
3645 handle_spot_order_submit_success(client_order_id, VenueOrderId::from("12345678"));
3646
3647 assert!(!dispatch_state.has_emitted_accepted(&client_order_id));
3648 assert!(rx.try_recv().is_err());
3649 let new_json = crate::common::testing::load_fixture_string(
3650 "spot/user_data_json/execution_report_new.json",
3651 );
3652 let report: BinanceSpotExecutionReport = serde_json::from_str(&new_json).unwrap();
3653 let identity = dispatch_state
3654 .order_identities
3655 .get(&client_order_id)
3656 .unwrap()
3657 .clone();
3658 dispatch_tracked_execution_report(
3659 &report,
3660 &emitter,
3661 AccountId::from("BINANCE-001"),
3662 false,
3663 &dispatch_state,
3664 &Arc::new(Mutex::new(FifoCache::new())),
3665 client_order_id,
3666 &identity,
3667 instrument_id,
3668 2,
3669 8,
3670 clock.get_time_ns(),
3671 );
3672
3673 assert!(matches!(
3674 rx.try_recv(),
3675 Ok(ExecutionEvent::Order(OrderEventAny::Accepted(event)))
3676 if event.client_order_id == client_order_id
3677 && event.venue_order_id == VenueOrderId::from("12345678")
3678 ));
3679 assert!(rx.try_recv().is_err());
3680 }
3681
3682 #[rstest]
3683 #[case::gtx(
3684 BinanceSpotHttpError::BinanceError {
3685 code: BINANCE_GTX_ORDER_REJECT_CODE,
3686 message: "Order would immediately trigger.".to_string(),
3687 },
3688 true,
3689 )]
3690 #[case::spot_post_only(
3691 BinanceSpotHttpError::BinanceError {
3692 code: BINANCE_NEW_ORDER_REJECTED_CODE,
3693 message: BINANCE_SPOT_POST_ONLY_REJECT_MSG.to_string(),
3694 },
3695 true,
3696 )]
3697 #[case::new_order_rejected_other_message(
3698 BinanceSpotHttpError::BinanceError {
3699 code: BINANCE_NEW_ORDER_REJECTED_CODE,
3700 message: "Insufficient balance.".to_string(),
3701 },
3702 false,
3703 )]
3704 #[case::unrelated_code(
3705 BinanceSpotHttpError::BinanceError {
3706 code: -2011,
3707 message: "Unknown order sent.".to_string(),
3708 },
3709 false,
3710 )]
3711 #[case::non_binance_error(
3712 BinanceSpotHttpError::NetworkError("connection reset".to_string()),
3713 false,
3714 )]
3715 fn test_is_spot_post_only_rejection(
3716 #[case] error: BinanceSpotHttpError,
3717 #[case] expected: bool,
3718 ) {
3719 assert_eq!(is_spot_post_only_rejection(&error), expected);
3720 }
3721
3722 #[rstest]
3723 #[case(BINANCE_UNEXPECTED_RESPONSE_CODE)]
3724 #[case(BINANCE_STATUS_UNKNOWN_CODE)]
3725 fn test_unknown_status_submit_error_is_ambiguous(#[case] code: i64) {
3726 let err = anyhow::Error::new(BinanceSpotHttpError::BinanceError {
3727 code,
3728 message: "test error".to_string(),
3729 });
3730 assert!(is_ambiguous_submit_error(&err));
3731 assert!(is_structured_venue_rejection(&err));
3732 }
3733
3734 #[rstest]
3735 fn test_other_structured_submit_error_is_not_ambiguous() {
3736 let err = anyhow::Error::new(BinanceSpotHttpError::BinanceError {
3737 code: BINANCE_GTX_ORDER_REJECT_CODE,
3738 message: "test error".to_string(),
3739 });
3740 assert!(!is_ambiguous_submit_error(&err));
3741 assert!(is_structured_venue_rejection(&err));
3742 }
3743
3744 #[rstest]
3745 fn test_dispatch_tracked_execution_report_trade_dedup() {
3746 let tasks = TaskGroup::new();
3747 let task_spawner = tasks.spawner().expect("task spawner");
3748 let clock = get_atomic_clock_realtime();
3749 let (emitter, mut rx) = create_test_emitter(clock);
3750 let http_client = create_test_http_client(clock);
3751 let client_order_id = ClientOrderId::from("x-TD67BGP9-T0000000000000");
3752 let dispatch_state = create_tracked_dispatch_state(
3753 ClientOrderId::from("O-20200101-000000-000-000-0"),
3754 InstrumentId::from("ETHUSDT.BINANCE"),
3755 );
3756 let ws_authenticated = tokio::sync::Notify::new();
3757 let ws_user_data_subscribed = tokio::sync::Notify::new();
3758 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
3759 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3760
3761 let trade_json = crate::common::testing::load_fixture_string(
3762 "spot/user_data_json/execution_report_trade.json",
3763 );
3764 let report: BinanceSpotExecutionReport = serde_json::from_str(&trade_json).unwrap();
3765
3766 dispatch_ws_trading_message(
3767 BinanceSpotWsTradingMessage::ExecutionReport(Box::new(report.clone())),
3768 &emitter,
3769 &http_client,
3770 AccountId::from("BINANCE-001"),
3771 false,
3772 clock,
3773 &dispatch_state,
3774 &ws_authenticated,
3775 &ws_user_data_subscribed,
3776 &ws_setup_error_tx,
3777 &seen_trade_ids,
3778 &task_spawner,
3779 );
3780 dispatch_ws_trading_message(
3781 BinanceSpotWsTradingMessage::ExecutionReport(Box::new(report)),
3782 &emitter,
3783 &http_client,
3784 AccountId::from("BINANCE-001"),
3785 false,
3786 clock,
3787 &dispatch_state,
3788 &ws_authenticated,
3789 &ws_user_data_subscribed,
3790 &ws_setup_error_tx,
3791 &seen_trade_ids,
3792 &task_spawner,
3793 );
3794
3795 let mut events = Vec::new();
3796 while let Ok(event) = rx.try_recv() {
3797 events.push(event);
3798 }
3799
3800 let fills: Vec<_> = events
3801 .iter()
3802 .filter(|e| matches!(e, ExecutionEvent::Order(OrderEventAny::Filled(_))))
3803 .collect();
3804 assert_eq!(fills.len(), 1, "duplicate trade should be deduped");
3805
3806 match fills[0] {
3807 ExecutionEvent::Order(OrderEventAny::Filled(fill)) => {
3808 assert_eq!(
3809 fill.client_order_id,
3810 ClientOrderId::from("O-20200101-000000-000-000-0"),
3811 );
3812 assert_eq!(fill.trade_id, TradeId::new("98765432"));
3813 assert_eq!(fill.liquidity_side, LiquiditySide::Maker);
3814 }
3815 _ => unreachable!(),
3816 }
3817 let _ = client_order_id;
3818 }
3819
3820 #[rstest]
3821 fn test_dispatch_tracked_execution_report_invalid_fill_qty_skips_filled_event() {
3822 let tasks = TaskGroup::new();
3823 let task_spawner = tasks.spawner().expect("task spawner");
3824 let clock = get_atomic_clock_realtime();
3825 let (emitter, mut rx) = create_test_emitter(clock);
3826 let http_client = create_test_http_client(clock);
3827 let dispatch_state = create_tracked_dispatch_state(
3828 ClientOrderId::from("O-20200101-000000-000-000-0"),
3829 InstrumentId::from("ETHUSDT.BINANCE"),
3830 );
3831 let ws_authenticated = tokio::sync::Notify::new();
3832 let ws_user_data_subscribed = tokio::sync::Notify::new();
3833 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
3834 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3835
3836 let trade_json = crate::common::testing::load_fixture_string(
3837 "spot/user_data_json/execution_report_trade.json",
3838 );
3839 let mut report: BinanceSpotExecutionReport = serde_json::from_str(&trade_json).unwrap();
3840 report.last_filled_qty = "not-a-number".to_string();
3841
3842 dispatch_ws_trading_message(
3843 BinanceSpotWsTradingMessage::ExecutionReport(Box::new(report)),
3844 &emitter,
3845 &http_client,
3846 AccountId::from("BINANCE-001"),
3847 false,
3848 clock,
3849 &dispatch_state,
3850 &ws_authenticated,
3851 &ws_user_data_subscribed,
3852 &ws_setup_error_tx,
3853 &seen_trade_ids,
3854 &task_spawner,
3855 );
3856
3857 let mut events = Vec::new();
3858 while let Ok(event) = rx.try_recv() {
3859 events.push(event);
3860 }
3861
3862 assert!(
3863 events
3864 .iter()
3865 .all(|e| !matches!(e, ExecutionEvent::Order(OrderEventAny::Filled(_)))),
3866 "invalid fill quantity must not emit OrderFilled",
3867 );
3868 }
3869
3870 #[rstest]
3871 fn test_dispatch_execution_report_invalid_client_order_id_emits_nothing() {
3872 let clock = get_atomic_clock_realtime();
3873 let (emitter, mut rx) = create_test_emitter(clock);
3874 let http_client = create_test_http_client(clock);
3875 let dispatch_state = WsDispatchState::default();
3876 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3877 let json = crate::common::testing::load_fixture_string(
3878 "spot/user_data_json/execution_report_new.json",
3879 );
3880 let mut report: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
3881 report.client_order_id = "x-TD67BGP9-R".to_string();
3882
3883 dispatch_execution_report(
3884 &report,
3885 &emitter,
3886 &http_client,
3887 AccountId::from("BINANCE-001"),
3888 false,
3889 &dispatch_state,
3890 &seen_trade_ids,
3891 clock.get_time_ns(),
3892 );
3893
3894 assert!(rx.try_recv().is_err());
3895 assert!(dispatch_state.order_identities.is_empty());
3896 }
3897
3898 #[rstest]
3899 #[case::as_expired(false, OrderStatus::Expired)]
3900 #[case::as_canceled(true, OrderStatus::Canceled)]
3901 fn test_normalize_spot_order_status_report_expired_respects_config(
3902 #[case] treat_expired_as_canceled: bool,
3903 #[case] expected: OrderStatus,
3904 ) {
3905 let clock = get_atomic_clock_realtime();
3906 let json = crate::common::testing::load_fixture_string(
3907 "spot/user_data_json/execution_report_expired.json",
3908 );
3909 let msg: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
3910 let mut report = parse_spot_exec_report_to_order_status(
3911 &msg,
3912 InstrumentId::from("ETHUSDT.BINANCE"),
3913 2,
3914 5,
3915 AccountId::from("BINANCE-001"),
3916 false,
3917 clock.get_time_ns(),
3918 )
3919 .unwrap();
3920 let mut reports = vec![report.clone()];
3921
3922 normalize_spot_order_status_report(&mut report, treat_expired_as_canceled);
3923 normalize_spot_order_status_reports(&mut reports, treat_expired_as_canceled);
3924
3925 assert_eq!(report.order_status, expected);
3926 assert_eq!(reports[0].order_status, expected);
3927 }
3928
3929 #[rstest]
3930 #[case::as_expired(false)]
3931 #[case::as_canceled(true)]
3932 fn test_dispatch_tracked_execution_report_expired_respects_config(
3933 #[case] treat_expired_as_canceled: bool,
3934 ) {
3935 let clock = get_atomic_clock_realtime();
3936 let (emitter, mut rx) = create_test_emitter(clock);
3937 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
3938 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
3939 let dispatch_state = WsDispatchState::default();
3940 dispatch_state.insert_accepted(client_order_id);
3941 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3942 let identity = OrderIdentity {
3943 instrument_id,
3944 strategy_id: StrategyId::from("TEST-STRATEGY"),
3945 order_side: OrderSide::Buy,
3946 order_type: OrderType::Limit,
3947 price: None,
3948 quantity: Quantity::from("1"),
3949 venue_position_id: None,
3950 };
3951
3952 let json = crate::common::testing::load_fixture_string(
3953 "spot/user_data_json/execution_report_expired.json",
3954 );
3955 let report: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
3956
3957 dispatch_tracked_execution_report(
3958 &report,
3959 &emitter,
3960 AccountId::from("BINANCE-001"),
3961 treat_expired_as_canceled,
3962 &dispatch_state,
3963 &seen_trade_ids,
3964 client_order_id,
3965 &identity,
3966 instrument_id,
3967 2,
3968 5,
3969 clock.get_time_ns(),
3970 );
3971
3972 let event = rx.try_recv().expect("terminal order event expected");
3973 match (treat_expired_as_canceled, event) {
3974 (true, ExecutionEvent::Order(OrderEventAny::Canceled(event))) => {
3975 assert_eq!(event.client_order_id, client_order_id);
3976 }
3977 (false, ExecutionEvent::Order(OrderEventAny::Expired(event))) => {
3978 assert_eq!(event.client_order_id, client_order_id);
3979 }
3980 (_, other) => panic!("Expected terminal expired/canceled event, was {other:?}"),
3981 }
3982 assert!(rx.try_recv().is_err());
3983 }
3984
3985 #[rstest]
3986 fn test_dispatch_tracked_execution_report_rejected_gtx_sets_post_only() {
3987 let tasks = TaskGroup::new();
3988 let task_spawner = tasks.spawner().expect("task spawner");
3989 let clock = get_atomic_clock_realtime();
3990 let (emitter, mut rx) = create_test_emitter(clock);
3991 let http_client = create_test_http_client(clock);
3992 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-1");
3993 let dispatch_state =
3994 create_tracked_dispatch_state(client_order_id, InstrumentId::from("ETHUSDT.BINANCE"));
3995 let ws_authenticated = tokio::sync::Notify::new();
3996 let ws_user_data_subscribed = tokio::sync::Notify::new();
3997 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
3998 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
3999
4000 let encoded = encode_broker_id(&client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
4001 let report_json = format!(
4002 r#"{{
4003 "e":"executionReport","E":1709654400000,"s":"ETHUSDT",
4004 "c":"{encoded}","S":"BUY","o":"LIMIT","f":"GTX",
4005 "q":"1.00000000","p":"2500.00000000","P":"0.00000000",
4006 "x":"REJECTED","X":"REJECTED","r":"NONE","i":12345678,
4007 "l":"0.00000000","z":"0.00000000","L":"0.00000000",
4008 "n":"0","N":null,"T":1709654400000,"t":-1,"w":false,"m":false,
4009 "O":1709654400000,"Z":"0.00000000","C":""
4010 }}"#,
4011 );
4012 let report: BinanceSpotExecutionReport = serde_json::from_str(&report_json).unwrap();
4013
4014 dispatch_ws_trading_message(
4015 BinanceSpotWsTradingMessage::ExecutionReport(Box::new(report)),
4016 &emitter,
4017 &http_client,
4018 AccountId::from("BINANCE-001"),
4019 false,
4020 clock,
4021 &dispatch_state,
4022 &ws_authenticated,
4023 &ws_user_data_subscribed,
4024 &ws_setup_error_tx,
4025 &seen_trade_ids,
4026 &task_spawner,
4027 );
4028
4029 match rx.try_recv().expect("OrderRejected event expected") {
4030 ExecutionEvent::Order(OrderEventAny::Rejected(event)) => {
4031 assert_eq!(event.client_order_id, client_order_id);
4032 assert_eq!(event.account_id, AccountId::from("BINANCE-001"));
4033 assert!(event.due_post_only);
4034 }
4035 other => panic!("Expected OrderRejected event, was {other:?}"),
4036 }
4037 }
4038}