1use std::{future::Future, sync::Arc, time::Duration};
19
20use ahash::{AHashMap, AHashSet};
21use anyhow::Context;
22use async_trait::async_trait;
23use jiff::Timestamp;
24use nautilus_common::{
25 cache::fifo::FifoCache,
26 clients::ExecutionClient,
27 enums::LogLevel,
28 live::runner::get_exec_event_sender,
29 messages::execution::{
30 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
31 GenerateFillReportsBuilder, GenerateOrderStatusReport, GenerateOrderStatusReports,
32 GenerateOrderStatusReportsBuilder, GeneratePositionStatusReports,
33 GeneratePositionStatusReportsBuilder, ModifyOrder, QueryAccount, QueryOrder, SubmitOrder,
34 SubmitOrderList,
35 },
36};
37use nautilus_core::{
38 DurationNanos, Params, UUID4, UnixNanos,
39 datetime::NANOSECONDS_IN_MILLISECOND,
40 string::secret::SecretString,
41 time::{AtomicTime, get_atomic_clock_realtime},
42};
43use nautilus_live::{
44 ExecutionClientCore, ExecutionEventEmitter, SocketControlFactory,
45 execution::failure::CommandFailure,
46 task::{TaskGroup, TaskGroupGuard, TaskRef, TaskSpawner},
47};
48use nautilus_model::{
49 accounts::AccountAny,
50 enums::{ContingencyType, LiquiditySide, OmsType, OrderStatus, OrderType},
51 events::{
52 AccountState, OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDeniedReason,
53 OrderEventAny, OrderExpired, OrderFilled, OrderModifyRejected, OrderRejected, OrderUpdated,
54 },
55 identifiers::{
56 AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, Venue,
57 VenueOrderId,
58 },
59 instruments::Instrument,
60 orders::{Order, OrderAny},
61 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
62 types::{AccountBalance, Currency, MarginBalance, Money, Price, Quantity},
63};
64use parking_lot::Mutex;
65use rust_decimal::Decimal;
66use ustr::Ustr;
67
68use super::websocket::trading::{
69 client::BinanceSpotWsTradingClient,
70 messages::BinanceSpotWsTradingMessage,
71 parse::{
72 parse_spot_account_position, parse_spot_exec_report_to_fill,
73 parse_spot_exec_report_to_order_status,
74 },
75 user_data::{BinanceSpotExecutionReport, BinanceSpotExecutionType},
76};
77use crate::{
78 common::{
79 consts::{
80 BINANCE_GTX_ORDER_REJECT_CODE, BINANCE_NAUTILUS_SPOT_BROKER_ID,
81 BINANCE_NEW_ORDER_REJECTED_CODE, BINANCE_SPOT_POST_ONLY_REJECT_MSG,
82 BINANCE_SPOT_SBE_WS_API_DEMO_URL, BINANCE_SPOT_SBE_WS_API_TESTNET_URL,
83 BINANCE_SPOT_SBE_WS_API_URL, BINANCE_VENUE, BINANCE_WS_HEARTBEAT_SECS,
84 },
85 credential::resolve_credentials,
86 dispatch::{
87 OrderIdentity, PendingOperation, PendingRequest, WsDispatchState,
88 ensure_accepted_emitted,
89 },
90 encoder::{decode_client_order_id, encode_broker_id},
91 enums::{BinanceEnvironment, BinanceSide, BinanceTimeInForce},
92 failure::{classify_spot_http_failure, classify_venue_failure, sanitize_reason},
93 parse::{
94 parse_micros_or_init, parse_millis_or_init, parse_required_decimal,
95 parse_required_price_at_precision, parse_required_quantity_at_precision,
96 },
97 urls::{get_http_base_url_with_us, get_spot_user_stream_url},
98 },
99 config::BinanceExecutionClientConfig,
100 spot::{
101 enums::{
102 BinanceCancelReplaceMode, BinanceOrderResponseType, BinanceSpotOrderType,
103 order_type_to_binance_spot, time_in_force_to_binance_spot,
104 },
105 http::{
106 client::BinanceSpotHttpClient,
107 error::BinanceSpotHttpError,
108 models::{
109 BinanceCancelOpenOrdersResponse, BinanceCancelOrderListResponse,
110 BinanceCancelOrderResponse,
111 },
112 query::{
113 CancelOrderParams, CancelReplaceOrderParams, NewOcoOrderListParams, NewOrderParams,
114 },
115 },
116 sbe::spot::{
117 list_order_status::ListOrderStatus as SbeListOrderStatus,
118 list_status_type::ListStatusType as SbeListStatusType,
119 order_status::OrderStatus as SbeOrderStatus,
120 },
121 },
122};
123
124const ACCOUNT_TRADES_MAX_INTERVAL_MS: i64 = 24 * 60 * 60 * 1_000;
125
126const ACCOUNT_TRADES_PAGE_LIMIT: u32 = 1_000;
127
128const WS_RECONNECT_SETUP_RETRY_DELAY: Duration = Duration::from_secs(1);
129
130#[derive(Debug)]
136pub struct BinanceSpotExecutionClient {
137 core: ExecutionClientCore,
138 clock: &'static AtomicTime,
139 config: BinanceExecutionClientConfig,
140 emitter: ExecutionEventEmitter,
141 dispatch_state: Arc<WsDispatchState>,
142 lifecycle_lock: Arc<Mutex<()>>,
143 http_client: BinanceSpotHttpClient,
144 socket_factory: SocketControlFactory,
145 ws_trading_client: Option<BinanceSpotWsTradingClient>,
146 ws_trading_dispatch: Option<TaskRef>,
147 ws_user_data_client: Arc<Mutex<Option<BinanceSpotWsTradingClient>>>,
148 ws_user_data_dispatch: Option<TaskRef>,
149 listen_key: Option<SecretString>,
150 us_credentials: Option<(SecretString, SecretString)>,
151 ws_authenticated: Arc<tokio::sync::Notify>,
152 ws_user_data_subscribed: Arc<tokio::sync::Notify>,
153 session_tasks: TaskGroup,
154 pending_tasks: TaskGroup,
155 shutdown_errors: Vec<String>,
156}
157
158impl BinanceSpotExecutionClient {
159 pub fn new(
165 core: ExecutionClientCore,
166 config: BinanceExecutionClientConfig,
167 ) -> anyhow::Result<Self> {
168 config.validate()?;
169 let (api_key, api_secret) = resolve_credentials(
170 config
171 .api_key
172 .as_ref()
173 .map(|value| value.expose_secret().to_owned()),
174 config
175 .api_secret
176 .as_ref()
177 .map(|value| value.expose_secret().to_owned()),
178 config.environment,
179 config.product_type,
180 )?;
181
182 let clock = get_atomic_clock_realtime();
183 let socket_factory = SocketControlFactory::new(core.client_id, Some(*BINANCE_VENUE));
184 let base_url_http = config.base_url_http.clone().or_else(|| {
185 config.us.then(|| {
186 get_http_base_url_with_us(config.product_type, config.environment, true).to_string()
187 })
188 });
189 let proxy_url = config
190 .proxy_url
191 .as_ref()
192 .map(|value| value.expose_secret().to_owned());
193
194 let http_client = BinanceSpotHttpClient::new_with_json_responses(
195 config.environment,
196 clock,
197 Some(api_key.clone()),
198 Some(api_secret.clone()),
199 base_url_http,
200 Some(config.recv_window_ms),
201 None, proxy_url.clone(),
203 config.us,
204 )
205 .context("failed to construct Binance Spot HTTP client")?
206 .with_retry_config(config.retry_config());
207 let emitter = ExecutionEventEmitter::new(
208 clock,
209 core.trader_id,
210 core.account_id,
211 core.account_type,
212 core.base_currency,
213 );
214
215 let ws_trading_client = if config.us {
216 None
217 } else {
218 let url = Some(Self::resolve_ws_trading_url(
219 config.base_url_ws_trading.clone(),
220 config.environment,
221 ));
222 Some(
223 BinanceSpotWsTradingClient::new(
224 url,
225 api_key.clone(),
226 api_secret.clone(),
227 Some(BINANCE_WS_HEARTBEAT_SECS),
228 config.transport_backend,
229 )
230 .with_proxy(proxy_url)
231 .with_recv_window(Some(config.recv_window_ms))
232 .with_socket_control(socket_factory.control("binance-spot-trading")),
233 )
234 };
235 let us_credentials = config
236 .us
237 .then_some((SecretString::from(api_key), SecretString::from(api_secret)));
238
239 let session_tasks = TaskGroup::new();
240 let pending_tasks = TaskGroup::new();
241
242 Ok(Self {
243 core,
244 clock,
245 config,
246 emitter,
247 dispatch_state: Arc::new(WsDispatchState::default()),
248 lifecycle_lock: Arc::new(Mutex::new(())),
249 http_client,
250 socket_factory,
251 ws_trading_client,
252 ws_trading_dispatch: None,
253 ws_user_data_client: Arc::new(Mutex::new(None)),
254 ws_user_data_dispatch: None,
255 listen_key: None,
256 us_credentials,
257 ws_authenticated: Arc::new(tokio::sync::Notify::new()),
258 ws_user_data_subscribed: Arc::new(tokio::sync::Notify::new()),
259 session_tasks,
260 pending_tasks,
261 shutdown_errors: Vec::new(),
262 })
263 }
264
265 fn resolve_ws_trading_url(base_url: Option<String>, environment: BinanceEnvironment) -> String {
266 base_url.unwrap_or_else(|| {
267 match environment {
268 BinanceEnvironment::Live => BINANCE_SPOT_SBE_WS_API_URL,
269 BinanceEnvironment::Testnet => BINANCE_SPOT_SBE_WS_API_TESTNET_URL,
270 BinanceEnvironment::Demo => BINANCE_SPOT_SBE_WS_API_DEMO_URL,
271 }
272 .to_string()
273 })
274 }
275
276 async fn refresh_account_state(&self) -> anyhow::Result<AccountState> {
277 self.http_client
278 .request_account_state(self.core.account_id)
279 .await
280 }
281
282 fn update_account_state(&self) {
283 let http_client = self.http_client.clone();
284 let account_id = self.core.account_id;
285 let emitter = self.emitter.clone();
286 let clock = self.clock;
287
288 self.spawn_task("query_account", async move {
289 let account_state = http_client.request_account_state(account_id).await?;
290 let ts_now = clock.get_time_ns();
291 emitter.emit_account_state(
292 account_state.balances.clone(),
293 account_state.margins.clone(),
294 account_state.is_reported,
295 ts_now,
296 account_state.info,
297 );
298 Ok(())
299 });
300 }
301
302 fn ws_user_data_active(&self) -> bool {
303 let dispatch_running = if self.config.us {
304 self.ws_user_data_dispatch
305 .as_ref()
306 .is_some_and(TaskRef::is_active)
307 } else {
308 self.ws_trading_dispatch
309 .as_ref()
310 .is_some_and(TaskRef::is_active)
311 };
312 let user_data_active = if self.config.us {
313 self.ws_user_data_client
314 .lock()
315 .as_ref()
316 .is_some_and(BinanceSpotWsTradingClient::is_user_data_active)
317 } else {
318 self.ws_trading_client
319 .as_ref()
320 .is_some_and(BinanceSpotWsTradingClient::is_user_data_active)
321 };
322
323 user_data_active && dispatch_running
324 }
325
326 fn ensure_ws_user_data_active(&self) -> anyhow::Result<()> {
327 anyhow::ensure!(
328 self.ws_user_data_active(),
329 "Binance Spot user data stream is not active",
330 );
331 Ok(())
332 }
333
334 fn ws_order_transport_active(&self) -> bool {
335 self.config.use_ws_trading && self.ws_trading_client.is_some() && self.ws_user_data_active()
336 }
337
338 fn submit_order_internal(
339 &self,
340 cmd: &SubmitOrder,
341 params: Option<NewOrderParams>,
342 ) -> anyhow::Result<()> {
343 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
344
345 let event_emitter = self.emitter.clone();
346 let trader_id = self.core.trader_id;
347 let account_id = self.core.account_id;
348 let client_order_id = order.client_order_id();
349 let strategy_id = order.strategy_id();
350 let instrument_id = order.instrument_id();
351 let order_side = order.order_side();
352 let order_type = order.order_type();
353 let quantity = order.quantity();
354 let time_in_force = order.time_in_force();
355 let price = order.price();
356 let trigger_price = order.trigger_price();
357 let is_post_only = order.is_post_only();
358 let is_quote_quantity = order.is_quote_quantity();
359 let display_qty = order.display_qty();
360 let use_gtd = self.config.use_gtd;
361 let clock = self.clock;
362 let ts_init = self.clock.get_time_ns();
363
364 self.dispatch_state.order_identities.insert(
366 client_order_id,
367 OrderIdentity {
368 instrument_id,
369 strategy_id,
370 order_side,
371 order_type,
372 price,
373 quantity,
374 venue_position_id: None,
375 },
376 );
377
378 if let Some(params) = params {
379 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
380 let dispatch_state = self.dispatch_state.clone();
381
382 let request_id = ws_client.next_request_id();
384 dispatch_state.pending_requests.insert(
385 request_id.clone(),
386 PendingRequest {
387 client_order_id,
388 venue_order_id: None,
389 operation: PendingOperation::Place,
390 },
391 );
392
393 self.spawn_task("submit_order_ws", async move {
394 if let Err(e) = ws_client
395 .place_order_with_id(request_id.clone(), params)
396 .await
397 {
398 dispatch_state.pending_requests.remove(&request_id);
399 log::warn!(
400 "WS submit request failed for {client_order_id}, awaiting reconciliation: {e}"
401 );
402 anyhow::bail!("WS submit order failed: {e}");
403 }
404 Ok(())
405 });
406 } else {
407 let http_client = self.http_client.clone();
408 let dispatch_state = self.dispatch_state.clone();
409 log::debug!("WS trading not active, falling back to HTTP for submit_order");
410
411 self.spawn_task("submit_order_http", async move {
412 let result = http_client
413 .submit_order(
414 account_id,
415 instrument_id,
416 client_order_id,
417 order_side,
418 order_type,
419 quantity,
420 time_in_force,
421 price,
422 trigger_price,
423 is_post_only,
424 is_quote_quantity,
425 display_qty,
426 use_gtd,
427 )
428 .await;
429
430 match result {
431 Ok(report) => handle_spot_order_submit_success(
432 client_order_id,
433 report.venue_order_id,
434 ),
435 Err(e) => {
436 let http_error = e.downcast_ref::<BinanceSpotHttpError>();
437 let failure = http_error.map_or_else(
438 || CommandFailure::Ambiguous(e.to_string()),
439 classify_spot_http_failure,
440 );
441
442 match failure {
443 CommandFailure::Ambiguous(reason) => {
444 log::warn!(
445 "Ambiguous submit failure for {client_order_id}, awaiting reconciliation: {reason}"
446 );
447 }
448 CommandFailure::NotSent(reason)
449 | CommandFailure::VenueRejected(reason) => {
450 let due_post_only =
451 http_error.is_some_and(is_spot_post_only_rejection);
452 dispatch_state.cleanup_terminal(client_order_id);
453 let rejected = OrderRejected::new(
454 trader_id,
455 strategy_id,
456 instrument_id,
457 client_order_id,
458 account_id,
459 format!("submit-order-error: {}", sanitize_reason(&reason))
460 .into(),
461 UUID4::new(),
462 ts_init,
463 clock.get_time_ns(),
464 false,
465 due_post_only,
466 );
467 event_emitter.send_order_event(OrderEventAny::Rejected(rejected));
468 }
469 }
470 return Err(e);
471 }
472 }
473 Ok(())
474 });
475 }
476
477 Ok(())
478 }
479
480 fn cancel_order_internal(&self, cmd: &CancelOrder) {
481 let event_emitter = self.emitter.clone();
482 let trader_id = self.core.trader_id;
483 let account_id = self.core.account_id;
484 let clock = self.clock;
485 let command = cmd.clone();
486 let prefer_client_order_id = self
487 .dispatch_state
488 .order_identities
489 .contains_key(&cmd.client_order_id);
490
491 if self.ws_order_transport_active() {
492 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
493 let dispatch_state = self.dispatch_state.clone();
494 let params = build_cancel_order_params(&command, prefer_client_order_id);
495
496 let request_id = ws_client.next_request_id();
498 dispatch_state.pending_requests.insert(
499 request_id.clone(),
500 PendingRequest {
501 client_order_id: command.client_order_id,
502 venue_order_id: command.venue_order_id,
503 operation: PendingOperation::Cancel,
504 },
505 );
506
507 self.spawn_task("cancel_order_ws", async move {
508 if let Err(e) = ws_client
509 .cancel_order_with_id(request_id.clone(), params)
510 .await
511 {
512 dispatch_state.pending_requests.remove(&request_id);
513 log::warn!(
514 "WS cancel request failed for {}, awaiting reconciliation: {e}",
515 command.client_order_id
516 );
517 anyhow::bail!("WS cancel order failed: {e}");
518 }
519 Ok(())
520 });
521 } else {
522 let http_client = self.http_client.clone();
523 let dispatch_state = self.dispatch_state.clone();
524 log::debug!("WS trading not active, falling back to HTTP for cancel_order");
525
526 self.spawn_task("cancel_order_http", async move {
527 let result = http_client
528 .cancel_order(
529 command.instrument_id,
530 if prefer_client_order_id { None } else { command.venue_order_id },
531 Some(command.client_order_id),
532 )
533 .await;
534
535 match result {
536 Ok(venue_order_id) => {
537 dispatch_state.cleanup_terminal(command.client_order_id);
538 let ts_now = clock.get_time_ns();
539 let canceled_event = OrderCanceled::new(
540 trader_id,
541 command.strategy_id,
542 command.instrument_id,
543 command.client_order_id,
544 UUID4::new(),
545 ts_now,
546 ts_now,
547 false,
548 Some(venue_order_id),
549 Some(account_id),
550 None,
551 );
552 event_emitter.send_order_event(OrderEventAny::Canceled(canceled_event));
553 }
554 Err(e) => {
555 let failure = e.downcast_ref::<BinanceSpotHttpError>().map_or_else(
556 || CommandFailure::Ambiguous(e.to_string()),
557 classify_spot_http_failure,
558 );
559
560 match failure {
561 CommandFailure::Ambiguous(reason) => {
562 log::warn!(
563 "Ambiguous cancel failure for {}, awaiting reconciliation: {reason}",
564 command.client_order_id
565 );
566 }
567 CommandFailure::NotSent(reason)
568 | CommandFailure::VenueRejected(reason) => {
569 let ts_now = clock.get_time_ns();
570 let rejected_event = OrderCancelRejected::new(
571 trader_id,
572 command.strategy_id,
573 command.instrument_id,
574 command.client_order_id,
575 format!("cancel-order-error: {}", sanitize_reason(&reason))
576 .into(),
577 UUID4::new(),
578 ts_now,
579 ts_now,
580 false,
581 command.venue_order_id,
582 Some(account_id),
583 );
584 event_emitter
585 .send_order_event(OrderEventAny::CancelRejected(rejected_event));
586 }
587 }
588 return Err(e);
589 }
590 }
591 Ok(())
592 });
593 }
594 }
595
596 fn spawn_task<F>(&self, description: &'static str, fut: F)
597 where
598 F: Future<Output = anyhow::Result<()>> + Send + 'static,
599 {
600 crate::common::execution::spawn_task(&self.pending_tasks, description, fut);
601 }
602
603 fn begin_generation_shutdown(&self) {
604 if let Some(client) = self.ws_trading_client.as_ref() {
605 client.mark_user_data_inactive();
606 client.begin_shutdown();
607 }
608
609 if let Some(client) = self.ws_user_data_client.lock().as_ref() {
610 client.mark_user_data_inactive();
611 client.begin_shutdown();
612 }
613
614 self.core.set_disconnected();
615 self.abort_session_tasks();
616 self.abort_pending_tasks();
617 }
618
619 fn abort_pending_tasks(&self) {
620 crate::common::execution::abort_pending_tasks(&self.pending_tasks);
621 }
622
623 fn abort_session_tasks(&self) {
624 self.session_tasks.begin_shutdown();
625 }
626
627 async fn await_pending_tasks(&self) -> anyhow::Result<()> {
628 crate::common::execution::await_pending_tasks(&self.pending_tasks).await
629 }
630
631 async fn await_session_tasks(&self) -> anyhow::Result<()> {
632 self.session_tasks.begin_shutdown();
633 self.session_tasks
634 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
635 .await
636 .map_err(|e| anyhow::anyhow!("Failed to terminate Binance Spot session tasks: {e}"))?;
637 Ok(())
638 }
639
640 async fn ws_setup_failure(
641 &mut self,
642 mut ws_trading: BinanceSpotWsTradingClient,
643 reason: String,
644 ) -> anyhow::Error {
645 ws_trading.mark_user_data_inactive();
646 log::error!("{reason}; Binance Spot private user data is required for execution");
647
648 self.abort_session_tasks();
649
650 if let Err(e) = ws_trading.disconnect().await {
651 log::warn!("Failed to stop Binance Spot trading WebSocket after setup failure: {e}");
652 }
653 self.disconnect_us_user_data().await;
654
655 if let Err(e) = self.await_session_tasks().await {
656 log::warn!("Failed to drain Binance Spot session tasks after setup failure: {e}");
657 }
658 self.ws_trading_client = Some(ws_trading);
659 anyhow::anyhow!(reason)
660 }
661
662 async fn connect_us_user_data(&mut self) -> anyhow::Result<()> {
663 let (api_key, api_secret) = self
664 .us_credentials
665 .clone()
666 .context("Binance US user data credentials are unavailable")?;
667 let response = self
668 .http_client
669 .inner()
670 .create_listen_key()
671 .await
672 .context("failed to create Binance US listen key")?;
673 let listen_key = response.into_listen_key();
674 self.listen_key = Some(listen_key.clone());
675 let url = get_spot_user_stream_url(
676 self.config.base_url_ws.as_deref(),
677 listen_key.expose_secret(),
678 );
679 let mut ws_user_data = BinanceSpotWsTradingClient::new(
680 Some(url),
681 api_key.expose_secret().to_owned(),
682 api_secret.expose_secret().to_owned(),
683 Some(BINANCE_WS_HEARTBEAT_SECS),
684 self.config.transport_backend,
685 )
686 .with_proxy(
687 self.config
688 .proxy_url
689 .as_ref()
690 .map(|value| value.expose_secret().to_owned()),
691 )
692 .with_socket_control(self.socket_factory.control("binance-spot-user-streams"));
693 *self.ws_user_data_client.lock() = Some(ws_user_data.clone());
694 ws_user_data
695 .connect()
696 .await
697 .context("failed to connect Binance US user data stream")?;
698
699 let ws_clone = ws_user_data.clone();
700 let emitter = self.emitter.clone();
701 let account_id = self.core.account_id;
702 let clock = self.clock;
703 let http_client = self.http_client.clone();
704 let dispatch_state = self.dispatch_state.clone();
705 let lifecycle_lock = self.lifecycle_lock.clone();
706 let treat_expired_as_canceled = self.config.treat_expired_as_canceled;
707 let ws_authenticated = self.ws_authenticated.clone();
708 let ws_user_data_subscribed = self.ws_user_data_subscribed.clone();
709 let (setup_error_tx, _setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
710 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
711 let task_spawner = self
712 .session_tasks
713 .spawner()
714 .context("Binance Spot session task admission is closed")?;
715
716 let future = async move {
717 while let Some(message) = ws_clone.recv().await {
718 if matches!(&message, BinanceSpotWsTradingMessage::Reconnected) {
719 ws_clone.mark_user_data_active();
720 }
721 let _lifecycle_guard = lifecycle_lock.lock();
722 dispatch_ws_trading_message(
723 message,
724 &emitter,
725 &http_client,
726 account_id,
727 treat_expired_as_canceled,
728 clock,
729 &dispatch_state,
730 &ws_authenticated,
731 &ws_user_data_subscribed,
732 &setup_error_tx,
733 &seen_trade_ids,
734 &task_spawner,
735 );
736 }
737 log::warn!("Binance US user data dispatch loop ended");
738 };
739 let dispatch = self
740 .session_tasks
741 .spawn_named("binance-spot-user-data-dispatch", future)?;
742 self.ws_user_data_dispatch = Some(dispatch);
743
744 let keepalive_http = self.http_client.clone();
745 let keepalive_key = listen_key.clone();
746
747 if let Err(e) = self.session_tasks.spawn(async move {
748 let mut interval = tokio::time::interval(Duration::from_secs(30 * 60));
749 interval.tick().await;
750
751 loop {
752 interval.tick().await;
753
754 if let Err(e) = keepalive_http
755 .inner()
756 .extend_listen_key(keepalive_key.expose_secret())
757 .await
758 {
759 log::warn!("Binance US listen key keepalive failed: {e}");
760 }
761 }
762 }) {
763 return Err(e.into());
764 }
765
766 ws_user_data.mark_user_data_active();
767 *self.ws_user_data_client.lock() = Some(ws_user_data);
768 Ok(())
769 }
770
771 async fn disconnect_us_user_data(&mut self) {
772 let mut client_drained = true;
773 let client = self.ws_user_data_client.lock().clone();
774
775 if let Some(mut client) = client {
776 client.mark_user_data_inactive();
777 if let Err(e) = client.disconnect().await {
778 client_drained = false;
779 self.shutdown_errors.push(format!(
780 "failed to stop Binance US user data WebSocket: {e}"
781 ));
782 }
783 }
784
785 if client_drained {
786 *self.ws_user_data_client.lock() = None;
787 }
788
789 if let Some(listen_key) = self.listen_key.clone() {
790 match self
791 .http_client
792 .inner()
793 .close_listen_key(listen_key.expose_secret())
794 .await
795 {
796 Ok(())
797 if self.listen_key.as_ref().map(SecretString::expose_secret)
798 == Some(listen_key.expose_secret()) =>
799 {
800 self.listen_key = None;
801 }
802 Ok(()) => {}
803 Err(e) => self
804 .shutdown_errors
805 .push(format!("failed to close Binance US listen key: {e}")),
806 }
807 }
808 }
809
810 async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
811 if let Some(client) = self.ws_trading_client.as_ref() {
812 client.mark_user_data_inactive();
813 client.begin_shutdown();
814 }
815
816 if let Some(client) = self.ws_user_data_client.lock().as_ref() {
817 client.mark_user_data_inactive();
818 client.begin_shutdown();
819 }
820
821 self.abort_session_tasks();
822 self.abort_pending_tasks();
823
824 if let Some(ref mut ws_trading) = self.ws_trading_client
825 && let Err(e) = ws_trading.disconnect().await
826 {
827 self.shutdown_errors
828 .push(format!("trading WebSocket shutdown failed: {e}"));
829 }
830 self.disconnect_us_user_data().await;
831
832 let (session_result, pending_result) =
833 tokio::join!(self.await_session_tasks(), self.await_pending_tasks());
834 self.core.set_disconnected();
835
836 if let Err(e) = session_result {
837 self.shutdown_errors.push(e.to_string());
838 }
839
840 if let Err(e) = pending_result {
841 self.shutdown_errors.push(e.to_string());
842 }
843
844 if !self.shutdown_errors.is_empty() {
845 let errors = std::mem::take(&mut self.shutdown_errors);
846 anyhow::bail!("Binance Spot shutdown failed: {}", errors.join("; "));
847 }
848 Ok(())
849 }
850}
851
852#[async_trait(?Send)]
853impl ExecutionClient for BinanceSpotExecutionClient {
854 fn is_connected(&self) -> bool {
855 self.core.is_connected()
856 }
857
858 fn client_id(&self) -> ClientId {
859 self.core.client_id
860 }
861
862 fn account_id(&self) -> AccountId {
863 self.core.account_id
864 }
865
866 fn venue(&self) -> Venue {
867 *BINANCE_VENUE
868 }
869
870 fn oms_type(&self) -> OmsType {
871 self.core.oms_type
872 }
873
874 fn get_account(&self) -> Option<AccountAny> {
875 self.core.cache().account_owned(&self.core.account_id)
876 }
877
878 async fn connect(&mut self) -> anyhow::Result<()> {
879 if self.core.is_connected() && self.session_tasks.is_open() && self.pending_tasks.is_open()
880 {
881 return Ok(());
882 }
883
884 if !self.pending_tasks.is_open() || !self.session_tasks.is_open() {
885 self.teardown_partial_connect().await?;
886 }
887
888 if !self.pending_tasks.is_open() {
889 self.await_pending_tasks().await?;
890 self.pending_tasks.start_generation().map_err(|e| {
891 anyhow::anyhow!("Failed to start Binance Spot task generation: {e}")
892 })?;
893 }
894
895 if !self.session_tasks.is_open() {
896 self.await_session_tasks().await?;
897 self.session_tasks.start_generation().map_err(|e| {
898 anyhow::anyhow!("Failed to start Binance Spot session generation: {e}")
899 })?;
900 }
901 let ws_trading_client = self.ws_trading_client.clone();
902 let ws_user_data_client = Arc::clone(&self.ws_user_data_client);
903 let setup_guard =
904 TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
905 if let Some(client) = ws_trading_client {
906 client.begin_shutdown();
907 }
908
909 if let Some(client) = ws_user_data_client.lock().as_ref() {
910 client.begin_shutdown();
911 }
912 });
913
914 let ws_setup_timeout = Duration::from_millis(self.config.ws_trading_setup_timeout_ms);
915
916 if !self.core.instruments_initialized() {
918 let instruments = self
919 .http_client
920 .request_instruments_with_config(&self.config.instrument_provider, self.config.us)
921 .await
922 .context("failed to request Binance Spot instruments")?;
923
924 if instruments.is_empty() {
925 log::warn!("No instruments returned for Binance Spot");
926 } else {
927 log::debug!("Loaded {} Spot instruments", instruments.len());
928 self.http_client.cache_instruments(instruments);
929 }
930
931 self.core.set_instruments_initialized();
932 }
933
934 let account_state = self
936 .refresh_account_state()
937 .await
938 .context("failed to request Binance account state")?;
939
940 if !account_state.balances.is_empty() {
941 log::debug!(
942 "Received account state with {} balance(s)",
943 account_state.balances.len()
944 );
945 }
946
947 self.emitter.send_account_state(account_state);
948
949 crate::common::execution::await_account_registered(&self.core, self.core.account_id, 30.0)
951 .await?;
952
953 let session_result = async {
954 if self.config.us {
955 self.connect_us_user_data().await?;
956 }
957
958 if let Some(mut ws_trading) = self.ws_trading_client.clone() {
959 match ws_trading.connect().await {
960 Ok(()) => {
961 log::debug!("Connected to Binance Spot WS trading API");
962
963 let ws_trading_clone = ws_trading.clone();
964 let emitter = self.emitter.clone();
965 let account_id = self.core.account_id;
966 let clock = self.clock;
967 let http_client = self.http_client.clone();
968 let dispatch_state = self.dispatch_state.clone();
969 let lifecycle_lock = self.lifecycle_lock.clone();
970 let treat_expired_as_canceled = self.config.treat_expired_as_canceled;
971 let ws_authenticated = self.ws_authenticated.clone();
972 let ws_user_data_subscribed = self.ws_user_data_subscribed.clone();
973 let (ws_setup_error_tx, mut ws_setup_error_rx) =
974 tokio::sync::mpsc::unbounded_channel();
975 let seen_trade_ids = std::sync::Arc::new(Mutex::new(FifoCache::new()));
976 let task_spawner = self
977 .session_tasks
978 .spawner()
979 .context("Binance Spot session task admission is closed")?;
980
981 let future = async move {
982 let mut resubscribing = false;
983
984 loop {
985 match ws_trading_clone.recv().await {
986 Some(msg) => {
987 match &msg {
988 BinanceSpotWsTradingMessage::Reconnected => {
989 ws_trading_clone.mark_user_data_inactive();
990 resubscribing = true;
991 if let Err(e) = ws_trading_clone.session_logon().await {
992 resubscribing = false;
993 log::error!(
994 "Failed to re-authenticate Binance Spot user data stream: {e}"
995 );
996 }
997 }
998 BinanceSpotWsTradingMessage::Authenticated if resubscribing => {
999 if let Err(e) =
1000 ws_trading_clone.subscribe_user_data().await
1001 {
1002 resubscribing = false;
1003 log::error!(
1004 "Failed to resubscribe Binance Spot user data stream: {e}"
1005 );
1006 }
1007 continue;
1008 }
1009 BinanceSpotWsTradingMessage::UserDataSubscribed { .. } => {
1010 let was_resubscribing = resubscribing;
1011 resubscribing = false;
1012 ws_trading_clone.mark_user_data_active();
1013
1014 if was_resubscribing {
1015 continue;
1016 }
1017 }
1018 BinanceSpotWsTradingMessage::AuthenticationRejected(reason)
1019 if resubscribing =>
1020 {
1021 log::warn!(
1022 "Binance Spot reconnect authentication failed; retrying: {reason}"
1023 );
1024 tokio::time::sleep(WS_RECONNECT_SETUP_RETRY_DELAY).await;
1025
1026 if let Err(e) = ws_trading_clone.session_logon().await {
1027 resubscribing = false;
1028 log::error!(
1029 "Failed to retry Binance Spot reconnect authentication: {e}"
1030 );
1031 }
1032 continue;
1033 }
1034 BinanceSpotWsTradingMessage::UserDataSubscriptionRejected(reason)
1035 if resubscribing =>
1036 {
1037 log::warn!(
1038 "Binance Spot reconnect user data subscription failed; retrying: {reason}"
1039 );
1040 tokio::time::sleep(WS_RECONNECT_SETUP_RETRY_DELAY).await;
1041
1042 if let Err(e) =
1043 ws_trading_clone.subscribe_user_data().await
1044 {
1045 resubscribing = false;
1046 log::error!(
1047 "Failed to retry Binance Spot user data subscription: {e}"
1048 );
1049 }
1050 continue;
1051 }
1052 _ => {}
1053 }
1054
1055 let _lifecycle_guard = lifecycle_lock.lock();
1056 dispatch_ws_trading_message(
1057 msg,
1058 &emitter,
1059 &http_client,
1060 account_id,
1061 treat_expired_as_canceled,
1062 clock,
1063 &dispatch_state,
1064 &ws_authenticated,
1065 &ws_user_data_subscribed,
1066 &ws_setup_error_tx,
1067 &seen_trade_ids,
1068 &task_spawner,
1069 );
1070 }
1071 None => {
1072 log::warn!("WS trading dispatch loop ended");
1073 break;
1074 }
1075 }
1076 }
1077 };
1078 let dispatch = self
1079 .session_tasks
1080 .spawn_named("binance-spot-trading-dispatch", future)?;
1081 self.ws_trading_dispatch = Some(dispatch);
1082
1083 if let Err(e) = ws_trading.session_logon().await {
1084 let reason = format!("WS session logon failed: {e}");
1085 return Err(self.ws_setup_failure(ws_trading, reason).await);
1086 } else {
1087 let auth_result = wait_for_ws_setup_response(
1088 ws_setup_timeout,
1089 self.ws_authenticated.notified(),
1090 &mut ws_setup_error_rx,
1091 "WS session authentication timed out",
1092 )
1093 .await;
1094
1095 if let Err(e) = auth_result {
1096 return Err(self.ws_setup_failure(ws_trading, e.to_string()).await);
1097 } else if let Err(e) = ws_trading.subscribe_user_data().await {
1098 let reason = format!("WS user data subscribe failed: {e}");
1099 return Err(self.ws_setup_failure(ws_trading, reason).await);
1100 } else {
1101 let subscribe_result = wait_for_ws_setup_response(
1102 ws_setup_timeout,
1103 self.ws_user_data_subscribed.notified(),
1104 &mut ws_setup_error_rx,
1105 "WS user data subscription timed out",
1106 )
1107 .await;
1108
1109 if let Err(e) = subscribe_result {
1110 return Err(self.ws_setup_failure(ws_trading, e.to_string()).await);
1111 } else {
1112 self.ws_trading_client = Some(ws_trading);
1113 }
1114 }
1115 }
1116 }
1117 Err(e) => {
1118 let reason = format!("Failed to connect WS trading API: {e}");
1119 return Err(self.ws_setup_failure(ws_trading, reason).await);
1120 }
1121 }
1122 }
1123
1124 let refresh_secs = self.config.instrument_refresh_interval_secs;
1125 if refresh_secs > 0 {
1126 let http_client = self.http_client.clone();
1127 let provider = self.config.instrument_provider.clone();
1128 let us = self.config.us;
1129
1130 self.session_tasks.spawn(async move {
1131 let mut interval = tokio::time::interval(Duration::from_secs(refresh_secs));
1132 interval.tick().await;
1133
1134 loop {
1135 interval.tick().await;
1136
1137 match http_client
1138 .request_instruments_with_config(&provider, us)
1139 .await
1140 {
1141 Ok(instruments) => log::debug!(
1142 "Refreshed Binance Spot execution instruments: count={}",
1143 instruments.len()
1144 ),
1145 Err(e) => {
1146 log::warn!("Binance Spot execution instrument refresh failed: {e}");
1147 }
1148 }
1149 }
1150 })?;
1151 }
1152
1153 Ok::<(), anyhow::Error>(())
1154 }
1155 .await;
1156
1157 if let Err(e) = session_result {
1158 if let Err(teardown_error) = self.teardown_partial_connect().await {
1159 return Err(e.context(format!(
1160 "Binance Spot startup teardown failed: {teardown_error}"
1161 )));
1162 }
1163 return Err(e);
1164 }
1165
1166 setup_guard.disarm();
1167 self.core.set_connected();
1168 log::info!("Connected: client_id={}", self.core.client_id);
1169 Ok(())
1170 }
1171
1172 async fn disconnect(&mut self) -> anyhow::Result<()> {
1173 self.teardown_partial_connect().await?;
1174 log::info!("Disconnected: client_id={}", self.core.client_id);
1175 Ok(())
1176 }
1177
1178 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
1179 self.update_account_state();
1180 Ok(())
1181 }
1182
1183 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
1184 log::debug!("query_order: client_order_id={}", cmd.client_order_id);
1185
1186 let http_client = self.http_client.clone();
1187 let command = cmd;
1188 let event_emitter = self.emitter.clone();
1189 let account_id = self.core.account_id;
1190 let treat_expired_as_canceled = self.config.treat_expired_as_canceled;
1191
1192 self.spawn_task("query_order", async move {
1193 let result = http_client
1194 .request_order_status_report(
1195 account_id,
1196 command.instrument_id,
1197 command.venue_order_id,
1198 Some(command.client_order_id),
1199 )
1200 .await;
1201
1202 match result {
1203 Ok(Some(mut report)) => {
1204 normalize_spot_order_status_report(&mut report, treat_expired_as_canceled);
1205 event_emitter.send_order_status_report(report);
1206 }
1207 Ok(None) => log::debug!(
1208 "No order status report returned: client_order_id={}",
1209 command.client_order_id
1210 ),
1211 Err(e) => log::warn!("Failed to query order status: {e}"),
1212 }
1213
1214 Ok(())
1215 });
1216
1217 Ok(())
1218 }
1219
1220 fn generate_account_state(
1221 &self,
1222 balances: Vec<AccountBalance>,
1223 margins: Vec<MarginBalance>,
1224 reported: bool,
1225 ts_event: UnixNanos,
1226 info: Option<Params>,
1227 ) -> anyhow::Result<()> {
1228 self.emitter
1229 .emit_account_state(balances, margins, reported, ts_event, info);
1230 Ok(())
1231 }
1232
1233 fn start(&mut self) -> anyhow::Result<()> {
1234 if self.core.is_started() {
1235 return Ok(());
1236 }
1237
1238 self.emitter.set_sender(get_exec_event_sender());
1239 self.core.set_started();
1240
1241 let http_client = self.http_client.clone();
1243 let provider = self.config.instrument_provider.clone();
1244 let us = self.config.us;
1245
1246 self.session_tasks.spawn(async move {
1247 match http_client
1248 .request_instruments_with_config(&provider, us)
1249 .await
1250 {
1251 Ok(instruments) => {
1252 if instruments.is_empty() {
1253 log::warn!("No instruments returned for Binance Spot");
1254 } else {
1255 http_client.cache_instruments(instruments);
1256 log::debug!("Instruments initialized");
1257 }
1258 }
1259 Err(e) => {
1260 log::error!("Failed to request Binance Spot instruments: {e}");
1261 }
1262 }
1263 })?;
1264
1265 log::info!(
1266 "Started: client_id={}, account_id={}, account_type={:?}, environment={:?}, product_type={:?}",
1267 self.core.client_id,
1268 self.core.account_id,
1269 self.core.account_type,
1270 self.config.environment,
1271 self.config.product_type,
1272 );
1273 Ok(())
1274 }
1275
1276 fn stop(&mut self) -> anyhow::Result<()> {
1277 let was_started = self.core.is_started();
1278 self.core.set_stopped();
1279 self.begin_generation_shutdown();
1280
1281 if was_started {
1282 log::info!("Stopped: client_id={}", self.core.client_id);
1283 }
1284 Ok(())
1285 }
1286
1287 fn reset(&mut self) -> anyhow::Result<()> {
1288 self.begin_generation_shutdown();
1289 Ok(())
1290 }
1291
1292 fn dispose(&mut self) -> anyhow::Result<()> {
1293 self.begin_generation_shutdown();
1294 Ok(())
1295 }
1296
1297 async fn generate_order_status_report(
1298 &self,
1299 cmd: &GenerateOrderStatusReport,
1300 ) -> anyhow::Result<Option<OrderStatusReport>> {
1301 let Some(instrument_id) = cmd.instrument_id else {
1302 log::warn!("generate_order_status_report requires instrument_id: {cmd}");
1303 return Ok(None);
1304 };
1305
1306 anyhow::ensure!(
1307 !self.config.instrument_provider.excludes(instrument_id),
1308 "Cannot query Binance Spot order for excluded instrument {instrument_id}"
1309 );
1310
1311 let venue_order_id = cmd
1313 .venue_order_id
1314 .as_ref()
1315 .map(|id| VenueOrderId::new(id.inner()));
1316
1317 let report = self
1318 .http_client
1319 .request_order_status_report(
1320 self.core.account_id,
1321 instrument_id,
1322 venue_order_id,
1323 cmd.client_order_id,
1324 )
1325 .await?;
1326
1327 Ok(report.map(|mut report| {
1328 normalize_spot_order_status_report(&mut report, self.config.treat_expired_as_canceled);
1329 report
1330 }))
1331 }
1332
1333 async fn generate_order_status_reports(
1334 &self,
1335 cmd: &GenerateOrderStatusReports,
1336 ) -> anyhow::Result<Vec<OrderStatusReport>> {
1337 let start_dt = cmd.start.map(|nanos| nanos.to_datetime_utc());
1338 let end_dt = cmd.end.map(|nanos| nanos.to_datetime_utc());
1339
1340 let mut reports = self
1341 .http_client
1342 .request_order_status_reports_scoped(
1343 self.core.account_id,
1344 cmd.instrument_id,
1345 start_dt,
1346 end_dt,
1347 cmd.open_only,
1348 None, Some(&self.config.instrument_provider),
1350 )
1351 .await?;
1352
1353 normalize_spot_order_status_reports(&mut reports, self.config.treat_expired_as_canceled);
1354
1355 crate::common::execution::log_report_receipt(
1356 reports.len(),
1357 "OrderStatusReport",
1358 cmd.log_receipt_level,
1359 );
1360 Ok(reports)
1361 }
1362
1363 async fn generate_fill_reports(
1364 &self,
1365 cmd: GenerateFillReports,
1366 ) -> anyhow::Result<Vec<FillReport>> {
1367 let Some(instrument_id) = cmd.instrument_id else {
1368 log::warn!("generate_fill_reports requires instrument_id for Binance Spot");
1369 return Ok(Vec::new());
1370 };
1371
1372 if self.config.instrument_provider.excludes(instrument_id) {
1373 log::debug!("Dropping out-of-scope Binance Spot report request for {instrument_id}");
1374 return Ok(Vec::new());
1375 }
1376
1377 let venue_order_id = cmd
1379 .venue_order_id
1380 .as_ref()
1381 .map(|id| VenueOrderId::new(id.inner()));
1382 let requested_start_time = cmd
1383 .start
1384 .map(|start| start.as_i64() / NANOSECONDS_IN_MILLISECOND as i64);
1385 let requested_end_time = cmd
1386 .end
1387 .map(|end| end.as_i64() / NANOSECONDS_IN_MILLISECOND as i64);
1388 if let (Some(start), Some(end)) = (requested_start_time, requested_end_time) {
1389 anyhow::ensure!(
1390 start <= end,
1391 "fill report start time must not exceed end time"
1392 );
1393 }
1394
1395 let mut reports = Vec::new();
1396 let mut seen_trade_ids = AHashSet::new();
1397
1398 if venue_order_id.is_some() {
1399 let mut from_id = 0;
1400
1401 loop {
1402 let page = self
1403 .http_client
1404 .request_fill_reports_with_cursor(
1405 self.core.account_id,
1406 instrument_id,
1407 venue_order_id,
1408 None,
1409 None,
1410 Some(from_id),
1411 Some(ACCOUNT_TRADES_PAGE_LIMIT),
1412 )
1413 .await?;
1414
1415 if page.is_empty() {
1416 break;
1417 }
1418
1419 let page_len = page.len();
1420 let max_trade_id = max_trade_id(&page)?;
1421 let passed_end = requested_end_time.is_some_and(|end_time| {
1422 page.iter().any(|report| report_time_ms(report) > end_time)
1423 });
1424
1425 reports.extend(page.into_iter().filter(|report| {
1426 requested_start_time
1427 .is_none_or(|start_time| report_time_ms(report) >= start_time)
1428 && requested_end_time
1429 .is_none_or(|end_time| report_time_ms(report) <= end_time)
1430 && seen_trade_ids.insert(report.trade_id)
1431 }));
1432
1433 if page_len < ACCOUNT_TRADES_PAGE_LIMIT as usize || passed_end {
1434 break;
1435 }
1436
1437 let next_from_id = max_trade_id
1438 .checked_add(1)
1439 .context("Binance Spot trade ID overflow during pagination")?;
1440 anyhow::ensure!(
1441 next_from_id > from_id,
1442 "Binance Spot account-trades pagination made no progress"
1443 );
1444 from_id = next_from_id;
1445 }
1446 } else if let Some(query_start_time) = requested_start_time {
1447 let query_end_time = requested_end_time.unwrap_or_else(|| {
1448 self.clock.get_time_ns().as_i64() / NANOSECONDS_IN_MILLISECOND as i64
1449 });
1450 anyhow::ensure!(
1451 query_start_time <= query_end_time,
1452 "fill report start time must not exceed end time"
1453 );
1454 let mut window_start = query_start_time;
1455
1456 loop {
1457 let window_end = window_start
1458 .saturating_add(ACCOUNT_TRADES_MAX_INTERVAL_MS)
1459 .min(query_end_time);
1460 let mut from_id = None;
1461
1462 loop {
1463 let start = if from_id.is_none() {
1464 Some(
1465 Timestamp::from_millisecond(window_start)
1466 .context("invalid Binance Spot account-trades start time")?,
1467 )
1468 } else {
1469 None
1470 };
1471 let end = if from_id.is_none() {
1472 Some(
1473 Timestamp::from_millisecond(window_end)
1474 .context("invalid Binance Spot account-trades end time")?,
1475 )
1476 } else {
1477 None
1478 };
1479 let page = self
1480 .http_client
1481 .request_fill_reports_with_cursor(
1482 self.core.account_id,
1483 instrument_id,
1484 None,
1485 start,
1486 end,
1487 from_id,
1488 Some(ACCOUNT_TRADES_PAGE_LIMIT),
1489 )
1490 .await?;
1491
1492 if page.is_empty() {
1493 break;
1494 }
1495
1496 let page_len = page.len();
1497 let max_trade_id = max_trade_id(&page)?;
1498 let passed_window_end = page
1499 .iter()
1500 .any(|report| report_time_ms(report) > window_end);
1501
1502 reports.extend(page.into_iter().filter(|report| {
1503 let report_time = report_time_ms(report);
1504 report_time >= window_start
1505 && report_time <= window_end
1506 && seen_trade_ids.insert(report.trade_id)
1507 }));
1508
1509 if page_len < ACCOUNT_TRADES_PAGE_LIMIT as usize || passed_window_end {
1510 break;
1511 }
1512
1513 let next_from_id = max_trade_id
1514 .checked_add(1)
1515 .context("Binance Spot trade ID overflow during pagination")?;
1516 anyhow::ensure!(
1517 from_id.is_none_or(|cursor| next_from_id > cursor),
1518 "Binance Spot account-trades pagination made no progress"
1519 );
1520 from_id = Some(next_from_id);
1521 }
1522
1523 if window_end >= query_end_time {
1524 break;
1525 }
1526 window_start = window_end.saturating_add(1);
1527 }
1528 } else {
1529 let mut from_id = 0;
1530
1531 loop {
1532 let page = self
1533 .http_client
1534 .request_fill_reports_with_cursor(
1535 self.core.account_id,
1536 instrument_id,
1537 None,
1538 None,
1539 None,
1540 Some(from_id),
1541 Some(ACCOUNT_TRADES_PAGE_LIMIT),
1542 )
1543 .await?;
1544
1545 if page.is_empty() {
1546 break;
1547 }
1548
1549 let page_len = page.len();
1550 let max_trade_id = max_trade_id(&page)?;
1551 let passed_end = requested_end_time.is_some_and(|end_time| {
1552 page.iter().any(|report| report_time_ms(report) > end_time)
1553 });
1554
1555 reports.extend(page.into_iter().filter(|report| {
1556 requested_end_time.is_none_or(|end_time| report_time_ms(report) <= end_time)
1557 && seen_trade_ids.insert(report.trade_id)
1558 }));
1559
1560 if page_len < ACCOUNT_TRADES_PAGE_LIMIT as usize || passed_end {
1561 break;
1562 }
1563
1564 let next_from_id = max_trade_id
1565 .checked_add(1)
1566 .context("Binance Spot trade ID overflow during pagination")?;
1567 anyhow::ensure!(
1568 next_from_id > from_id,
1569 "Binance Spot account-trades pagination made no progress"
1570 );
1571 from_id = next_from_id;
1572 }
1573 }
1574
1575 let mut reports_with_trade_ids = reports
1576 .into_iter()
1577 .map(|report| parse_trade_id(&report).map(|trade_id| (report, trade_id)))
1578 .collect::<anyhow::Result<Vec<_>>>()?;
1579 reports_with_trade_ids
1580 .sort_unstable_by_key(|(report, trade_id)| (report.ts_event, *trade_id));
1581 crate::common::execution::log_report_receipt(
1582 reports_with_trade_ids.len(),
1583 "FillReport",
1584 cmd.log_receipt_level,
1585 );
1586 Ok(reports_with_trade_ids
1587 .into_iter()
1588 .map(|(report, _)| report)
1589 .collect())
1590 }
1591
1592 async fn generate_position_status_reports(
1593 &self,
1594 cmd: &GeneratePositionStatusReports,
1595 ) -> anyhow::Result<Vec<PositionStatusReport>> {
1596 crate::common::execution::log_report_receipt(
1599 0,
1600 "PositionStatusReport",
1601 cmd.log_receipt_level,
1602 );
1603 Ok(Vec::new())
1604 }
1605
1606 async fn generate_mass_status(
1607 &self,
1608 lookback_mins: Option<u64>,
1609 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
1610 log::info!("Generating ExecutionMassStatus (lookback_mins={lookback_mins:?})");
1611
1612 let ts_now = self.clock.get_time_ns();
1613
1614 let start = lookback_mins
1615 .map(DurationNanos::try_from_mins)
1616 .transpose()?
1617 .map(|lookback| ts_now.saturating_sub(lookback));
1618
1619 let order_cmd = GenerateOrderStatusReportsBuilder::default()
1622 .log_receipt_level(LogLevel::Off)
1623 .ts_init(ts_now)
1624 .open_only(true)
1625 .start(start)
1626 .build()
1627 .map_err(|e| anyhow::anyhow!("{e}"))?;
1628
1629 let position_cmd = GeneratePositionStatusReportsBuilder::default()
1630 .log_receipt_level(LogLevel::Off)
1631 .ts_init(ts_now)
1632 .start(start)
1633 .build()
1634 .map_err(|e| anyhow::anyhow!("{e}"))?;
1635
1636 let (order_reports, position_reports) = tokio::try_join!(
1637 self.generate_order_status_reports(&order_cmd),
1638 self.generate_position_status_reports(&position_cmd),
1639 )?;
1640
1641 let mut instrument_ids: Vec<_> = order_reports
1642 .iter()
1643 .map(|report| report.instrument_id)
1644 .collect();
1645 {
1646 let cache = self.core.cache();
1647 instrument_ids.extend(
1648 cache
1649 .orders_open(
1650 Some(&BINANCE_VENUE),
1651 None,
1652 None,
1653 Some(&self.core.account_id),
1654 None,
1655 )
1656 .into_iter()
1657 .chain(cache.orders_inflight(
1658 Some(&BINANCE_VENUE),
1659 None,
1660 None,
1661 Some(&self.core.account_id),
1662 None,
1663 ))
1664 .map(|order| order.instrument_id())
1665 .filter(|instrument_id| {
1666 self.http_client
1667 .get_instrument(&instrument_id.symbol.inner())
1668 .is_some_and(|instrument| instrument.id() == *instrument_id)
1669 }),
1670 );
1671 }
1672 instrument_ids.sort_unstable();
1673 instrument_ids.dedup();
1674
1675 let mut fill_reports = Vec::new();
1676
1677 for instrument_id in instrument_ids {
1678 let fill_cmd = GenerateFillReportsBuilder::default()
1679 .log_receipt_level(LogLevel::Off)
1680 .ts_init(ts_now)
1681 .instrument_id(Some(instrument_id))
1682 .start(start)
1683 .end(start.map(|_| ts_now))
1684 .build()
1685 .map_err(|e| anyhow::anyhow!("{e}"))?;
1686 fill_reports.extend(self.generate_fill_reports(fill_cmd).await?);
1687 }
1688
1689 log::info!("Received {} OrderStatusReports", order_reports.len());
1690 log::info!("Received {} FillReports", fill_reports.len());
1691 log::info!("Received {} PositionReports", position_reports.len());
1692
1693 let mut mass_status = ExecutionMassStatus::new(
1694 self.core.client_id,
1695 self.core.account_id,
1696 *BINANCE_VENUE,
1697 ts_now,
1698 None,
1699 );
1700
1701 let reported_order_ids: AHashSet<_> = order_reports
1702 .iter()
1703 .map(|report| report.venue_order_id)
1704 .collect();
1705 let cache = self.core.cache();
1706 let reports_complete = fill_reports.iter().all(|fill| {
1707 reported_order_ids.contains(&fill.venue_order_id)
1708 || cache
1709 .client_order_id(&fill.venue_order_id)
1710 .and_then(|client_order_id| cache.order(client_order_id))
1711 .is_some_and(|order| {
1712 order.instrument_id() == fill.instrument_id
1713 && order.account_id() == Some(fill.account_id)
1714 && order.order_side() == fill.order_side
1715 })
1716 });
1717 mass_status.set_report_window(start, reports_complete);
1718 mass_status.add_order_reports(order_reports);
1719 mass_status.add_fill_reports(fill_reports);
1720 mass_status.add_position_reports(position_reports);
1721
1722 Ok(Some(mass_status))
1723 }
1724
1725 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
1726 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
1727
1728 if order.is_closed() {
1729 let client_order_id = order.client_order_id();
1730 log::warn!("Cannot submit closed order {client_order_id}");
1731 return Ok(());
1732 }
1733
1734 if let Err(reason) = validate_order(&order, self.config.use_gtd) {
1735 self.emitter.emit_order_denied(&order, &reason.to_string());
1736 return Ok(());
1737 }
1738
1739 self.ensure_ws_user_data_active()?;
1740
1741 let params = if self.ws_order_transport_active() {
1742 match build_new_order_params(
1743 &order,
1744 order.client_order_id(),
1745 order.is_post_only(),
1746 order.is_quote_quantity(),
1747 self.config.use_gtd,
1748 ) {
1749 Ok(params) => Some(params),
1750 Err(e) => {
1751 let reason = OrderDeniedReason::ValidationFailed {
1752 detail: e.to_string(),
1753 };
1754 self.emitter.emit_order_denied(&order, &reason.to_string());
1755 return Ok(());
1756 }
1757 }
1758 } else {
1759 None
1760 };
1761
1762 log::debug!("OrderSubmitted client_order_id={}", order.client_order_id());
1763 self.emitter.emit_order_submitted(&order);
1764
1765 self.submit_order_internal(&cmd, params)
1766 }
1767
1768 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
1769 if cmd.order_list.client_order_ids.is_empty() {
1770 log::debug!("submit_order_list called with empty order list");
1771 return Ok(());
1772 }
1773
1774 self.ensure_ws_user_data_active()?;
1775
1776 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
1777
1778 if let Some(reason) = orders
1779 .iter()
1780 .find_map(|order| validate_order(order, self.config.use_gtd).err())
1781 {
1782 let reason = reason.to_string();
1783
1784 for order in &orders {
1785 self.emitter.emit_order_denied(order, &reason);
1786 }
1787 return Ok(());
1788 }
1789
1790 if let Some(order) = orders.iter().find(|order| order.is_closed()) {
1791 let reason = format!("Cannot submit closed order {}", order.client_order_id());
1792 for order in &orders {
1793 self.emitter.emit_order_denied(order, &reason);
1794 }
1795 return Ok(());
1796 }
1797
1798 let params = match build_spot_order_list_params(
1799 cmd.order_list.id.as_ref(),
1800 &orders,
1801 self.config.use_gtd,
1802 ) {
1803 Ok(request) => request,
1804 Err(reason) => {
1805 for order in &orders {
1806 self.emitter.emit_order_denied(order, &reason);
1807 }
1808 return Ok(());
1809 }
1810 };
1811
1812 for order in &orders {
1813 self.dispatch_state.order_identities.insert(
1814 order.client_order_id(),
1815 OrderIdentity {
1816 instrument_id: order.instrument_id(),
1817 strategy_id: order.strategy_id(),
1818 order_side: order.order_side(),
1819 order_type: order.order_type(),
1820 price: order.price(),
1821 quantity: order.quantity(),
1822 venue_position_id: None,
1823 },
1824 );
1825 self.emitter.emit_order_submitted(order);
1826 }
1827
1828 let event_emitter = self.emitter.clone();
1829 let trader_id = self.core.trader_id;
1830 let account_id = self.core.account_id;
1831 let clock = self.clock;
1832 let http_client = self.http_client.clone();
1833 let dispatch_state = self.dispatch_state.clone();
1834
1835 self.spawn_task("submit_order_list_http", async move {
1836 if let Err(e) = submit_spot_order_list(&http_client, ¶ms).await {
1837 handle_spot_order_list_submit_error(
1838 &event_emitter,
1839 &dispatch_state,
1840 trader_id,
1841 account_id,
1842 clock,
1843 &orders,
1844 e,
1845 )?;
1846 }
1847
1848 Ok(())
1849 });
1850
1851 Ok(())
1852 }
1853
1854 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1855 let order = self
1859 .core
1860 .cache()
1861 .order(&cmd.client_order_id)
1862 .map(|o| o.clone());
1863
1864 let Some(order) = order else {
1865 log::warn!(
1866 "Cannot modify order {}: not found in cache",
1867 cmd.client_order_id
1868 );
1869 let ts_init = self.clock.get_time_ns();
1870 let rejected_event = OrderModifyRejected::new(
1871 self.core.trader_id,
1872 cmd.strategy_id,
1873 cmd.instrument_id,
1874 cmd.client_order_id,
1875 "Order not found in cache for modify".into(),
1876 UUID4::new(),
1877 ts_init, ts_init,
1879 false,
1880 cmd.venue_order_id,
1881 Some(self.core.account_id),
1882 );
1883
1884 self.emitter
1885 .send_order_event(OrderEventAny::ModifyRejected(rejected_event));
1886 return Ok(());
1887 };
1888
1889 let event_emitter = self.emitter.clone();
1890 let trader_id = self.core.trader_id;
1891 let account_id = self.core.account_id;
1892 let clock = self.clock;
1893
1894 let order_side = order.order_side();
1895 let order_type = order.order_type();
1896 let time_in_force = order.time_in_force();
1897 let quantity = cmd.quantity.unwrap_or_else(|| order.quantity());
1898 let use_gtd = self.config.use_gtd;
1899 let dispatch_state = self.dispatch_state.clone();
1900
1901 if self.ws_order_transport_active() {
1902 let command = cmd;
1903 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
1904 let cancel_id = cancel_replace_cancel_id();
1905 dispatch_state.insert_cancel_replace(cancel_id.clone());
1906 let params = build_cancel_replace_params(
1907 &command,
1908 &order,
1909 quantity,
1910 use_gtd,
1911 cancel_id.clone(),
1912 )?;
1913
1914 if let Some(venue_order_id) = command.venue_order_id {
1915 dispatch_state.begin_replace(command.client_order_id, venue_order_id);
1916 }
1917
1918 let request_id = ws_client.next_request_id();
1920 dispatch_state.pending_requests.insert(
1921 request_id.clone(),
1922 PendingRequest {
1923 client_order_id: command.client_order_id,
1924 venue_order_id: command.venue_order_id,
1925 operation: PendingOperation::Modify,
1926 },
1927 );
1928 dispatch_state
1929 .cancel_replace_request_ids
1930 .insert(request_id.clone(), cancel_id);
1931
1932 self.spawn_task("modify_order_ws", async move {
1933 if let Err(e) = ws_client
1934 .cancel_replace_order_with_id(request_id.clone(), params)
1935 .await
1936 {
1937 dispatch_state.pending_requests.remove(&request_id);
1938 dispatch_state
1939 .cancel_replace_request_ids
1940 .remove(&request_id);
1941 log::warn!(
1942 "WS modify request failed for {}, awaiting reconciliation: {e}",
1943 command.client_order_id
1944 );
1945 anyhow::bail!("WS modify order failed: {e}");
1946 }
1947 Ok(())
1948 });
1949 } else {
1950 let command = cmd;
1951 let http_client = self.http_client.clone();
1952 log::debug!("WS trading not active, falling back to HTTP for modify_order");
1953 let dispatch_state = self.dispatch_state.clone();
1954 let cancel_id = cancel_replace_cancel_id();
1955 dispatch_state.insert_cancel_replace(cancel_id.clone());
1956
1957 if let Some(venue_order_id) = command.venue_order_id {
1958 dispatch_state.begin_replace(command.client_order_id, venue_order_id);
1959 }
1960
1961 self.spawn_task("modify_order_http", async move {
1962 let result = match command.venue_order_id {
1963 Some(venue_order_id) => {
1964 http_client
1965 .modify_order(
1966 account_id,
1967 command.instrument_id,
1968 venue_order_id,
1969 command.client_order_id,
1970 order_side,
1971 order_type,
1972 quantity,
1973 time_in_force,
1974 command.price,
1975 use_gtd,
1976 &cancel_id,
1977 )
1978 .await
1979 }
1980 None => Err(anyhow::anyhow!(BinanceSpotHttpError::ValidationError(
1981 "venue_order_id required for modify".to_string()
1982 ))),
1983 };
1984
1985 match result {
1986 Ok(report) => {
1987 let Some(price) = report.price else {
1988 anyhow::bail!("Spot replacement response has no price");
1989 };
1990
1991 if !dispatch_state.record_order_update(
1992 command.client_order_id,
1993 report.venue_order_id,
1994 report.quantity,
1995 price,
1996 report.trigger_price,
1997 ) {
1998 return Ok(());
1999 }
2000 let ts_now = clock.get_time_ns();
2001 let updated_event = OrderUpdated::new(
2002 trader_id,
2003 command.strategy_id,
2004 command.instrument_id,
2005 command.client_order_id,
2006 report.quantity,
2007 UUID4::new(),
2008 ts_now,
2009 ts_now,
2010 false,
2011 Some(report.venue_order_id),
2012 Some(account_id),
2013 report.price,
2014 None, None, false, );
2018 event_emitter.send_order_event(OrderEventAny::Updated(updated_event));
2019 }
2020 Err(e) => {
2021 handle_http_modify_failure(
2022 &e,
2023 &command,
2024 &cancel_id,
2025 &event_emitter,
2026 &dispatch_state,
2027 account_id,
2028 clock.get_time_ns(),
2029 );
2030 return Err(e);
2031 }
2032 }
2033 Ok(())
2034 });
2035 }
2036
2037 Ok(())
2038 }
2039
2040 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
2041 self.cancel_order_internal(&cmd);
2042 Ok(())
2043 }
2044
2045 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
2046 if cmd.order_side.is_some() {
2047 let cancels: Vec<CancelOrder> = {
2049 let cache = self.core.cache();
2050 cache
2051 .orders_open(None, Some(&cmd.instrument_id), None, None, cmd.order_side)
2052 .into_iter()
2053 .map(|order| CancelOrder {
2054 trader_id: order.trader_id(),
2055 client_id: cmd.client_id,
2056 strategy_id: order.strategy_id(),
2057 instrument_id: order.instrument_id(),
2058 client_order_id: order.client_order_id(),
2059 venue_order_id: order.venue_order_id(),
2060 command_id: cmd.command_id,
2061 ts_init: cmd.ts_init,
2062 params: cmd.params.clone(),
2063 correlation_id: cmd.correlation_id,
2064 causation_id: cmd.causation_id,
2065 })
2066 .collect()
2067 };
2068
2069 if cancels.is_empty() {
2070 log::debug!("No open orders to cancel for {}", cmd.instrument_id);
2071 return Ok(());
2072 }
2073
2074 return self.batch_cancel_orders(BatchCancelOrders {
2075 trader_id: cmd.trader_id,
2076 client_id: cmd.client_id,
2077 strategy_id: cmd.strategy_id,
2078 instrument_id: cmd.instrument_id,
2079 cancels,
2080 command_id: cmd.command_id,
2081 ts_init: cmd.ts_init,
2082 params: cmd.params,
2083 correlation_id: cmd.correlation_id,
2084 causation_id: cmd.causation_id,
2085 });
2086 }
2087
2088 let event_emitter = self.emitter.clone();
2089 let trader_id = self.core.trader_id;
2090 let account_id = self.core.account_id;
2091 let clock = self.clock;
2092
2093 if self.ws_order_transport_active() {
2094 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
2095 let symbol = cmd.instrument_id.symbol.to_string();
2096
2097 self.spawn_task("cancel_all_orders_ws", async move {
2098 if let Err(e) = ws_client.cancel_all_orders(symbol).await {
2099 log::error!("WS cancel_all_orders failed: {e}");
2100 }
2101 Ok(())
2103 });
2104
2105 return Ok(());
2106 }
2107
2108 log::debug!("WS trading not active, falling back to HTTP for cancel_all_orders");
2109 let http_client = self.http_client.clone();
2110 let dispatch_state = self.dispatch_state.clone();
2111 let lifecycle_lock = self.lifecycle_lock.clone();
2112
2113 let strategy_lookup: AHashMap<ClientOrderId, StrategyId> = {
2115 let cache = self.core.cache();
2116 cache
2117 .orders_open(None, Some(&cmd.instrument_id), None, None, None)
2118 .into_iter()
2119 .map(|order| (order.client_order_id(), order.strategy_id()))
2120 .collect()
2121 };
2122
2123 let command = cmd;
2124 self.spawn_task("cancel_all_orders_http", async move {
2125 let responses = http_client
2126 .cancel_all_order_responses(command.instrument_id)
2127 .await?;
2128 let canceled_orders = prepare_cancel_all_orders(responses)?;
2129 anyhow::ensure!(
2130 canceled_orders
2131 .iter()
2132 .all(|order| order.instrument_id == command.instrument_id),
2133 "cancel-all response contains an order for a different instrument",
2134 );
2135 let _lifecycle_guard = lifecycle_lock.lock();
2136
2137 for canceled_order in canceled_orders {
2138 let client_order_id = canceled_order.client_order_id;
2139 if canceled_order.order_list {
2140 dispatch_order_list_canceled(
2141 &canceled_order,
2142 &event_emitter,
2143 account_id,
2144 &dispatch_state,
2145 clock.get_time_ns(),
2146 );
2147 continue;
2148 }
2149 let strategy_id = strategy_lookup
2150 .get(&client_order_id)
2151 .copied()
2152 .unwrap_or(command.strategy_id);
2153
2154 let canceled_event = OrderCanceled::new(
2155 trader_id,
2156 strategy_id,
2157 command.instrument_id,
2158 client_order_id,
2159 UUID4::new(),
2160 command.ts_init,
2161 clock.get_time_ns(),
2162 false,
2163 Some(canceled_order.venue_order_id),
2164 Some(account_id),
2165 None,
2166 );
2167
2168 event_emitter.send_order_event(OrderEventAny::Canceled(canceled_event));
2169 }
2170
2171 Ok(())
2172 });
2173
2174 Ok(())
2175 }
2176
2177 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
2178 if cmd.cancels.is_empty() {
2179 return Ok(());
2180 }
2181
2182 for cancel in &cmd.cancels {
2184 self.cancel_order_internal(cancel);
2185 }
2186
2187 Ok(())
2188 }
2189}
2190
2191fn validate_order(order: &impl Order, use_gtd: bool) -> Result<(), OrderDeniedReason> {
2192 if order.is_reduce_only() {
2193 return Err(OrderDeniedReason::UnsupportedReduceOnly);
2194 }
2195
2196 let order_type = order.order_type();
2197 let venue_type = order_type_to_binance_spot(order_type, order.is_post_only())
2198 .map_err(|_| OrderDeniedReason::UnsupportedOrderType { order_type })?;
2199
2200 if matches!(
2201 venue_type,
2202 BinanceSpotOrderType::Limit
2203 | BinanceSpotOrderType::StopLossLimit
2204 | BinanceSpotOrderType::TakeProfitLimit
2205 ) {
2206 time_in_force_to_binance_spot(order.time_in_force(), use_gtd)
2207 .map_err(|_| OrderDeniedReason::UnsupportedTimeInForce(order.time_in_force()))?;
2208 }
2209
2210 if matches!(
2211 order_type,
2212 OrderType::StopMarket
2213 | OrderType::StopLimit
2214 | OrderType::MarketIfTouched
2215 | OrderType::LimitIfTouched
2216 ) && order.trigger_price().is_none()
2217 {
2218 return Err(OrderDeniedReason::ValidationFailed {
2219 detail: "Conditional orders require a trigger price".to_string(),
2220 });
2221 }
2222
2223 if order.is_quote_quantity() && order_type != OrderType::Market {
2224 return Err(OrderDeniedReason::ValidationFailed {
2225 detail: "Quote quantity requires a MARKET order on Binance Spot".to_string(),
2226 });
2227 }
2228
2229 Ok(())
2230}
2231
2232fn max_trade_id(reports: &[FillReport]) -> anyhow::Result<i64> {
2233 let mut max_trade_id = None;
2234
2235 for report in reports {
2236 let trade_id = parse_trade_id(report)?;
2237 max_trade_id = Some(max_trade_id.map_or(trade_id, |current: i64| current.max(trade_id)));
2238 }
2239
2240 max_trade_id.context("Binance Spot account-trades page was empty")
2241}
2242
2243fn parse_trade_id(report: &FillReport) -> anyhow::Result<i64> {
2244 report
2245 .trade_id
2246 .to_string()
2247 .parse::<i64>()
2248 .with_context(|| format!("invalid Binance Spot trade ID {}", report.trade_id))
2249}
2250
2251fn report_time_ms(report: &FillReport) -> i64 {
2252 report.ts_event.as_i64() / NANOSECONDS_IN_MILLISECOND as i64
2253}
2254
2255fn normalize_spot_order_status_report(
2256 report: &mut OrderStatusReport,
2257 treat_expired_as_canceled: bool,
2258) {
2259 if treat_expired_as_canceled && report.order_status == OrderStatus::Expired {
2260 report.order_status = OrderStatus::Canceled;
2261 }
2262}
2263
2264fn normalize_spot_order_status_reports(
2265 reports: &mut [OrderStatusReport],
2266 treat_expired_as_canceled: bool,
2267) {
2268 for report in reports {
2269 normalize_spot_order_status_report(report, treat_expired_as_canceled);
2270 }
2271}
2272
2273async fn wait_for_ws_setup_response(
2274 timeout: Duration,
2275 success: impl Future<Output = ()>,
2276 setup_errors: &mut tokio::sync::mpsc::UnboundedReceiver<String>,
2277 timeout_message: &'static str,
2278) -> anyhow::Result<()> {
2279 tokio::pin!(success);
2280
2281 let result = tokio::time::timeout(timeout, async {
2282 tokio::select! {
2283 () = &mut success => Ok(()),
2284 err = setup_errors.recv() => {
2285 anyhow::bail!(
2286 "{}",
2287 err.unwrap_or_else(|| "WS setup error channel closed".to_string()),
2288 )
2289 }
2290 }
2291 })
2292 .await;
2293
2294 result.map_err(|_| anyhow::anyhow!(timeout_message))?
2295}
2296
2297#[expect(clippy::too_many_arguments)]
2298fn dispatch_ws_trading_message(
2299 msg: BinanceSpotWsTradingMessage,
2300 emitter: &ExecutionEventEmitter,
2301 http_client: &BinanceSpotHttpClient,
2302 account_id: AccountId,
2303 treat_expired_as_canceled: bool,
2304 clock: &'static AtomicTime,
2305 dispatch_state: &WsDispatchState,
2306 ws_authenticated: &tokio::sync::Notify,
2307 ws_user_data_subscribed: &tokio::sync::Notify,
2308 ws_setup_error_tx: &tokio::sync::mpsc::UnboundedSender<String>,
2309 seen_trade_ids: &std::sync::Arc<Mutex<FifoCache<(Ustr, i64), 10_000>>>,
2310 task_spawner: &TaskSpawner,
2311) {
2312 match msg {
2313 BinanceSpotWsTradingMessage::OrderAccepted {
2314 request_id,
2315 response,
2316 } => {
2317 dispatch_state.pending_requests.remove(&request_id);
2318 log::debug!(
2319 "WS order accepted: request_id={request_id}, order_id={}",
2320 response.order_id
2321 );
2322 }
2324 BinanceSpotWsTradingMessage::OrderRejected {
2325 request_id,
2326 status,
2327 code,
2328 msg,
2329 } => {
2330 log::debug!(
2331 "WS order rejected: request_id={request_id}, status={status}, code={code}, msg={msg}"
2332 );
2333
2334 if let Some((_, pending)) = dispatch_state.pending_requests.remove(&request_id) {
2335 let code_i64 = i64::from(code);
2336 let reason = format!("code={code}: {msg}");
2337
2338 match classify_venue_failure(Some(code_i64), Some(status), &reason) {
2339 CommandFailure::Ambiguous(_) => {
2340 log::warn!(
2341 "Ambiguous WS submit failure for {}, awaiting reconciliation: {reason}",
2342 pending.client_order_id,
2343 );
2344 return;
2345 }
2346 CommandFailure::NotSent(_) | CommandFailure::VenueRejected(_) => {}
2347 }
2348
2349 let identity = dispatch_state
2351 .order_identities
2352 .get(&pending.client_order_id)
2353 .map(|r| r.clone());
2354
2355 if let Some(identity) = identity {
2356 let due_post_only = code_i64 == BINANCE_GTX_ORDER_REJECT_CODE
2357 || (code_i64 == BINANCE_NEW_ORDER_REJECTED_CODE
2358 && msg == BINANCE_SPOT_POST_ONLY_REJECT_MSG);
2359 let ts_now = clock.get_time_ns();
2360 let rejected = OrderRejected::new(
2361 emitter.trader_id(),
2362 identity.strategy_id,
2363 identity.instrument_id,
2364 pending.client_order_id,
2365 account_id,
2366 Ustr::from(&sanitize_reason(&reason)),
2367 UUID4::new(),
2368 ts_now,
2369 ts_now,
2370 false,
2371 due_post_only,
2372 );
2373 dispatch_state.cleanup_terminal(pending.client_order_id);
2374 emitter.send_order_event(OrderEventAny::Rejected(rejected));
2375 } else {
2376 log::warn!(
2377 "No order identity for {}, cannot emit OrderRejected",
2378 pending.client_order_id
2379 );
2380 }
2381 } else {
2382 log::warn!("No pending request for {request_id}, cannot emit OrderRejected");
2383 }
2384 }
2385 BinanceSpotWsTradingMessage::OrderCanceled {
2386 request_id,
2387 response,
2388 } => {
2389 dispatch_state.pending_requests.remove(&request_id);
2390 log::debug!(
2391 "WS order canceled: request_id={request_id}, order_id={}",
2392 response.order_id
2393 );
2394 }
2396 BinanceSpotWsTradingMessage::CancelRejected {
2397 request_id,
2398 status,
2399 code,
2400 msg,
2401 } => {
2402 log::debug!(
2403 "WS cancel rejected: request_id={request_id}, status={status}, code={code}, msg={msg}"
2404 );
2405
2406 if let Some((_, pending)) = dispatch_state.pending_requests.remove(&request_id) {
2407 let reason = format!("code={code}: {msg}");
2408
2409 match classify_venue_failure(Some(i64::from(code)), Some(status), &reason) {
2410 CommandFailure::Ambiguous(_) => {
2411 log::warn!(
2412 "Ambiguous WS cancel failure for {}, awaiting reconciliation: {reason}",
2413 pending.client_order_id,
2414 );
2415 return;
2416 }
2417 CommandFailure::NotSent(_) | CommandFailure::VenueRejected(_) => {}
2418 }
2419
2420 if let Some(identity) = dispatch_state
2421 .order_identities
2422 .get(&pending.client_order_id)
2423 {
2424 let ts_now = clock.get_time_ns();
2425 let rejected = OrderCancelRejected::new(
2426 emitter.trader_id(),
2427 identity.strategy_id,
2428 identity.instrument_id,
2429 pending.client_order_id,
2430 Ustr::from(&sanitize_reason(&reason)),
2431 UUID4::new(),
2432 ts_now,
2433 ts_now,
2434 false,
2435 pending.venue_order_id,
2436 Some(account_id),
2437 );
2438 emitter.send_order_event(OrderEventAny::CancelRejected(rejected));
2439 }
2440 }
2441 }
2442 BinanceSpotWsTradingMessage::CancelReplaceAccepted {
2443 request_id,
2444 cancel_response,
2445 new_order_response,
2446 } => {
2447 dispatch_state.pending_requests.remove(&request_id);
2448 dispatch_state
2449 .cancel_replace_request_ids
2450 .remove(&request_id);
2451 log::debug!(
2452 "WS cancel-replace accepted: request_id={request_id}, \
2453 canceled_id={}, new_id={}",
2454 cancel_response.order_id,
2455 new_order_response.order_id,
2456 );
2457 }
2459 BinanceSpotWsTradingMessage::CancelReplaceRejected {
2460 request_id,
2461 status,
2462 code,
2463 msg,
2464 } => {
2465 log::debug!(
2466 "WS cancel-replace rejected: request_id={request_id}, status={status}, code={code}, msg={msg}"
2467 );
2468
2469 let reason = format!("code={code}: {msg}");
2470 let is_ambiguous = matches!(
2471 classify_venue_failure(Some(i64::from(code)), Some(status), &reason),
2472 CommandFailure::Ambiguous(_)
2473 );
2474 let confirmed_cancel = dispatch_state
2475 .cancel_replace_request_ids
2476 .remove(&request_id)
2477 .filter(|_| !is_ambiguous)
2478 .and_then(|(_, cancel_id)| dispatch_state.on_cancel_replace_rejected(&cancel_id));
2479 let pending = dispatch_state
2480 .pending_requests
2481 .remove(&request_id)
2482 .map(|(_, pending)| pending);
2483
2484 if is_ambiguous {
2485 log::warn!("Ambiguous WS modify failure, awaiting reconciliation: {reason}");
2486 return;
2487 }
2488
2489 if let Some(pending) = &pending {
2490 if let Some(canceled) = dispatch_state.reject_replace(pending.client_order_id) {
2491 dispatch_state.cleanup_terminal(pending.client_order_id);
2492 emitter.send_order_event(OrderEventAny::Canceled(canceled));
2493 return;
2494 }
2495
2496 if let Some(identity) = dispatch_state
2497 .order_identities
2498 .get(&pending.client_order_id)
2499 {
2500 let ts_now = clock.get_time_ns();
2501 let rejected = OrderModifyRejected::new(
2502 emitter.trader_id(),
2503 identity.strategy_id,
2504 identity.instrument_id,
2505 pending.client_order_id,
2506 Ustr::from(&sanitize_reason(&reason)),
2507 UUID4::new(),
2508 ts_now,
2509 ts_now,
2510 false,
2511 pending.venue_order_id,
2512 Some(account_id),
2513 );
2514 emitter.send_order_event(OrderEventAny::ModifyRejected(rejected));
2515 }
2516 }
2517
2518 if let Some(report) = confirmed_cancel {
2519 emit_canceled_after_rejected_replace(
2520 emitter,
2521 dispatch_state,
2522 account_id,
2523 report,
2524 clock.get_time_ns(),
2525 );
2526 }
2527 }
2528 BinanceSpotWsTradingMessage::RequestFailed { request_id, msg } => {
2529 dispatch_state.pending_requests.remove(&request_id);
2530 dispatch_state
2531 .cancel_replace_request_ids
2532 .remove(&request_id);
2533 log::error!(
2534 "WS trading request failed without structured venue response: request_id={request_id}, {msg}"
2535 );
2536 }
2537 BinanceSpotWsTradingMessage::AllOrdersCanceled {
2538 request_id,
2539 responses,
2540 } => {
2541 dispatch_state.pending_requests.remove(&request_id);
2542 log::debug!(
2543 "WS all orders canceled: request_id={request_id}, count={}",
2544 responses.len()
2545 );
2546
2547 match prepare_cancel_all_orders(responses) {
2548 Ok(canceled_orders) => {
2549 let ts_init = clock.get_time_ns();
2550
2551 for canceled_order in canceled_orders
2552 .iter()
2553 .filter(|canceled_order| canceled_order.order_list)
2554 {
2555 dispatch_order_list_canceled(
2556 canceled_order,
2557 emitter,
2558 account_id,
2559 dispatch_state,
2560 ts_init,
2561 );
2562 }
2563 }
2564 Err(e) => log::error!(
2565 "Ignoring invalid WS cancel-all response {request_id} to avoid false terminal events: {e}"
2566 ),
2567 }
2568 }
2569 BinanceSpotWsTradingMessage::UserDataSubscribed { subscription_id } => {
2570 log::debug!("User data stream subscribed: id={subscription_id}");
2571 ws_user_data_subscribed.notify_one();
2572 }
2573 BinanceSpotWsTradingMessage::ExecutionReport(report) => {
2574 let ts_init = clock.get_time_ns();
2575 dispatch_execution_report(
2576 &report,
2577 emitter,
2578 http_client,
2579 account_id,
2580 treat_expired_as_canceled,
2581 dispatch_state,
2582 seen_trade_ids,
2583 ts_init,
2584 );
2585 }
2586 BinanceSpotWsTradingMessage::AccountPosition(position) => {
2587 let ts_init = clock.get_time_ns();
2588 let state = parse_spot_account_position(&position, account_id, ts_init);
2589 emitter.send_account_state(state);
2590 }
2591 BinanceSpotWsTradingMessage::BalanceUpdate(update) => {
2592 log::debug!(
2593 "Balance update: asset={}, delta={}",
2594 update.asset,
2595 update.delta,
2596 );
2597 let http_client = http_client.clone();
2598 let emitter = emitter.clone();
2599
2600 if let Err(e) = task_spawner.spawn(async move {
2601 match http_client.request_account_state(account_id).await {
2602 Ok(state) => emitter.send_account_state(state),
2603 Err(e) => {
2604 log::error!("Failed to refresh account state after balance update: {e}");
2605 }
2606 }
2607 }) {
2608 log::warn!("Skipping Binance Spot balance refresh after shutdown began: {e}");
2609 }
2610 }
2611 BinanceSpotWsTradingMessage::Connected => {
2612 log::debug!("WS trading API connected");
2613 }
2614 BinanceSpotWsTradingMessage::Authenticated => {
2615 log::debug!("WS trading API authenticated");
2616 ws_authenticated.notify_one();
2617 }
2618 BinanceSpotWsTradingMessage::AuthenticationRejected(reason) => {
2619 log::error!("WS trading API authentication failed: {reason}");
2620 let _ = ws_setup_error_tx.send(reason);
2621 }
2622 BinanceSpotWsTradingMessage::Reconnected => {
2623 log::info!("WS trading API reconnected");
2624 }
2625 BinanceSpotWsTradingMessage::ServerShutdown { event_time } => {
2626 log::warn!(
2627 "WS trading API server shutdown notice (event_time={event_time}); reconnect expected within ~10 minutes"
2628 );
2629 }
2630 BinanceSpotWsTradingMessage::Error(err) => {
2631 log::error!("WS trading API error: {err}");
2632 let _ = ws_setup_error_tx.send(err);
2633 }
2634 BinanceSpotWsTradingMessage::UserDataSubscriptionRejected(reason) => {
2635 log::error!("WS trading API user data subscription failed: {reason}");
2636 let _ = ws_setup_error_tx.send(reason);
2637 }
2638 }
2639}
2640
2641#[derive(Debug)]
2642struct PreparedCancelOrder {
2643 venue_order_id: VenueOrderId,
2644 transaction_time: i64,
2645 client_order_id: ClientOrderId,
2646 instrument_id: InstrumentId,
2647 order_list: bool,
2648}
2649
2650fn prepare_cancel_all_orders(
2651 responses: Vec<BinanceCancelOpenOrdersResponse>,
2652) -> anyhow::Result<Vec<PreparedCancelOrder>> {
2653 let mut prepared = Vec::new();
2654 let mut order_ids = AHashSet::new();
2655 let mut client_order_ids = AHashSet::new();
2656
2657 for response in responses {
2658 match response {
2659 BinanceCancelOpenOrdersResponse::Order(response) => prepare_cancel_order(
2660 &response,
2661 false,
2662 &mut order_ids,
2663 &mut client_order_ids,
2664 &mut prepared,
2665 )?,
2666 BinanceCancelOpenOrdersResponse::OrderList(response) => {
2667 prepare_cancel_order_list(
2668 response,
2669 &mut order_ids,
2670 &mut client_order_ids,
2671 &mut prepared,
2672 )?;
2673 }
2674 }
2675 }
2676
2677 Ok(prepared)
2678}
2679
2680fn prepare_cancel_order_list(
2681 response: BinanceCancelOrderListResponse,
2682 order_ids: &mut AHashSet<(InstrumentId, i64)>,
2683 client_order_ids: &mut AHashSet<ClientOrderId>,
2684 prepared: &mut Vec<PreparedCancelOrder>,
2685) -> anyhow::Result<()> {
2686 anyhow::ensure!(
2687 response.order_list_id >= 0 && !response.symbol.is_empty(),
2688 "order list has an invalid list ID or empty symbol",
2689 );
2690 anyhow::ensure!(
2691 response.list_status_type == SbeListStatusType::AllDone
2692 && response.list_order_status == SbeListOrderStatus::AllDone,
2693 "order list {} was not fully canceled: status={:?}, order_status={:?}",
2694 response.order_list_id,
2695 response.list_status_type,
2696 response.list_order_status,
2697 );
2698 anyhow::ensure!(
2699 !response.orders.is_empty() && response.orders.len() == response.order_reports.len(),
2700 "order list {} has {} orders and {} reports",
2701 response.order_list_id,
2702 response.orders.len(),
2703 response.order_reports.len(),
2704 );
2705
2706 let mut list_orders = AHashMap::with_capacity(response.orders.len());
2707 for order in response.orders {
2708 anyhow::ensure!(
2709 order.symbol == response.symbol,
2710 "order list {} contains order {} for symbol {}, expected {}",
2711 response.order_list_id,
2712 order.order_id,
2713 order.symbol,
2714 response.symbol,
2715 );
2716 anyhow::ensure!(
2717 list_orders.insert(order.order_id, order).is_none(),
2718 "order list {} contains a duplicate order ID",
2719 response.order_list_id,
2720 );
2721 }
2722
2723 for report in response.order_reports {
2724 anyhow::ensure!(
2725 report.order_list_id == Some(response.order_list_id),
2726 "order {} reports order-list ID {:?}, expected {}",
2727 report.order_id,
2728 report.order_list_id,
2729 response.order_list_id,
2730 );
2731 let order = list_orders.remove(&report.order_id).with_context(|| {
2732 format!(
2733 "order-list report {} is absent from list {}",
2734 report.order_id, response.order_list_id
2735 )
2736 })?;
2737 anyhow::ensure!(
2738 report.symbol == response.symbol
2739 && report.symbol == order.symbol
2740 && report.orig_client_order_id == order.client_order_id,
2741 "order-list report {} does not match its order identity",
2742 report.order_id,
2743 );
2744 prepare_cancel_order(&report, true, order_ids, client_order_ids, prepared)?;
2745 }
2746
2747 anyhow::ensure!(
2748 list_orders.is_empty(),
2749 "order list {} is missing {} child reports",
2750 response.order_list_id,
2751 list_orders.len(),
2752 );
2753 Ok(())
2754}
2755
2756fn prepare_cancel_order(
2757 response: &BinanceCancelOrderResponse,
2758 order_list: bool,
2759 order_ids: &mut AHashSet<(InstrumentId, i64)>,
2760 client_order_ids: &mut AHashSet<ClientOrderId>,
2761 prepared: &mut Vec<PreparedCancelOrder>,
2762) -> anyhow::Result<()> {
2763 anyhow::ensure!(
2764 response.order_id >= 0
2765 && !response.symbol.is_empty()
2766 && !response.orig_client_order_id.is_empty(),
2767 "cancel-all response has an invalid order ID, symbol, or original client order ID",
2768 );
2769 anyhow::ensure!(
2770 response.status == SbeOrderStatus::Canceled,
2771 "order {} reports status {:?}, expected Canceled",
2772 response.order_id,
2773 response.status,
2774 );
2775 let client_order_id = decode_client_order_id(
2776 &response.orig_client_order_id,
2777 BINANCE_NAUTILUS_SPOT_BROKER_ID,
2778 )?;
2779 let instrument_id = InstrumentId::new(response.symbol.as_str().into(), *BINANCE_VENUE);
2780 anyhow::ensure!(
2781 order_ids.insert((instrument_id, response.order_id)),
2782 "cancel-all response contains duplicate order ID {} for {}",
2783 response.order_id,
2784 instrument_id,
2785 );
2786 anyhow::ensure!(
2787 client_order_ids.insert(client_order_id),
2788 "cancel-all response contains duplicate client order ID {client_order_id}",
2789 );
2790 prepared.push(PreparedCancelOrder {
2791 venue_order_id: VenueOrderId::new(response.order_id.to_string()),
2792 transaction_time: response.transact_time,
2793 client_order_id,
2794 instrument_id,
2795 order_list,
2796 });
2797 Ok(())
2798}
2799
2800fn dispatch_order_list_canceled(
2801 canceled_order: &PreparedCancelOrder,
2802 emitter: &ExecutionEventEmitter,
2803 account_id: AccountId,
2804 state: &WsDispatchState,
2805 ts_init: UnixNanos,
2806) {
2807 let client_order_id = canceled_order.client_order_id;
2808 let identity = state
2809 .order_identities
2810 .remove_if(&client_order_id, |_, identity| {
2811 identity.instrument_id == canceled_order.instrument_id
2812 });
2813
2814 let Some((_, identity)) = identity else {
2815 if state.order_identities.contains_key(&client_order_id) {
2816 log::error!(
2817 "Ignoring cancel-all result for {client_order_id}: tracked instrument does not match {}",
2818 canceled_order.instrument_id,
2819 );
2820 } else {
2821 state
2822 .pending_requests
2823 .retain(|_, pending| pending.client_order_id != client_order_id);
2824 log::debug!("Skipping duplicate cancel-all result for {client_order_id}");
2825 }
2826 return;
2827 };
2828
2829 state
2830 .pending_requests
2831 .retain(|_, pending| pending.client_order_id != client_order_id);
2832 let ts_event = parse_micros_or_init(
2833 canceled_order.transaction_time,
2834 "Spot cancel-all transaction time",
2835 ts_init,
2836 );
2837 ensure_accepted_emitted(
2838 client_order_id,
2839 account_id,
2840 canceled_order.venue_order_id,
2841 &identity,
2842 emitter,
2843 state,
2844 ts_event,
2845 );
2846 state.cleanup_terminal(client_order_id);
2847 let canceled = OrderCanceled::new(
2848 emitter.trader_id(),
2849 identity.strategy_id,
2850 identity.instrument_id,
2851 client_order_id,
2852 UUID4::new(),
2853 ts_event,
2854 ts_init,
2855 false,
2856 Some(canceled_order.venue_order_id),
2857 Some(account_id),
2858 None,
2859 );
2860 emitter.send_order_event(OrderEventAny::Canceled(canceled));
2861}
2862
2863fn build_new_order_params(
2864 order: &impl Order,
2865 client_order_id: ClientOrderId,
2866 is_post_only: bool,
2867 is_quote_quantity: bool,
2868 use_gtd: bool,
2869) -> anyhow::Result<NewOrderParams> {
2870 let binance_side = BinanceSide::try_from(order.order_side())?;
2871 let binance_order_type = order_type_to_binance_spot(order.order_type(), is_post_only)?;
2872
2873 let requires_trigger = matches!(
2874 order.order_type(),
2875 OrderType::StopMarket
2876 | OrderType::StopLimit
2877 | OrderType::MarketIfTouched
2878 | OrderType::LimitIfTouched
2879 );
2880
2881 if requires_trigger && order.trigger_price().is_none() {
2882 anyhow::bail!("Conditional orders require a trigger price");
2883 }
2884
2885 let supports_tif = matches!(
2886 binance_order_type,
2887 BinanceSpotOrderType::Limit
2888 | BinanceSpotOrderType::StopLossLimit
2889 | BinanceSpotOrderType::TakeProfitLimit
2890 );
2891 let binance_tif = if supports_tif {
2892 Some(time_in_force_to_binance_spot(
2893 order.time_in_force(),
2894 use_gtd,
2895 )?)
2896 } else {
2897 None
2898 };
2899
2900 let qty_str = order.quantity().to_string();
2901 let (base_qty, quote_qty) = if is_quote_quantity {
2902 (None, Some(qty_str))
2903 } else {
2904 (Some(qty_str), None)
2905 };
2906
2907 let client_id_str = encode_broker_id(&client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
2908
2909 Ok(NewOrderParams {
2910 symbol: order.instrument_id().symbol.to_string(),
2911 side: binance_side,
2912 order_type: binance_order_type,
2913 time_in_force: binance_tif,
2914 quantity: base_qty,
2915 quote_order_qty: quote_qty,
2916 price: order.price().map(|p| p.to_string()),
2917 new_client_order_id: Some(client_id_str),
2918 stop_price: order.trigger_price().map(|p| p.to_string()),
2919 trailing_delta: None,
2920 iceberg_qty: order.display_qty().map(|q| q.to_string()),
2921 new_order_resp_type: Some(BinanceOrderResponseType::Full),
2922 self_trade_prevention_mode: None,
2923 strategy_id: None,
2924 strategy_type: None,
2925 })
2926}
2927
2928fn build_spot_order_list_params(
2929 order_list_id: &str,
2930 orders: &[OrderAny],
2931 use_gtd: bool,
2932) -> Result<NewOcoOrderListParams, String> {
2933 let has_grouped_order = orders.iter().any(is_grouped_order);
2934
2935 if has_grouped_order {
2936 return build_spot_oco_order_list_params(order_list_id, orders, use_gtd);
2937 }
2938
2939 Err("Binance Spot order-list submission currently supports only OCO lists".to_string())
2940}
2941
2942fn build_spot_oco_order_list_params(
2943 order_list_id: &str,
2944 orders: &[OrderAny],
2945 use_gtd: bool,
2946) -> Result<NewOcoOrderListParams, String> {
2947 if orders.len() != 2 {
2948 return Err(format!(
2949 "Binance Spot OCO order-list submission requires exactly 2 orders, was {}",
2950 orders.len()
2951 ));
2952 }
2953
2954 if orders
2955 .iter()
2956 .any(|order| order.contingency_type() != Some(ContingencyType::Oco))
2957 {
2958 return Err(
2959 "Binance Spot grouped order-list submission currently supports only OCO lists"
2960 .to_string(),
2961 );
2962 }
2963
2964 let first = &orders[0];
2965 let second = &orders[1];
2966 if first.instrument_id() != second.instrument_id() {
2967 return Err("Binance Spot OCO order-list legs must use the same instrument".to_string());
2968 }
2969
2970 if first.order_side() != second.order_side() {
2971 return Err("Binance Spot OCO order-list legs must use the same side".to_string());
2972 }
2973
2974 if first.quantity() != second.quantity() {
2975 return Err("Binance Spot OCO order-list legs must use the same quantity".to_string());
2976 }
2977
2978 if first.is_quote_quantity() || second.is_quote_quantity() {
2979 return Err("Binance Spot OCO order-list legs do not support quote quantity".to_string());
2980 }
2981
2982 let mut above = None;
2983 let mut below = None;
2984
2985 for order in orders {
2986 let params = build_new_order_params(
2987 order,
2988 order.client_order_id(),
2989 order.is_post_only(),
2990 false,
2991 use_gtd,
2992 )
2993 .map_err(|e| e.to_string())?;
2994
2995 match spot_oco_leg_position(params.side, params.order_type)? {
2996 SpotOcoLegPosition::Above => {
2997 if above.replace(params).is_some() {
2998 return Err(
2999 "Binance Spot OCO order-list resolved more than one above leg".to_string(),
3000 );
3001 }
3002 }
3003 SpotOcoLegPosition::Below => {
3004 if below.replace(params).is_some() {
3005 return Err(
3006 "Binance Spot OCO order-list resolved more than one below leg".to_string(),
3007 );
3008 }
3009 }
3010 }
3011 }
3012
3013 let above = above.ok_or_else(|| "Binance Spot OCO order-list missing above leg".to_string())?;
3014 let below = below.ok_or_else(|| "Binance Spot OCO order-list missing below leg".to_string())?;
3015 let quantity = above
3016 .quantity
3017 .clone()
3018 .ok_or_else(|| "Binance Spot OCO order-list requires base quantity".to_string())?;
3019
3020 Ok(NewOcoOrderListParams {
3021 symbol: first.instrument_id().symbol.to_string(),
3022 list_client_order_id: Some(order_list_id.to_string()),
3023 side: above.side,
3024 quantity,
3025 above_type: above.order_type,
3026 above_client_order_id: above.new_client_order_id,
3027 above_iceberg_qty: above.iceberg_qty,
3028 above_price: above.price,
3029 above_stop_price: above.stop_price,
3030 above_time_in_force: above.time_in_force,
3031 below_type: below.order_type,
3032 below_client_order_id: below.new_client_order_id,
3033 below_iceberg_qty: below.iceberg_qty,
3034 below_price: below.price,
3035 below_stop_price: below.stop_price,
3036 below_time_in_force: below.time_in_force,
3037 new_order_resp_type: Some(BinanceOrderResponseType::Full),
3038 self_trade_prevention_mode: None,
3039 })
3040}
3041
3042enum SpotOcoLegPosition {
3043 Above,
3044 Below,
3045}
3046
3047fn spot_oco_leg_position(
3048 side: BinanceSide,
3049 order_type: BinanceSpotOrderType,
3050) -> Result<SpotOcoLegPosition, String> {
3051 match (side, order_type) {
3052 (
3053 BinanceSide::Sell,
3054 BinanceSpotOrderType::LimitMaker
3055 | BinanceSpotOrderType::TakeProfit
3056 | BinanceSpotOrderType::TakeProfitLimit,
3057 )
3058 | (
3059 BinanceSide::Buy,
3060 BinanceSpotOrderType::StopLoss | BinanceSpotOrderType::StopLossLimit,
3061 ) => Ok(SpotOcoLegPosition::Above),
3062 (
3063 BinanceSide::Sell,
3064 BinanceSpotOrderType::StopLoss | BinanceSpotOrderType::StopLossLimit,
3065 )
3066 | (
3067 BinanceSide::Buy,
3068 BinanceSpotOrderType::LimitMaker
3069 | BinanceSpotOrderType::TakeProfit
3070 | BinanceSpotOrderType::TakeProfitLimit,
3071 ) => Ok(SpotOcoLegPosition::Below),
3072 (_, unsupported) => Err(format!(
3073 "Unsupported Binance Spot OCO leg order type: {unsupported:?}"
3074 )),
3075 }
3076}
3077
3078fn is_grouped_order(order: &OrderAny) -> bool {
3079 order.contingency_type().is_some()
3080 || order
3081 .linked_order_ids()
3082 .is_some_and(|linked_order_ids| !linked_order_ids.is_empty())
3083}
3084
3085fn handle_spot_order_submit_success(client_order_id: ClientOrderId, venue_order_id: VenueOrderId) {
3086 log::debug!(
3087 "Order submit succeeded: client_order_id={client_order_id}, venue_order_id={venue_order_id}",
3088 );
3089}
3090
3091async fn submit_spot_order_list(
3092 http_client: &BinanceSpotHttpClient,
3093 params: &NewOcoOrderListParams,
3094) -> Result<(), BinanceSpotHttpError> {
3095 let response = http_client.submit_oco_order_list(params).await?;
3096 log::debug!(
3097 "Order list submit succeeded: order_list_id={}, order_count={}",
3098 response.order_list_id,
3099 response.orders.len(),
3100 );
3101 Ok(())
3102}
3103
3104fn handle_spot_order_list_submit_error(
3105 event_emitter: &ExecutionEventEmitter,
3106 dispatch_state: &WsDispatchState,
3107 trader_id: TraderId,
3108 account_id: AccountId,
3109 clock: &'static AtomicTime,
3110 orders: &[OrderAny],
3111 error: BinanceSpotHttpError,
3112) -> anyhow::Result<()> {
3113 match classify_spot_http_failure(&error) {
3114 CommandFailure::Ambiguous(reason) => {
3115 log::error!("Ambiguous order-list submit failure, awaiting reconciliation: {reason}");
3116 }
3117 CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
3118 let ts_now = clock.get_time_ns();
3121 let reason = format!("submit-order-list-error: {}", sanitize_reason(&reason));
3122 for order in orders {
3123 let client_order_id = order.client_order_id();
3124 dispatch_state.cleanup_terminal(client_order_id);
3125 let rejected = OrderRejected::new(
3126 trader_id,
3127 order.strategy_id(),
3128 order.instrument_id(),
3129 client_order_id,
3130 account_id,
3131 reason.clone().into(),
3132 UUID4::new(),
3133 ts_now,
3134 ts_now,
3135 false,
3136 false,
3137 );
3138 event_emitter.send_order_event(OrderEventAny::Rejected(rejected));
3139 }
3140 }
3141 }
3142
3143 Err(error.into())
3144}
3145
3146fn build_cancel_order_params(cmd: &CancelOrder, prefer_client_order_id: bool) -> CancelOrderParams {
3147 let order_id = cmd
3148 .venue_order_id
3149 .and_then(|id| id.inner().parse::<i64>().ok());
3150
3151 if let Some(order_id) = order_id
3152 && !prefer_client_order_id
3153 {
3154 CancelOrderParams::by_order_id(cmd.instrument_id.symbol.to_string(), order_id)
3155 } else {
3156 let client_id_str = encode_broker_id(&cmd.client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
3157 CancelOrderParams::by_client_order_id(cmd.instrument_id.symbol.to_string(), client_id_str)
3158 }
3159}
3160
3161fn cancel_replace_cancel_id() -> String {
3166 format!("CR-{}", UUID4::new().as_str().replace('-', ""))
3167}
3168
3169fn emit_canceled_after_rejected_replace(
3174 emitter: &ExecutionEventEmitter,
3175 state: &WsDispatchState,
3176 account_id: AccountId,
3177 report: OrderStatusReport,
3178 ts_init: UnixNanos,
3179) {
3180 let identity = report
3181 .client_order_id
3182 .and_then(|cid| state.order_identities.get(&cid).map(|r| r.clone()));
3183 if let Some(client_order_id) = report.client_order_id {
3184 state.cleanup_terminal(client_order_id);
3185 }
3186 let (Some(client_order_id), Some(identity)) = (report.client_order_id, identity) else {
3187 emitter.send_order_status_report(report);
3188 return;
3189 };
3190
3191 let canceled = OrderCanceled::new(
3192 emitter.trader_id(),
3193 identity.strategy_id,
3194 identity.instrument_id,
3195 client_order_id,
3196 UUID4::new(),
3197 report.ts_last,
3198 ts_init,
3199 false,
3200 Some(report.venue_order_id),
3201 Some(account_id),
3202 None,
3203 );
3204 emitter.send_order_event(OrderEventAny::Canceled(canceled));
3205}
3206
3207fn handle_http_modify_failure(
3211 error: &anyhow::Error,
3212 command: &ModifyOrder,
3213 cancel_id: &str,
3214 emitter: &ExecutionEventEmitter,
3215 state: &WsDispatchState,
3216 account_id: AccountId,
3217 ts_now: UnixNanos,
3218) {
3219 let failure = error.downcast_ref::<BinanceSpotHttpError>().map_or_else(
3220 || CommandFailure::Ambiguous(error.to_string()),
3221 classify_spot_http_failure,
3222 );
3223
3224 match failure {
3225 CommandFailure::Ambiguous(reason) => {
3226 log::warn!(
3227 "Ambiguous modify failure for {}, awaiting reconciliation: {reason}",
3228 command.client_order_id
3229 );
3230 }
3231 CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
3232 if let Some(canceled) = state.reject_replace(command.client_order_id) {
3233 state.cleanup_terminal(command.client_order_id);
3234 emitter.send_order_event(OrderEventAny::Canceled(canceled));
3235 return;
3236 }
3237
3238 let rejected_event = OrderModifyRejected::new(
3239 emitter.trader_id(),
3240 command.strategy_id,
3241 command.instrument_id,
3242 command.client_order_id,
3243 format!("modify-order-error: {}", sanitize_reason(&reason)).into(),
3244 UUID4::new(),
3245 ts_now,
3246 ts_now,
3247 false,
3248 command.venue_order_id,
3249 Some(account_id),
3250 );
3251 emitter.send_order_event(OrderEventAny::ModifyRejected(rejected_event));
3252
3253 if let Some(report) = state.on_cancel_replace_rejected(cancel_id) {
3254 emit_canceled_after_rejected_replace(emitter, state, account_id, report, ts_now);
3255 }
3256 }
3257 }
3258}
3259
3260fn build_cancel_replace_params(
3261 cmd: &ModifyOrder,
3262 order: &impl Order,
3263 quantity: Quantity,
3264 use_gtd: bool,
3265 cancel_new_client_order_id: String,
3266) -> anyhow::Result<CancelReplaceOrderParams> {
3267 let binance_side = BinanceSide::try_from(order.order_side())?;
3268 let binance_order_type = order_type_to_binance_spot(order.order_type(), false)?;
3269 let binance_tif = time_in_force_to_binance_spot(order.time_in_force(), use_gtd)?;
3270
3271 let cancel_order_id: Option<i64> = cmd
3272 .venue_order_id
3273 .map(|id| {
3274 id.inner()
3275 .parse::<i64>()
3276 .map_err(|_| anyhow::anyhow!("Invalid venue order ID: {id}"))
3277 })
3278 .transpose()?;
3279
3280 let client_id_str = encode_broker_id(&cmd.client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
3281
3282 Ok(CancelReplaceOrderParams {
3283 symbol: cmd.instrument_id.symbol.to_string(),
3284 side: binance_side,
3285 order_type: binance_order_type,
3286 cancel_replace_mode: BinanceCancelReplaceMode::StopOnFailure,
3287 time_in_force: Some(binance_tif),
3288 quantity: Some(quantity.to_string()),
3289 quote_order_qty: None,
3290 price: cmd.price.map(|p| p.to_string()),
3291 cancel_order_id,
3292 cancel_orig_client_order_id: if cancel_order_id.is_none() {
3293 Some(client_id_str.clone())
3294 } else {
3295 None
3296 },
3297 cancel_new_client_order_id: Some(cancel_new_client_order_id),
3298 new_client_order_id: Some(client_id_str),
3299 stop_price: None,
3300 trailing_delta: None,
3301 iceberg_qty: None,
3302 new_order_resp_type: Some(BinanceOrderResponseType::Full),
3303 self_trade_prevention_mode: None,
3304 })
3305}
3306
3307#[expect(clippy::too_many_arguments)]
3312fn dispatch_execution_report(
3313 report: &BinanceSpotExecutionReport,
3314 emitter: &ExecutionEventEmitter,
3315 http_client: &BinanceSpotHttpClient,
3316 account_id: AccountId,
3317 treat_expired_as_canceled: bool,
3318 dispatch_state: &WsDispatchState,
3319 seen_trade_ids: &std::sync::Arc<Mutex<FifoCache<(Ustr, i64), 10_000>>>,
3320 ts_init: UnixNanos,
3321) {
3322 let symbol = report.symbol;
3323 let instrument_id = InstrumentId::new(symbol.into(), *BINANCE_VENUE);
3324 let Some(instrument) = http_client.get_instrument(&symbol) else {
3325 log::error!(
3326 "Cannot dispatch Spot execution report for uncached instrument {instrument_id}"
3327 );
3328 return;
3329 };
3330 let (price_precision, size_precision) =
3331 (instrument.price_precision(), instrument.size_precision());
3332
3333 if report.execution_type == BinanceSpotExecutionType::Canceled
3334 && dispatch_state.has_cancel_replace(&report.client_order_id)
3335 {
3336 match parse_spot_exec_report_to_order_status(
3337 report,
3338 instrument_id,
3339 price_precision,
3340 size_precision,
3341 account_id,
3342 treat_expired_as_canceled,
3343 ts_init,
3344 ) {
3345 Ok(status) => {
3346 if dispatch_state.on_cancel_replace_canceled(&report.client_order_id, status) {
3347 log::debug!(
3349 "Withholding cancel-replace cancel report: client_order_id={}, venue_order_id={}",
3350 report.order_client_order_id(),
3351 report.order_id
3352 );
3353 return;
3354 }
3355 }
3356 Err(e) => log::warn!(
3357 "Cannot withhold cancel-replace cancel report, dispatching normally: {e}"
3358 ),
3359 }
3360 }
3361
3362 let client_order_id = match decode_client_order_id(
3363 report.order_client_order_id(),
3364 BINANCE_NAUTILUS_SPOT_BROKER_ID,
3365 ) {
3366 Ok(client_order_id) => client_order_id,
3367 Err(e) => {
3368 log::warn!("Skipping Spot execution report with invalid client order ID: {e}");
3369 return;
3370 }
3371 };
3372
3373 let identity = dispatch_state
3374 .order_identities
3375 .get(&client_order_id)
3376 .map(|r| r.clone());
3377
3378 if let Some(identity) = identity {
3379 dispatch_tracked_execution_report(
3380 report,
3381 emitter,
3382 account_id,
3383 treat_expired_as_canceled,
3384 dispatch_state,
3385 seen_trade_ids,
3386 client_order_id,
3387 &identity,
3388 instrument_id,
3389 price_precision,
3390 size_precision,
3391 instrument.quote_currency(),
3392 ts_init,
3393 );
3394 } else {
3395 dispatch_untracked_execution_report(
3396 report,
3397 emitter,
3398 http_client,
3399 account_id,
3400 treat_expired_as_canceled,
3401 seen_trade_ids,
3402 instrument_id,
3403 price_precision,
3404 size_precision,
3405 ts_init,
3406 );
3407 }
3408}
3409
3410#[expect(clippy::too_many_arguments)]
3412fn dispatch_tracked_execution_report(
3413 report: &BinanceSpotExecutionReport,
3414 emitter: &ExecutionEventEmitter,
3415 account_id: AccountId,
3416 treat_expired_as_canceled: bool,
3417 state: &WsDispatchState,
3418 seen_trade_ids: &std::sync::Arc<Mutex<FifoCache<(Ustr, i64), 10_000>>>,
3419 client_order_id: ClientOrderId,
3420 identity: &OrderIdentity,
3421 instrument_id: InstrumentId,
3422 price_precision: u8,
3423 size_precision: u8,
3424 quote_currency: Currency,
3425 ts_init: UnixNanos,
3426) {
3427 let venue_order_id = VenueOrderId::new(report.order_id.to_string());
3428 let ts_event = parse_millis_or_init(report.event_time, "Spot execution event time", ts_init);
3429
3430 match report.execution_type {
3431 BinanceSpotExecutionType::Unknown => {
3432 log::warn!("Skipping unknown Spot execution type for {}", report.symbol);
3433 }
3434 BinanceSpotExecutionType::New => {
3435 if state.has_filled(&client_order_id) {
3436 log::debug!("Skipping New for already-filled {client_order_id}");
3437 return;
3438 }
3439
3440 let Some(price) =
3441 parse_spot_execution_report_price(report, &report.price, price_precision, "price")
3442 else {
3443 return;
3444 };
3445 let Some(quantity) = parse_spot_execution_report_quantity(
3446 report,
3447 &report.original_qty,
3448 size_precision,
3449 "original_qty",
3450 ) else {
3451 return;
3452 };
3453 let Some(stop_price) =
3454 parse_spot_execution_report_decimal(report, &report.stop_price, "stop_price")
3455 else {
3456 return;
3457 };
3458 let trigger = if stop_price > Decimal::ZERO {
3459 let Some(trigger_price) = parse_spot_execution_report_price(
3460 report,
3461 &report.stop_price,
3462 price_precision,
3463 "stop_price",
3464 ) else {
3465 return;
3466 };
3467 Some(trigger_price)
3468 } else {
3469 None
3470 };
3471 let changed = state.record_order_update(
3472 client_order_id,
3473 venue_order_id,
3474 quantity,
3475 price,
3476 trigger,
3477 );
3478
3479 if state.has_emitted_accepted(&client_order_id) {
3480 if !changed {
3481 return;
3482 }
3483 let updated = OrderUpdated::new(
3484 emitter.trader_id(),
3485 identity.strategy_id,
3486 identity.instrument_id,
3487 client_order_id,
3488 quantity,
3489 UUID4::new(),
3490 ts_event,
3491 ts_init,
3492 false,
3493 Some(venue_order_id),
3494 Some(account_id),
3495 Some(price),
3496 trigger,
3497 None, false, );
3500 emitter.send_order_event(OrderEventAny::Updated(updated));
3501 return;
3502 }
3503 state.insert_accepted(client_order_id);
3504 let accepted = OrderAccepted::new(
3505 emitter.trader_id(),
3506 identity.strategy_id,
3507 identity.instrument_id,
3508 client_order_id,
3509 venue_order_id,
3510 account_id,
3511 UUID4::new(),
3512 ts_event,
3513 ts_init,
3514 false,
3515 );
3516 emitter.send_order_event(OrderEventAny::Accepted(accepted));
3517 }
3518 BinanceSpotExecutionType::Trade => {
3519 let dedup_key = (report.symbol, report.trade_id);
3520 let is_duplicate = seen_trade_ids.lock().contains(&dedup_key);
3521
3522 if is_duplicate {
3523 log::debug!(
3524 "Duplicate trade_id={} for {}, skipping",
3525 report.trade_id,
3526 report.symbol
3527 );
3528 return;
3529 }
3530
3531 ensure_accepted_emitted(
3532 client_order_id,
3533 account_id,
3534 venue_order_id,
3535 identity,
3536 emitter,
3537 state,
3538 ts_init,
3539 );
3540
3541 let Some(last_qty) = parse_spot_execution_report_quantity(
3542 report,
3543 &report.last_filled_qty,
3544 size_precision,
3545 "last_filled_qty",
3546 ) else {
3547 return;
3548 };
3549 let Some(last_px) = parse_spot_execution_report_price(
3550 report,
3551 &report.last_filled_price,
3552 price_precision,
3553 "last_filled_price",
3554 ) else {
3555 return;
3556 };
3557 let Some(commission) =
3558 parse_spot_execution_report_decimal(report, &report.commission, "commission")
3559 else {
3560 return;
3561 };
3562 let commission_currency = report
3563 .commission_asset
3564 .as_ref()
3565 .map_or_else(Currency::USDT, |a| {
3566 Currency::get_or_create_crypto(a.as_str())
3567 });
3568 let commission_money = match Money::from_decimal(commission, commission_currency) {
3569 Ok(money) => money,
3570 Err(e) => {
3571 log::warn!(
3572 "Failed to build Spot commission money for symbol={}, order_id={}, \
3573 trade_id={}: {e}",
3574 report.symbol,
3575 report.order_id,
3576 report.trade_id,
3577 );
3578 return;
3579 }
3580 };
3581
3582 let liquidity_side = if report.is_maker {
3583 LiquiditySide::Maker
3584 } else {
3585 LiquiditySide::Taker
3586 };
3587
3588 let filled = OrderFilled::new(
3589 emitter.trader_id(),
3590 identity.strategy_id,
3591 instrument_id,
3592 client_order_id,
3593 venue_order_id,
3594 account_id,
3595 TradeId::new(report.trade_id.to_string()),
3596 identity.order_side,
3597 identity.order_type,
3598 last_qty,
3599 last_px,
3600 quote_currency,
3601 liquidity_side,
3602 UUID4::new(),
3603 ts_event,
3604 ts_init,
3605 false,
3606 None,
3607 Some(commission_money),
3608 None,
3609 );
3610
3611 state.insert_filled(client_order_id);
3612 emitter.send_order_event(OrderEventAny::Filled(filled));
3613 seen_trade_ids.lock().add(dedup_key);
3614
3615 let cumulative_qty = parse_spot_execution_report_decimal(
3616 report,
3617 &report.cumulative_filled_qty,
3618 "cumulative_filled_qty",
3619 );
3620 let original_qty =
3621 parse_spot_execution_report_decimal(report, &report.original_qty, "original_qty");
3622 if let (Some(original_qty), Some(cumulative_qty)) = (original_qty, cumulative_qty)
3623 && original_qty <= cumulative_qty
3624 {
3625 state.cleanup_terminal(client_order_id);
3626 }
3627 }
3628 BinanceSpotExecutionType::Replaced => {
3629 log::debug!(
3632 "Order replaced: client_order_id={client_order_id}, venue_order_id={venue_order_id}"
3633 );
3634 }
3635 BinanceSpotExecutionType::Canceled | BinanceSpotExecutionType::TradePrevention => {
3636 ensure_accepted_emitted(
3637 client_order_id,
3638 account_id,
3639 venue_order_id,
3640 identity,
3641 emitter,
3642 state,
3643 ts_init,
3644 );
3645 let canceled = OrderCanceled::new(
3646 emitter.trader_id(),
3647 identity.strategy_id,
3648 identity.instrument_id,
3649 client_order_id,
3650 UUID4::new(),
3651 ts_event,
3652 ts_init,
3653 false,
3654 Some(venue_order_id),
3655 Some(account_id),
3656 None,
3657 );
3658
3659 if state.defer_replace_cancel(canceled) {
3660 return;
3661 }
3662 state.cleanup_terminal(client_order_id);
3663 emitter.send_order_event(OrderEventAny::Canceled(canceled));
3664 }
3665 BinanceSpotExecutionType::Expired => {
3666 ensure_accepted_emitted(
3667 client_order_id,
3668 account_id,
3669 venue_order_id,
3670 identity,
3671 emitter,
3672 state,
3673 ts_init,
3674 );
3675 state.cleanup_terminal(client_order_id);
3676
3677 if treat_expired_as_canceled {
3678 let canceled = OrderCanceled::new(
3679 emitter.trader_id(),
3680 identity.strategy_id,
3681 identity.instrument_id,
3682 client_order_id,
3683 UUID4::new(),
3684 ts_event,
3685 ts_init,
3686 false,
3687 Some(venue_order_id),
3688 Some(account_id),
3689 None,
3690 );
3691 emitter.send_order_event(OrderEventAny::Canceled(canceled));
3692 } else {
3693 let expired = OrderExpired::new(
3694 emitter.trader_id(),
3695 identity.strategy_id,
3696 identity.instrument_id,
3697 client_order_id,
3698 UUID4::new(),
3699 ts_event,
3700 ts_init,
3701 false,
3702 Some(venue_order_id),
3703 Some(account_id),
3704 );
3705 emitter.send_order_event(OrderEventAny::Expired(expired));
3706 }
3707 }
3708 BinanceSpotExecutionType::Rejected => {
3709 let reason = if report.reject_reason.is_empty() {
3710 Ustr::from("Order rejected by venue")
3711 } else {
3712 Ustr::from(&report.reject_reason)
3713 };
3714 let due_post_only = report.time_in_force == BinanceTimeInForce::Gtx
3715 || (report.order_type == "LIMIT_MAKER"
3716 && (report.reject_reason.is_empty() || report.reject_reason == "NONE"));
3717 state.cleanup_terminal(client_order_id);
3718 emitter.emit_order_rejected_event(
3719 identity.strategy_id,
3720 identity.instrument_id,
3721 client_order_id,
3722 reason.as_str(),
3723 ts_init,
3724 due_post_only,
3725 );
3726 }
3727 }
3728}
3729
3730fn parse_spot_execution_report_quantity(
3731 report: &BinanceSpotExecutionReport,
3732 raw: &str,
3733 precision: u8,
3734 field: &str,
3735) -> Option<Quantity> {
3736 match parse_required_quantity_at_precision(raw, precision, field) {
3737 Ok(value) => Some(value),
3738 Err(e) => {
3739 warn_invalid_spot_execution_report_field(report, field, &e);
3740 None
3741 }
3742 }
3743}
3744
3745fn parse_spot_execution_report_price(
3746 report: &BinanceSpotExecutionReport,
3747 raw: &str,
3748 precision: u8,
3749 field: &str,
3750) -> Option<Price> {
3751 match parse_required_price_at_precision(raw, precision, field) {
3752 Ok(value) => Some(value),
3753 Err(e) => {
3754 warn_invalid_spot_execution_report_field(report, field, &e);
3755 None
3756 }
3757 }
3758}
3759
3760fn parse_spot_execution_report_decimal(
3761 report: &BinanceSpotExecutionReport,
3762 raw: &str,
3763 field: &str,
3764) -> Option<Decimal> {
3765 match parse_required_decimal(raw, field) {
3766 Ok(value) => Some(value),
3767 Err(e) => {
3768 warn_invalid_spot_execution_report_field(report, field, &e);
3769 None
3770 }
3771 }
3772}
3773
3774fn warn_invalid_spot_execution_report_field(
3775 report: &BinanceSpotExecutionReport,
3776 field: &str,
3777 error: &anyhow::Error,
3778) {
3779 log::warn!(
3780 "Failed to parse Spot execution report {field} for symbol={}, order_id={}, \
3781 trade_id={}, client_order_id={}: {error}",
3782 report.symbol,
3783 report.order_id,
3784 report.trade_id,
3785 report.client_order_id,
3786 );
3787}
3788
3789#[expect(clippy::too_many_arguments)]
3791fn dispatch_untracked_execution_report(
3792 report: &BinanceSpotExecutionReport,
3793 emitter: &ExecutionEventEmitter,
3794 _http_client: &BinanceSpotHttpClient,
3795 account_id: AccountId,
3796 treat_expired_as_canceled: bool,
3797 seen_trade_ids: &std::sync::Arc<Mutex<FifoCache<(Ustr, i64), 10_000>>>,
3798 instrument_id: InstrumentId,
3799 price_precision: u8,
3800 size_precision: u8,
3801 ts_init: UnixNanos,
3802) {
3803 match report.execution_type {
3804 BinanceSpotExecutionType::Unknown => {
3805 log::warn!("Skipping unknown Spot execution type for {}", report.symbol);
3806 }
3807 BinanceSpotExecutionType::Trade => {
3808 let dedup_key = (report.symbol, report.trade_id);
3809 let is_duplicate = seen_trade_ids.lock().contains(&dedup_key);
3810
3811 if is_duplicate {
3812 log::debug!(
3813 "Duplicate trade_id={} for {}, skipping",
3814 report.trade_id,
3815 report.symbol
3816 );
3817 return;
3818 }
3819
3820 match parse_spot_exec_report_to_order_status(
3821 report,
3822 instrument_id,
3823 price_precision,
3824 size_precision,
3825 account_id,
3826 treat_expired_as_canceled,
3827 ts_init,
3828 ) {
3829 Ok(status) => emitter.send_order_status_report(status),
3830 Err(e) => log::error!("Failed to parse order status report: {e}"),
3831 }
3832
3833 match parse_spot_exec_report_to_fill(
3834 report,
3835 instrument_id,
3836 price_precision,
3837 size_precision,
3838 account_id,
3839 ts_init,
3840 ) {
3841 Ok(fill) => {
3842 emitter.send_fill_report(fill);
3843 seen_trade_ids.lock().add(dedup_key);
3844 }
3845 Err(e) => log::error!("Failed to parse fill report: {e}"),
3846 }
3847 }
3848 BinanceSpotExecutionType::New
3849 | BinanceSpotExecutionType::Canceled
3850 | BinanceSpotExecutionType::Replaced
3851 | BinanceSpotExecutionType::Rejected
3852 | BinanceSpotExecutionType::Expired
3853 | BinanceSpotExecutionType::TradePrevention => {
3854 match parse_spot_exec_report_to_order_status(
3855 report,
3856 instrument_id,
3857 price_precision,
3858 size_precision,
3859 account_id,
3860 treat_expired_as_canceled,
3861 ts_init,
3862 ) {
3863 Ok(status) => emitter.send_order_status_report(status),
3864 Err(e) => log::error!("Failed to parse order status report: {e}"),
3865 }
3866 }
3867 }
3868}
3869
3870fn is_spot_post_only_rejection(error: &BinanceSpotHttpError) -> bool {
3872 match error {
3873 BinanceSpotHttpError::BinanceError { code, message, .. } => {
3874 *code == BINANCE_GTX_ORDER_REJECT_CODE
3875 || (*code == BINANCE_NEW_ORDER_REJECTED_CODE
3876 && message == BINANCE_SPOT_POST_ONLY_REJECT_MSG)
3877 }
3878 _ => false,
3879 }
3880}
3881
3882#[cfg(test)]
3883mod tests {
3884 use std::{
3885 cell::RefCell,
3886 rc::Rc,
3887 sync::{
3888 Arc,
3889 atomic::{AtomicUsize, Ordering},
3890 },
3891 };
3892
3893 use nautilus_common::{
3894 cache::Cache,
3895 clients::ExecutionClient,
3896 messages::{
3897 ExecutionEvent,
3898 execution::{CancelAllOrders, ExecutionReport},
3899 },
3900 };
3901 use nautilus_core::{UUID4, UnixNanos, time::get_atomic_clock_realtime};
3902 use nautilus_live::ExecutionClientCore;
3903 use nautilus_model::{
3904 enums::{AccountType, LiquiditySide, OmsType, OrderSide, OrderType, TimeInForce},
3905 events::OrderEventAny,
3906 identifiers::{AccountId, ClientOrderId, InstrumentId, StrategyId, TraderId, VenueOrderId},
3907 orders::{OrderTestBuilder, stubs::TestOrderEventStubs},
3908 types::{Price, Quantity},
3909 };
3910 use rstest::rstest;
3911
3912 use super::*;
3913 use crate::{
3914 common::{
3915 consts::{
3916 BINANCE_CLIENT_ID, BINANCE_STATUS_UNKNOWN_CODE, BINANCE_UNEXPECTED_RESPONSE_CODE,
3917 BINANCE_VENUE,
3918 },
3919 enums::BinanceEnvironment,
3920 },
3921 config::BinanceExecutionClientConfig,
3922 spot::{
3923 http::models::BinanceCancelOrderListOrder,
3924 sbe::spot::{
3925 contingency_type::ContingencyType as SbeContingencyType,
3926 order_side::OrderSide as SbeOrderSide, order_type::OrderType as SbeOrderType,
3927 self_trade_prevention_mode::SelfTradePreventionMode as SbeStp,
3928 time_in_force::TimeInForce as SbeTimeInForce,
3929 },
3930 },
3931 };
3932
3933 #[rstest]
3934 #[case::unsupported_type(OrderType::MarketToLimit, false, false, Some(OrderDeniedReason::UnsupportedOrderType { order_type: OrderType::MarketToLimit }))]
3935 #[case::reduce_only(
3936 OrderType::Market,
3937 true,
3938 false,
3939 Some(OrderDeniedReason::UnsupportedReduceOnly)
3940 )]
3941 #[case::limit_quote_quantity(OrderType::Limit, false, true, Some(OrderDeniedReason::ValidationFailed { detail: "Quote quantity requires a MARKET order on Binance Spot".to_string() }))]
3942 #[case::market_quote_quantity(OrderType::Market, false, true, None)]
3943 fn test_validate_order_fields(
3944 #[case] order_type: OrderType,
3945 #[case] reduce_only: bool,
3946 #[case] quote_quantity: bool,
3947 #[case] expected: Option<OrderDeniedReason>,
3948 ) {
3949 let order = OrderTestBuilder::new(order_type)
3950 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
3951 .quantity(Quantity::from("1"))
3952 .price(Price::from("100"))
3953 .reduce_only(reduce_only)
3954 .quote_quantity(quote_quantity)
3955 .build();
3956 assert_eq!(validate_order(&order, false).err(), expected);
3957 }
3958
3959 #[rstest]
3960 #[case::limit(
3961 false,
3962 Some(OrderDeniedReason::UnsupportedTimeInForce(TimeInForce::Day))
3963 )]
3964 #[case::maker(true, None)]
3965 fn test_validate_order_time_in_force_only_when_sent(
3966 #[case] post_only: bool,
3967 #[case] expected: Option<OrderDeniedReason>,
3968 ) {
3969 let order = OrderTestBuilder::new(OrderType::Limit)
3970 .instrument_id(InstrumentId::from("BTCUSDT.BINANCE"))
3971 .quantity(Quantity::from("1"))
3972 .price(Price::from("100"))
3973 .post_only(post_only)
3974 .time_in_force(TimeInForce::Day)
3975 .build();
3976 assert_eq!(validate_order(&order, false).err(), expected);
3977
3978 if post_only {
3979 let params =
3980 build_new_order_params(&order, order.client_order_id(), true, false, false)
3981 .unwrap();
3982 assert_eq!(params.time_in_force, None);
3983 }
3984 }
3985
3986 #[rstest]
3987 #[case::live(BinanceEnvironment::Live, BINANCE_SPOT_SBE_WS_API_URL)]
3988 #[case::testnet(BinanceEnvironment::Testnet, BINANCE_SPOT_SBE_WS_API_TESTNET_URL)]
3989 #[case::demo(BinanceEnvironment::Demo, BINANCE_SPOT_SBE_WS_API_DEMO_URL)]
3990 fn test_resolve_ws_trading_url_uses_environment_default(
3991 #[case] environment: BinanceEnvironment,
3992 #[case] expected: &str,
3993 ) {
3994 assert_eq!(
3995 BinanceSpotExecutionClient::resolve_ws_trading_url(None, environment),
3996 expected
3997 );
3998 }
3999
4000 #[rstest]
4001 fn test_resolve_ws_trading_url_preserves_override() {
4002 let expected = "wss://example.com/ws-api/v3";
4003
4004 assert_eq!(
4005 BinanceSpotExecutionClient::resolve_ws_trading_url(
4006 Some(expected.to_string()),
4007 BinanceEnvironment::Testnet,
4008 ),
4009 expected
4010 );
4011 }
4012
4013 #[rstest]
4014 fn test_dispatch_ws_trading_message_emits_cancel_rejected_and_clears_pending_request() {
4015 let tasks = TaskGroup::new();
4016 let task_spawner = tasks.spawner().expect("task spawner");
4017 let clock = get_atomic_clock_realtime();
4018 let (emitter, mut rx) = create_test_emitter(clock);
4019 let http_client = create_test_http_client(clock);
4020 let dispatch_state = create_tracked_dispatch_state(
4021 ClientOrderId::from("TEST"),
4022 InstrumentId::from("BTCUSDT.BINANCE"),
4023 );
4024 let ws_authenticated = tokio::sync::Notify::new();
4025 let ws_user_data_subscribed = tokio::sync::Notify::new();
4026 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
4027 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4028
4029 dispatch_state.pending_requests.insert(
4030 "req-cancel".to_string(),
4031 PendingRequest {
4032 client_order_id: ClientOrderId::from("TEST"),
4033 venue_order_id: Some(VenueOrderId::from("12345")),
4034 operation: PendingOperation::Cancel,
4035 },
4036 );
4037
4038 dispatch_ws_trading_message(
4039 BinanceSpotWsTradingMessage::CancelRejected {
4040 request_id: "req-cancel".to_string(),
4041 status: 400,
4042 code: -2011,
4043 msg: "Unknown order sent".to_string(),
4044 },
4045 &emitter,
4046 &http_client,
4047 AccountId::from("BINANCE-001"),
4048 false,
4049 clock,
4050 &dispatch_state,
4051 &ws_authenticated,
4052 &ws_user_data_subscribed,
4053 &ws_setup_error_tx,
4054 &seen_trade_ids,
4055 &task_spawner,
4056 );
4057
4058 assert!(dispatch_state.pending_requests.get("req-cancel").is_none());
4059
4060 match rx
4061 .try_recv()
4062 .expect("Cancel rejection event should be emitted")
4063 {
4064 ExecutionEvent::Order(OrderEventAny::CancelRejected(event)) => {
4065 assert_eq!(event.client_order_id, ClientOrderId::from("TEST"));
4066 assert_eq!(event.account_id, Some(AccountId::from("BINANCE-001")));
4067 assert!(event.reason.contains("code=-2011"));
4068 }
4069 other => panic!("Expected CancelRejected event, was {other:?}"),
4070 }
4071 }
4072
4073 #[rstest]
4074 #[case(
4075 BINANCE_UNEXPECTED_RESPONSE_CODE,
4076 "An unexpected response was received from the message bus"
4077 )]
4078 #[case(
4079 BINANCE_STATUS_UNKNOWN_CODE,
4080 "Timeout waiting for response from backend server"
4081 )]
4082 fn test_dispatch_ws_trading_message_unknown_status_keeps_order_registered(
4083 #[case] code: i64,
4084 #[case] msg: &str,
4085 ) {
4086 let tasks = TaskGroup::new();
4087 let task_spawner = tasks.spawner().expect("task spawner");
4088 let clock = get_atomic_clock_realtime();
4089 let (emitter, mut rx) = create_test_emitter(clock);
4090 let http_client = create_test_http_client(clock);
4091 let client_order_id = ClientOrderId::from("TEST");
4092 let dispatch_state =
4093 create_tracked_dispatch_state(client_order_id, InstrumentId::from("BTCUSDT.BINANCE"));
4094 let ws_authenticated = tokio::sync::Notify::new();
4095 let ws_user_data_subscribed = tokio::sync::Notify::new();
4096 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
4097 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4098
4099 dispatch_state.pending_requests.insert(
4100 "req-submit".to_string(),
4101 PendingRequest {
4102 client_order_id,
4103 venue_order_id: None,
4104 operation: PendingOperation::Place,
4105 },
4106 );
4107
4108 dispatch_ws_trading_message(
4109 BinanceSpotWsTradingMessage::OrderRejected {
4110 request_id: "req-submit".to_string(),
4111 status: 400,
4112 code: code as i32,
4113 msg: msg.to_string(),
4114 },
4115 &emitter,
4116 &http_client,
4117 AccountId::from("BINANCE-001"),
4118 false,
4119 clock,
4120 &dispatch_state,
4121 &ws_authenticated,
4122 &ws_user_data_subscribed,
4123 &ws_setup_error_tx,
4124 &seen_trade_ids,
4125 &task_spawner,
4126 );
4127
4128 assert!(dispatch_state.pending_requests.get("req-submit").is_none());
4129 assert!(
4130 dispatch_state
4131 .order_identities
4132 .get(&client_order_id)
4133 .is_some()
4134 );
4135 assert!(rx.try_recv().is_err());
4136 }
4137
4138 #[rstest]
4139 fn test_dispatch_ws_trading_message_definite_submit_rejection_emits_order_rejected() {
4140 let tasks = TaskGroup::new();
4141 let task_spawner = tasks.spawner().expect("task spawner");
4142 let clock = get_atomic_clock_realtime();
4143 let (emitter, mut rx) = create_test_emitter(clock);
4144 let http_client = create_test_http_client(clock);
4145 let client_order_id = ClientOrderId::from("TEST");
4146 let dispatch_state =
4147 create_tracked_dispatch_state(client_order_id, InstrumentId::from("BTCUSDT.BINANCE"));
4148 let ws_authenticated = tokio::sync::Notify::new();
4149 let ws_user_data_subscribed = tokio::sync::Notify::new();
4150 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
4151 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4152
4153 dispatch_state.pending_requests.insert(
4154 "req-submit".to_string(),
4155 PendingRequest {
4156 client_order_id,
4157 venue_order_id: None,
4158 operation: PendingOperation::Place,
4159 },
4160 );
4161
4162 dispatch_ws_trading_message(
4163 BinanceSpotWsTradingMessage::OrderRejected {
4164 request_id: "req-submit".to_string(),
4165 status: 400,
4166 code: BINANCE_NEW_ORDER_REJECTED_CODE as i32,
4167 msg: BINANCE_SPOT_POST_ONLY_REJECT_MSG.to_string(),
4168 },
4169 &emitter,
4170 &http_client,
4171 AccountId::from("BINANCE-001"),
4172 false,
4173 clock,
4174 &dispatch_state,
4175 &ws_authenticated,
4176 &ws_user_data_subscribed,
4177 &ws_setup_error_tx,
4178 &seen_trade_ids,
4179 &task_spawner,
4180 );
4181
4182 assert!(dispatch_state.pending_requests.get("req-submit").is_none());
4183 assert!(
4184 dispatch_state
4185 .order_identities
4186 .get(&client_order_id)
4187 .is_none()
4188 );
4189
4190 match rx
4191 .try_recv()
4192 .expect("OrderRejected event should be emitted")
4193 {
4194 ExecutionEvent::Order(OrderEventAny::Rejected(event)) => {
4195 assert_eq!(event.client_order_id, client_order_id);
4196 assert_eq!(event.account_id, AccountId::from("BINANCE-001"));
4197 assert!(event.reason.contains("code=-2010"));
4198 assert!(event.due_post_only);
4199 }
4200 other => panic!("Expected OrderRejected event, was {other:?}"),
4201 }
4202 }
4203
4204 #[rstest]
4205 fn test_dispatch_ws_trading_message_emits_modify_rejected_and_clears_pending_request() {
4206 let tasks = TaskGroup::new();
4207 let task_spawner = tasks.spawner().expect("task spawner");
4208 let clock = get_atomic_clock_realtime();
4209 let (emitter, mut rx) = create_test_emitter(clock);
4210 let http_client = create_test_http_client(clock);
4211 let dispatch_state = create_tracked_dispatch_state(
4212 ClientOrderId::from("TEST"),
4213 InstrumentId::from("BTCUSDT.BINANCE"),
4214 );
4215 let ws_authenticated = tokio::sync::Notify::new();
4216 let ws_user_data_subscribed = tokio::sync::Notify::new();
4217 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
4218 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4219
4220 dispatch_state.pending_requests.insert(
4221 "req-modify".to_string(),
4222 PendingRequest {
4223 client_order_id: ClientOrderId::from("TEST"),
4224 venue_order_id: Some(VenueOrderId::from("12345")),
4225 operation: PendingOperation::Modify,
4226 },
4227 );
4228
4229 dispatch_ws_trading_message(
4230 BinanceSpotWsTradingMessage::CancelReplaceRejected {
4231 request_id: "req-modify".to_string(),
4232 status: 400,
4233 code: -2021,
4234 msg: "Order cancel-replace partially failed".to_string(),
4235 },
4236 &emitter,
4237 &http_client,
4238 AccountId::from("BINANCE-001"),
4239 false,
4240 clock,
4241 &dispatch_state,
4242 &ws_authenticated,
4243 &ws_user_data_subscribed,
4244 &ws_setup_error_tx,
4245 &seen_trade_ids,
4246 &task_spawner,
4247 );
4248
4249 assert!(dispatch_state.pending_requests.get("req-modify").is_none());
4250
4251 match rx
4252 .try_recv()
4253 .expect("Modify rejection event should be emitted")
4254 {
4255 ExecutionEvent::Order(OrderEventAny::ModifyRejected(event)) => {
4256 assert_eq!(event.client_order_id, ClientOrderId::from("TEST"));
4257 assert_eq!(event.account_id, Some(AccountId::from("BINANCE-001")));
4258 assert!(event.reason.contains("code=-2021"));
4259 }
4260 other => panic!("Expected ModifyRejected event, was {other:?}"),
4261 }
4262 }
4263
4264 fn create_test_emitter(
4265 clock: &'static AtomicTime,
4266 ) -> (
4267 ExecutionEventEmitter,
4268 tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
4269 ) {
4270 let mut emitter = ExecutionEventEmitter::new(
4271 clock,
4272 TraderId::from("TESTER-001"),
4273 AccountId::from("BINANCE-001"),
4274 AccountType::Cash,
4275 None,
4276 );
4277 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
4278 emitter.set_sender(tx);
4279 (emitter, rx)
4280 }
4281
4282 fn create_test_http_client(clock: &'static AtomicTime) -> BinanceSpotHttpClient {
4283 let client = BinanceSpotHttpClient::new(
4284 BinanceEnvironment::Live,
4285 clock,
4286 None,
4287 None,
4288 None,
4289 None,
4290 None,
4291 None,
4292 )
4293 .expect("Test HTTP client should be created");
4294 client.cache_instruments(vec![
4295 nautilus_model::instruments::stubs::currency_pair_ethusdt().into(),
4296 nautilus_model::instruments::stubs::currency_pair_btcusdt().into(),
4297 ]);
4298 client
4299 }
4300
4301 fn create_tracked_dispatch_state(
4302 client_order_id: ClientOrderId,
4303 instrument_id: InstrumentId,
4304 ) -> WsDispatchState {
4305 let dispatch_state = WsDispatchState::default();
4306 dispatch_state.order_identities.insert(
4307 client_order_id,
4308 OrderIdentity {
4309 instrument_id,
4310 strategy_id: StrategyId::from("TEST-STRATEGY"),
4311 order_side: OrderSide::Buy,
4312 order_type: OrderType::Limit,
4313 price: None,
4314 quantity: Quantity::from("1"),
4315 venue_position_id: None,
4316 },
4317 );
4318 dispatch_state
4319 }
4320
4321 fn create_cancel_order_list_response(
4322 children: &[(ClientOrderId, i64)],
4323 ) -> BinanceCancelOpenOrdersResponse {
4324 let symbol = "BTCUSDT";
4325 let order_list_id = 44;
4326 let orders = children
4327 .iter()
4328 .map(|(client_order_id, order_id)| BinanceCancelOrderListOrder {
4329 symbol: symbol.to_string(),
4330 order_id: *order_id,
4331 client_order_id: encode_broker_id(client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID),
4332 })
4333 .collect();
4334 let order_reports = children
4335 .iter()
4336 .map(|(client_order_id, order_id)| BinanceCancelOrderResponse {
4337 price_exponent: -8,
4338 qty_exponent: -8,
4339 order_id: *order_id,
4340 order_list_id: Some(order_list_id),
4341 transact_time: 1_734_300_000_000_000,
4342 price_mantissa: 100_000_000_000,
4343 orig_qty_mantissa: 10_000_000,
4344 executed_qty_mantissa: 0,
4345 cummulative_quote_qty_mantissa: 0,
4346 status: SbeOrderStatus::Canceled,
4347 time_in_force: SbeTimeInForce::Gtc,
4348 order_type: SbeOrderType::Limit,
4349 side: SbeOrderSide::Sell,
4350 self_trade_prevention_mode: SbeStp::None,
4351 client_order_id: "cancel-list".to_string(),
4352 orig_client_order_id: encode_broker_id(
4353 client_order_id,
4354 BINANCE_NAUTILUS_SPOT_BROKER_ID,
4355 ),
4356 symbol: symbol.to_string(),
4357 })
4358 .collect();
4359 BinanceCancelOpenOrdersResponse::OrderList(BinanceCancelOrderListResponse {
4360 order_list_id,
4361 contingency_type: SbeContingencyType::Oco,
4362 list_status_type: SbeListStatusType::AllDone,
4363 list_order_status: SbeListOrderStatus::AllDone,
4364 transaction_time: 1_734_300_000_000_000,
4365 list_client_order_id: "list-44".to_string(),
4366 symbol: symbol.to_string(),
4367 orders,
4368 order_reports,
4369 })
4370 }
4371
4372 #[rstest]
4373 fn test_dispatch_order_list_canceled_emits_ordered_events_once_and_cleans_state() {
4374 let clock = get_atomic_clock_realtime();
4375 let (emitter, mut rx) = create_test_emitter(clock);
4376 let account_id = AccountId::from("BINANCE-001");
4377 let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
4378 let first = ClientOrderId::from("O-20200101-000000-001-001-1");
4379 let second = ClientOrderId::from("O-20200101-000000-002-002-2");
4380 let state = create_tracked_dispatch_state(first, instrument_id);
4381 state.order_identities.insert(
4382 second,
4383 OrderIdentity {
4384 instrument_id,
4385 strategy_id: StrategyId::from("TEST-STRATEGY"),
4386 order_side: OrderSide::Sell,
4387 order_type: OrderType::Limit,
4388 price: None,
4389 quantity: Quantity::from("1"),
4390 venue_position_id: None,
4391 },
4392 );
4393
4394 for (request_id, client_order_id) in [("cancel-1", first), ("cancel-2", second)] {
4395 state.pending_requests.insert(
4396 request_id.to_string(),
4397 PendingRequest {
4398 client_order_id,
4399 venue_order_id: None,
4400 operation: PendingOperation::Cancel,
4401 },
4402 );
4403 }
4404
4405 let prepared = prepare_cancel_all_orders(vec![create_cancel_order_list_response(&[
4406 (first, 111),
4407 (second, 222),
4408 ])])
4409 .unwrap();
4410
4411 for canceled_order in &prepared {
4412 dispatch_order_list_canceled(
4413 canceled_order,
4414 &emitter,
4415 account_id,
4416 &state,
4417 clock.get_time_ns(),
4418 );
4419 }
4420
4421 for canceled_order in &prepared {
4422 dispatch_order_list_canceled(
4423 canceled_order,
4424 &emitter,
4425 account_id,
4426 &state,
4427 clock.get_time_ns(),
4428 );
4429 }
4430
4431 let events: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()).collect();
4432 assert_eq!(events.len(), 4);
4433
4434 for (accepted_index, canceled_index, client_order_id) in [(0, 1, first), (2, 3, second)] {
4435 assert!(matches!(
4436 &events[accepted_index],
4437 ExecutionEvent::Order(OrderEventAny::Accepted(event))
4438 if event.client_order_id == client_order_id
4439 ));
4440 assert!(matches!(
4441 &events[canceled_index],
4442 ExecutionEvent::Order(OrderEventAny::Canceled(event))
4443 if event.client_order_id == client_order_id
4444 ));
4445 }
4446 assert!(state.order_identities.is_empty());
4447 assert!(state.pending_requests.is_empty());
4448 assert!(!state.has_emitted_accepted(&first));
4449 assert!(!state.has_emitted_accepted(&second));
4450 }
4451
4452 #[rstest]
4453 fn test_prepare_cancel_order_list_rejects_mismatched_child_without_events() {
4454 let clock = get_atomic_clock_realtime();
4455 let (_emitter, mut rx) = create_test_emitter(clock);
4456 let client_order_id = ClientOrderId::from("O-20200101-000000-001-001-1");
4457 let mut response = create_cancel_order_list_response(&[(client_order_id, 111)]);
4458 let BinanceCancelOpenOrdersResponse::OrderList(order_list) = &mut response else {
4459 unreachable!();
4460 };
4461 order_list.order_reports[0].orig_client_order_id = encode_broker_id(
4462 &ClientOrderId::from("OTHER"),
4463 BINANCE_NAUTILUS_SPOT_BROKER_ID,
4464 );
4465
4466 let error = prepare_cancel_all_orders(vec![response]).unwrap_err();
4467
4468 assert!(
4469 error
4470 .to_string()
4471 .contains("does not match its order identity")
4472 );
4473 assert!(rx.try_recv().is_err());
4474 }
4475
4476 #[rstest]
4477 fn test_prepare_cancel_order_list_rejects_non_canceled_child_without_state_change() {
4478 let clock = get_atomic_clock_realtime();
4479 let (_emitter, mut rx) = create_test_emitter(clock);
4480 let client_order_id = ClientOrderId::from("O-20200101-000000-001-001-1");
4481 let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
4482 let state = create_tracked_dispatch_state(client_order_id, instrument_id);
4483 let mut response = create_cancel_order_list_response(&[(client_order_id, 111)]);
4484 let BinanceCancelOpenOrdersResponse::OrderList(order_list) = &mut response else {
4485 unreachable!();
4486 };
4487 order_list.order_reports[0].status = SbeOrderStatus::Filled;
4488
4489 let error = prepare_cancel_all_orders(vec![response]).unwrap_err();
4490
4491 assert!(error.to_string().contains("reports status Filled"));
4492 assert!(state.order_identities.contains_key(&client_order_id));
4493 assert!(rx.try_recv().is_err());
4494 }
4495
4496 #[rstest]
4497 fn test_dispatch_order_list_canceled_concurrent_duplicate_emits_once() {
4498 let clock = get_atomic_clock_realtime();
4499 let (emitter, mut rx) = create_test_emitter(clock);
4500 let account_id = AccountId::from("BINANCE-001");
4501 let client_order_id = ClientOrderId::from("O-20200101-000000-001-001-1");
4502 let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
4503 let state = create_tracked_dispatch_state(client_order_id, instrument_id);
4504 let prepared = prepare_cancel_all_orders(vec![create_cancel_order_list_response(&[(
4505 client_order_id,
4506 111,
4507 )])])
4508 .unwrap();
4509
4510 std::thread::scope(|scope| {
4511 for _ in 0..2 {
4512 scope.spawn(|| {
4513 dispatch_order_list_canceled(
4514 &prepared[0],
4515 &emitter,
4516 account_id,
4517 &state,
4518 clock.get_time_ns(),
4519 );
4520 });
4521 }
4522 });
4523
4524 let events: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()).collect();
4525 assert_eq!(events.len(), 2);
4526 assert!(matches!(
4527 &events[0],
4528 ExecutionEvent::Order(OrderEventAny::Accepted(event))
4529 if event.client_order_id == client_order_id
4530 ));
4531 assert!(matches!(
4532 &events[1],
4533 ExecutionEvent::Order(OrderEventAny::Canceled(event))
4534 if event.client_order_id == client_order_id
4535 ));
4536 assert!(state.order_identities.is_empty());
4537 }
4538
4539 #[rstest]
4540 fn test_http_submit_success_defers_acceptance_to_user_stream() {
4541 let clock = get_atomic_clock_realtime();
4542 let (emitter, mut rx) = create_test_emitter(clock);
4543 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
4544 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
4545 let dispatch_state = Arc::new(create_tracked_dispatch_state(
4546 client_order_id,
4547 instrument_id,
4548 ));
4549 handle_spot_order_submit_success(client_order_id, VenueOrderId::from("12345678"));
4550
4551 assert!(!dispatch_state.has_emitted_accepted(&client_order_id));
4552 assert!(rx.try_recv().is_err());
4553 let new_json = crate::common::testing::load_fixture_string(
4554 "spot/user_data_json/execution_report_new.json",
4555 );
4556 let report: BinanceSpotExecutionReport = serde_json::from_str(&new_json).unwrap();
4557 let identity = dispatch_state
4558 .order_identities
4559 .get(&client_order_id)
4560 .unwrap()
4561 .clone();
4562 dispatch_tracked_execution_report(
4563 &report,
4564 &emitter,
4565 AccountId::from("BINANCE-001"),
4566 false,
4567 &dispatch_state,
4568 &Arc::new(Mutex::new(FifoCache::new())),
4569 client_order_id,
4570 &identity,
4571 instrument_id,
4572 2,
4573 8,
4574 Currency::USDT(),
4575 clock.get_time_ns(),
4576 );
4577
4578 assert!(matches!(
4579 rx.try_recv(),
4580 Ok(ExecutionEvent::Order(OrderEventAny::Accepted(event)))
4581 if event.client_order_id == client_order_id
4582 && event.venue_order_id == VenueOrderId::from("12345678")
4583 ));
4584 assert!(rx.try_recv().is_err());
4585 }
4586
4587 #[rstest]
4588 #[case::gtx(
4589 BinanceSpotHttpError::BinanceError {
4590 code: BINANCE_GTX_ORDER_REJECT_CODE,
4591 message: "Order would immediately trigger.".to_string(),
4592 status: 400,
4593 retry_after: None,
4594 },
4595 true,
4596 )]
4597 #[case::spot_post_only(
4598 BinanceSpotHttpError::BinanceError {
4599 code: BINANCE_NEW_ORDER_REJECTED_CODE,
4600 message: BINANCE_SPOT_POST_ONLY_REJECT_MSG.to_string(),
4601 status: 400,
4602 retry_after: None,
4603 },
4604 true,
4605 )]
4606 #[case::new_order_rejected_other_message(
4607 BinanceSpotHttpError::BinanceError {
4608 code: BINANCE_NEW_ORDER_REJECTED_CODE,
4609 message: "Insufficient balance.".to_string(),
4610 status: 400,
4611 retry_after: None,
4612 },
4613 false,
4614 )]
4615 #[case::unrelated_code(
4616 BinanceSpotHttpError::BinanceError {
4617 code: -2011,
4618 message: "Unknown order sent.".to_string(),
4619 status: 400,
4620 retry_after: None,
4621 },
4622 false,
4623 )]
4624 #[case::non_binance_error(
4625 BinanceSpotHttpError::NetworkError("connection reset".to_string()),
4626 false,
4627 )]
4628 fn test_is_spot_post_only_rejection(
4629 #[case] error: BinanceSpotHttpError,
4630 #[case] expected: bool,
4631 ) {
4632 assert_eq!(is_spot_post_only_rejection(&error), expected);
4633 }
4634
4635 #[rstest]
4636 fn test_dispatch_tracked_execution_report_trade_dedup() {
4637 let tasks = TaskGroup::new();
4638 let task_spawner = tasks.spawner().expect("task spawner");
4639 let clock = get_atomic_clock_realtime();
4640 let (emitter, mut rx) = create_test_emitter(clock);
4641 let http_client = create_test_http_client(clock);
4642 let client_order_id = ClientOrderId::from("x-TD67BGP9-T0000000000000");
4643 let dispatch_state = create_tracked_dispatch_state(
4644 ClientOrderId::from("O-20200101-000000-000-000-0"),
4645 InstrumentId::from("ETHUSDT.BINANCE"),
4646 );
4647 let ws_authenticated = tokio::sync::Notify::new();
4648 let ws_user_data_subscribed = tokio::sync::Notify::new();
4649 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
4650 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4651
4652 let trade_json = crate::common::testing::load_fixture_string(
4653 "spot/user_data_json/execution_report_trade.json",
4654 );
4655 let report: BinanceSpotExecutionReport = serde_json::from_str(&trade_json).unwrap();
4656
4657 dispatch_ws_trading_message(
4658 BinanceSpotWsTradingMessage::ExecutionReport(Box::new(report.clone())),
4659 &emitter,
4660 &http_client,
4661 AccountId::from("BINANCE-001"),
4662 false,
4663 clock,
4664 &dispatch_state,
4665 &ws_authenticated,
4666 &ws_user_data_subscribed,
4667 &ws_setup_error_tx,
4668 &seen_trade_ids,
4669 &task_spawner,
4670 );
4671 dispatch_ws_trading_message(
4672 BinanceSpotWsTradingMessage::ExecutionReport(Box::new(report)),
4673 &emitter,
4674 &http_client,
4675 AccountId::from("BINANCE-001"),
4676 false,
4677 clock,
4678 &dispatch_state,
4679 &ws_authenticated,
4680 &ws_user_data_subscribed,
4681 &ws_setup_error_tx,
4682 &seen_trade_ids,
4683 &task_spawner,
4684 );
4685
4686 let mut events = Vec::new();
4687 while let Ok(event) = rx.try_recv() {
4688 events.push(event);
4689 }
4690
4691 let fills: Vec<_> = events
4692 .iter()
4693 .filter(|e| matches!(e, ExecutionEvent::Order(OrderEventAny::Filled(_))))
4694 .collect();
4695 assert_eq!(fills.len(), 1, "duplicate trade should be deduped");
4696
4697 match fills[0] {
4698 ExecutionEvent::Order(OrderEventAny::Filled(fill)) => {
4699 assert_eq!(
4700 fill.client_order_id,
4701 ClientOrderId::from("O-20200101-000000-000-000-0"),
4702 );
4703 assert_eq!(fill.trade_id, TradeId::new("98765432"));
4704 assert_eq!(fill.liquidity_side, LiquiditySide::Maker);
4705 assert_eq!(fill.currency, Currency::USDT());
4706 assert_eq!(fill.last_px, Price::from("2500.00"));
4707 assert_eq!(fill.last_qty, Quantity::from("1.00000"));
4708 assert_eq!(
4709 fill.commission,
4710 Some(Money::from_decimal(Decimal::new(1, 3), Currency::ETH()).unwrap()),
4711 );
4712 }
4713 _ => unreachable!(),
4714 }
4715 let _ = client_order_id;
4716 }
4717
4718 #[rstest]
4719 fn test_dispatch_tracked_execution_report_invalid_fill_qty_skips_filled_event() {
4720 let tasks = TaskGroup::new();
4721 let task_spawner = tasks.spawner().expect("task spawner");
4722 let clock = get_atomic_clock_realtime();
4723 let (emitter, mut rx) = create_test_emitter(clock);
4724 let http_client = create_test_http_client(clock);
4725 let dispatch_state = create_tracked_dispatch_state(
4726 ClientOrderId::from("O-20200101-000000-000-000-0"),
4727 InstrumentId::from("ETHUSDT.BINANCE"),
4728 );
4729 let ws_authenticated = tokio::sync::Notify::new();
4730 let ws_user_data_subscribed = tokio::sync::Notify::new();
4731 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
4732 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4733
4734 let trade_json = crate::common::testing::load_fixture_string(
4735 "spot/user_data_json/execution_report_trade.json",
4736 );
4737 let mut report: BinanceSpotExecutionReport = serde_json::from_str(&trade_json).unwrap();
4738 report.last_filled_qty = "not-a-number".to_string();
4739
4740 dispatch_ws_trading_message(
4741 BinanceSpotWsTradingMessage::ExecutionReport(Box::new(report)),
4742 &emitter,
4743 &http_client,
4744 AccountId::from("BINANCE-001"),
4745 false,
4746 clock,
4747 &dispatch_state,
4748 &ws_authenticated,
4749 &ws_user_data_subscribed,
4750 &ws_setup_error_tx,
4751 &seen_trade_ids,
4752 &task_spawner,
4753 );
4754
4755 let mut events = Vec::new();
4756 while let Ok(event) = rx.try_recv() {
4757 events.push(event);
4758 }
4759
4760 assert!(
4761 events
4762 .iter()
4763 .all(|e| !matches!(e, ExecutionEvent::Order(OrderEventAny::Filled(_)))),
4764 "invalid fill quantity must not emit OrderFilled",
4765 );
4766 }
4767
4768 #[rstest]
4769 fn test_dispatch_execution_report_invalid_client_order_id_emits_nothing() {
4770 let clock = get_atomic_clock_realtime();
4771 let (emitter, mut rx) = create_test_emitter(clock);
4772 let http_client = create_test_http_client(clock);
4773 let dispatch_state = WsDispatchState::default();
4774 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4775 let json = crate::common::testing::load_fixture_string(
4776 "spot/user_data_json/execution_report_new.json",
4777 );
4778 let mut report: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
4779 report.client_order_id = "x-TD67BGP9-R".to_string();
4780
4781 dispatch_execution_report(
4782 &report,
4783 &emitter,
4784 &http_client,
4785 AccountId::from("BINANCE-001"),
4786 false,
4787 &dispatch_state,
4788 &seen_trade_ids,
4789 clock.get_time_ns(),
4790 );
4791
4792 assert!(rx.try_recv().is_err());
4793 assert!(dispatch_state.order_identities.is_empty());
4794 }
4795
4796 #[rstest]
4797 fn test_dispatch_execution_report_canceled_resolves_order_by_orig_client_order_id() {
4798 let clock = get_atomic_clock_realtime();
4799 let (emitter, mut rx) = create_test_emitter(clock);
4800 let http_client = create_test_http_client(clock);
4801 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
4802 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
4803 let dispatch_state = WsDispatchState::default();
4804 dispatch_state.insert_accepted(client_order_id);
4805 dispatch_state.order_identities.insert(
4806 client_order_id,
4807 OrderIdentity {
4808 instrument_id,
4809 strategy_id: StrategyId::from("TEST-STRATEGY"),
4810 order_side: OrderSide::Buy,
4811 order_type: OrderType::Limit,
4812 price: None,
4813 quantity: Quantity::from("1"),
4814 venue_position_id: None,
4815 },
4816 );
4817 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4818 let json = crate::common::testing::load_fixture_string(
4819 "spot/user_data_json/execution_report_canceled.json",
4820 );
4821 let mut report: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
4822 report.original_client_order_id = Some(report.client_order_id.clone());
4823 report.client_order_id = "web_9f8e7d6c5b4a".to_string();
4824
4825 dispatch_execution_report(
4826 &report,
4827 &emitter,
4828 &http_client,
4829 AccountId::from("BINANCE-001"),
4830 false,
4831 &dispatch_state,
4832 &seen_trade_ids,
4833 clock.get_time_ns(),
4834 );
4835
4836 match rx.try_recv().expect("OrderCanceled expected") {
4837 ExecutionEvent::Order(OrderEventAny::Canceled(event)) => {
4838 assert_eq!(event.client_order_id, client_order_id);
4839 }
4840 other => panic!("Expected OrderCanceled, was {other:?}"),
4841 }
4842 assert!(rx.try_recv().is_err());
4843 assert!(dispatch_state.order_identities.is_empty());
4844 }
4845
4846 #[rstest]
4847 fn test_build_cancel_replace_params_sends_cancel_new_client_order_id() {
4848 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
4849 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
4850 let mut builder = OrderTestBuilder::new(OrderType::Limit);
4851 let order = builder
4852 .instrument_id(instrument_id)
4853 .client_order_id(client_order_id)
4854 .side(OrderSide::Buy)
4855 .quantity(Quantity::from("1"))
4856 .price(Price::from("2500.00"))
4857 .build();
4858 let cmd = ModifyOrder::new(
4859 TraderId::from("TESTER-001"),
4860 None,
4861 StrategyId::from("TEST-STRATEGY"),
4862 instrument_id,
4863 client_order_id,
4864 Some(VenueOrderId::from("12345678")),
4865 Some(Quantity::from("2")),
4866 Some(Price::from("2400.00")),
4867 None,
4868 UUID4::new(),
4869 UnixNanos::default(),
4870 None,
4871 None,
4872 );
4873 let cancel_id = cancel_replace_cancel_id();
4874
4875 let params = build_cancel_replace_params(
4876 &cmd,
4877 &order,
4878 Quantity::from("2"),
4879 false,
4880 cancel_id.clone(),
4881 )
4882 .unwrap();
4883
4884 assert!(cancel_id.len() <= 36);
4885 assert_eq!(params.cancel_order_id, Some(12345678));
4886 assert_eq!(
4887 params.cancel_new_client_order_id.as_deref(),
4888 Some(cancel_id.as_str())
4889 );
4890 let json = serde_json::to_value(¶ms).unwrap();
4891 assert_eq!(json["cancelNewClientOrderId"], cancel_id);
4892 }
4893
4894 #[rstest]
4895 fn test_dispatch_execution_report_cancel_replace_skips_cancel_then_updates_on_new() {
4896 let clock = get_atomic_clock_realtime();
4897 let (emitter, mut rx) = create_test_emitter(clock);
4898 let http_client = create_test_http_client(clock);
4899 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
4900 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
4901 let dispatch_state = WsDispatchState::default();
4902 dispatch_state.insert_accepted(client_order_id);
4903 dispatch_state.order_identities.insert(
4904 client_order_id,
4905 OrderIdentity {
4906 instrument_id,
4907 strategy_id: StrategyId::from("TEST-STRATEGY"),
4908 order_side: OrderSide::Buy,
4909 order_type: OrderType::Limit,
4910 price: None,
4911 quantity: Quantity::from("1"),
4912 venue_position_id: None,
4913 },
4914 );
4915 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4916 let account_id = AccountId::from("BINANCE-001");
4917
4918 let json = crate::common::testing::load_fixture_string(
4919 "spot/user_data_json/execution_report_canceled.json",
4920 );
4921 let mut canceled: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
4922 canceled.original_client_order_id = Some(canceled.client_order_id.clone());
4923 let cancel_id = cancel_replace_cancel_id();
4924 dispatch_state.insert_cancel_replace(cancel_id.clone());
4925 canceled.client_order_id = cancel_id;
4926
4927 dispatch_execution_report(
4928 &canceled,
4929 &emitter,
4930 &http_client,
4931 account_id,
4932 false,
4933 &dispatch_state,
4934 &seen_trade_ids,
4935 clock.get_time_ns(),
4936 );
4937
4938 assert!(rx.try_recv().is_err(), "cancel half must not emit");
4939 assert!(
4940 dispatch_state
4941 .order_identities
4942 .contains_key(&client_order_id)
4943 );
4944
4945 let json = crate::common::testing::load_fixture_string(
4946 "spot/user_data_json/execution_report_new.json",
4947 );
4948 let mut new: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
4949 new.order_id = canceled.order_id + 1;
4950
4951 dispatch_execution_report(
4952 &new,
4953 &emitter,
4954 &http_client,
4955 account_id,
4956 false,
4957 &dispatch_state,
4958 &seen_trade_ids,
4959 clock.get_time_ns(),
4960 );
4961
4962 match rx.try_recv().expect("OrderUpdated expected") {
4963 ExecutionEvent::Order(OrderEventAny::Updated(event)) => {
4964 assert_eq!(event.client_order_id, client_order_id);
4965 assert_eq!(
4966 event.venue_order_id,
4967 Some(VenueOrderId::new(new.order_id.to_string()))
4968 );
4969 }
4970 other => panic!("Expected OrderUpdated, was {other:?}"),
4971 }
4972 assert!(rx.try_recv().is_err());
4973 }
4974
4975 fn cancel_replace_cancel_report(cancel_id: &str) -> BinanceSpotExecutionReport {
4976 let json = crate::common::testing::load_fixture_string(
4977 "spot/user_data_json/execution_report_canceled.json",
4978 );
4979 let mut report: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
4980 report.original_client_order_id = Some(report.client_order_id.clone());
4981 report.client_order_id = cancel_id.to_string();
4982 report
4983 }
4984
4985 fn reject_cancel_replace(
4986 request_id: &str,
4987 code: i32,
4988 emitter: &ExecutionEventEmitter,
4989 http_client: &BinanceSpotHttpClient,
4990 dispatch_state: &WsDispatchState,
4991 task_spawner: &TaskSpawner,
4992 clock: &'static AtomicTime,
4993 ) {
4994 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
4995 let ws_authenticated = tokio::sync::Notify::new();
4996 let ws_user_data_subscribed = tokio::sync::Notify::new();
4997 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
4998
4999 dispatch_ws_trading_message(
5000 BinanceSpotWsTradingMessage::CancelReplaceRejected {
5001 request_id: request_id.to_string(),
5002 status: 400,
5003 code,
5004 msg: "Order cancel-replace partially failed".to_string(),
5005 },
5006 emitter,
5007 http_client,
5008 AccountId::from("BINANCE-001"),
5009 false,
5010 clock,
5011 dispatch_state,
5012 &ws_authenticated,
5013 &ws_user_data_subscribed,
5014 &ws_setup_error_tx,
5015 &seen_trade_ids,
5016 task_spawner,
5017 );
5018 }
5019
5020 #[rstest]
5021 #[case::cancel_report_first(true)]
5022 #[case::rejection_first(false)]
5023 fn test_cancel_replace_rejected_after_confirmed_cancel_emits_canceled(
5024 #[case] cancel_report_first: bool,
5025 ) {
5026 let tasks = TaskGroup::new();
5027 let task_spawner = tasks.spawner().expect("task spawner");
5028 let clock = get_atomic_clock_realtime();
5029 let (emitter, mut rx) = create_test_emitter(clock);
5030 let http_client = create_test_http_client(clock);
5031 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
5032 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
5033 let dispatch_state = create_tracked_dispatch_state(client_order_id, instrument_id);
5034 dispatch_state.insert_accepted(client_order_id);
5035 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
5036 let account_id = AccountId::from("BINANCE-001");
5037
5038 let cancel_id = cancel_replace_cancel_id();
5039 dispatch_state.insert_cancel_replace(cancel_id.clone());
5040 dispatch_state.pending_requests.insert(
5041 "req-modify".to_string(),
5042 PendingRequest {
5043 client_order_id,
5044 venue_order_id: Some(VenueOrderId::from("12345678")),
5045 operation: PendingOperation::Modify,
5046 },
5047 );
5048 dispatch_state
5049 .cancel_replace_request_ids
5050 .insert("req-modify".to_string(), cancel_id.clone());
5051 dispatch_state.begin_replace(client_order_id, VenueOrderId::from("12345678"));
5052 let canceled = cancel_replace_cancel_report(&cancel_id);
5053
5054 if cancel_report_first {
5055 dispatch_execution_report(
5056 &canceled,
5057 &emitter,
5058 &http_client,
5059 account_id,
5060 false,
5061 &dispatch_state,
5062 &seen_trade_ids,
5063 clock.get_time_ns(),
5064 );
5065 assert!(rx.try_recv().is_err(), "cancel report must be withheld");
5066 }
5067
5068 reject_cancel_replace(
5069 "req-modify",
5070 -2021,
5071 &emitter,
5072 &http_client,
5073 &dispatch_state,
5074 &task_spawner,
5075 clock,
5076 );
5077
5078 match rx.try_recv().expect("OrderModifyRejected expected") {
5079 ExecutionEvent::Order(OrderEventAny::ModifyRejected(event)) => {
5080 assert_eq!(event.client_order_id, client_order_id);
5081 }
5082 other => panic!("Expected OrderModifyRejected, was {other:?}"),
5083 }
5084
5085 if !cancel_report_first {
5086 assert!(rx.try_recv().is_err(), "no cancel confirmed yet");
5087 dispatch_execution_report(
5088 &canceled,
5089 &emitter,
5090 &http_client,
5091 account_id,
5092 false,
5093 &dispatch_state,
5094 &seen_trade_ids,
5095 clock.get_time_ns(),
5096 );
5097 }
5098
5099 match rx.try_recv().expect("OrderCanceled expected") {
5100 ExecutionEvent::Order(OrderEventAny::Canceled(event)) => {
5101 assert_eq!(event.client_order_id, client_order_id);
5102 assert_eq!(event.venue_order_id, Some(VenueOrderId::from("12345678")));
5103 assert_eq!(
5104 event.ts_event,
5105 UnixNanos::from(1_709_654_402_000_000_000u64)
5106 );
5107 }
5108 other => panic!("Expected OrderCanceled, was {other:?}"),
5109 }
5110 assert!(rx.try_recv().is_err());
5111 assert!(dispatch_state.order_identities.is_empty());
5112 assert!(dispatch_state.pending_requests.is_empty());
5113 assert!(dispatch_state.cancel_replace_request_ids.is_empty());
5114 }
5115
5116 #[rstest]
5117 #[case::partial_failure(-2021, 409, true)]
5118 #[case::ambiguous(BINANCE_STATUS_UNKNOWN_CODE, 500, false)]
5119 fn test_http_modify_failure_replays_withheld_cancel_without_identity(
5120 #[case] code: i64,
5121 #[case] status: u16,
5122 #[case] definitive: bool,
5123 ) {
5124 let clock = get_atomic_clock_realtime();
5125 let (emitter, mut rx) = create_test_emitter(clock);
5126 let http_client = create_test_http_client(clock);
5127 let dispatch_state = WsDispatchState::default();
5128 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
5129 let account_id = AccountId::from("BINANCE-001");
5130 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
5131 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
5132
5133 let cancel_id = cancel_replace_cancel_id();
5134 dispatch_state.insert_cancel_replace(cancel_id.clone());
5135 let canceled = cancel_replace_cancel_report(&cancel_id);
5136
5137 dispatch_execution_report(
5138 &canceled,
5139 &emitter,
5140 &http_client,
5141 account_id,
5142 false,
5143 &dispatch_state,
5144 &seen_trade_ids,
5145 clock.get_time_ns(),
5146 );
5147 assert!(rx.try_recv().is_err(), "cancel report must be withheld");
5148
5149 let command = ModifyOrder::new(
5150 TraderId::from("TESTER-001"),
5151 None,
5152 StrategyId::from("TEST-STRATEGY"),
5153 instrument_id,
5154 client_order_id,
5155 Some(VenueOrderId::from("12345678")),
5156 Some(Quantity::from("2")),
5157 Some(Price::from("2400.00")),
5158 None,
5159 UUID4::new(),
5160 UnixNanos::default(),
5161 None,
5162 None,
5163 );
5164 let error = anyhow::anyhow!(BinanceSpotHttpError::BinanceError {
5165 code,
5166 message: "Order cancel-replace partially failed.".to_string(),
5167 status,
5168 retry_after: None,
5169 });
5170
5171 handle_http_modify_failure(
5172 &error,
5173 &command,
5174 &cancel_id,
5175 &emitter,
5176 &dispatch_state,
5177 account_id,
5178 clock.get_time_ns(),
5179 );
5180
5181 if !definitive {
5182 assert!(
5183 rx.try_recv().is_err(),
5184 "ambiguous outcome must remain pending"
5185 );
5186 return;
5187 }
5188
5189 match rx.try_recv().expect("OrderModifyRejected expected") {
5190 ExecutionEvent::Order(OrderEventAny::ModifyRejected(event)) => {
5191 assert_eq!(event.client_order_id, client_order_id);
5192 }
5193 other => panic!("Expected OrderModifyRejected, was {other:?}"),
5194 }
5195
5196 match rx.try_recv().expect("cancel status report expected") {
5197 ExecutionEvent::Report(ExecutionReport::Order(report)) => {
5198 assert_eq!(report.client_order_id, Some(client_order_id));
5199 assert_eq!(report.venue_order_id, VenueOrderId::from("12345678"));
5200 assert_eq!(report.order_status, OrderStatus::Canceled);
5201 assert_eq!(
5202 report.ts_last,
5203 UnixNanos::from(1_709_654_402_000_000_000u64)
5204 );
5205 }
5206 other => panic!("Expected order status report, was {other:?}"),
5207 }
5208 assert!(rx.try_recv().is_err());
5209 }
5210
5211 #[rstest]
5212 #[case::unexpected_response(BINANCE_UNEXPECTED_RESPONSE_CODE)]
5213 #[case::status_unknown(BINANCE_STATUS_UNKNOWN_CODE)]
5214 fn test_cancel_replace_ambiguous_rejection_keeps_withholding_cancel(#[case] code: i64) {
5215 let tasks = TaskGroup::new();
5216 let task_spawner = tasks.spawner().expect("task spawner");
5217 let clock = get_atomic_clock_realtime();
5218 let (emitter, mut rx) = create_test_emitter(clock);
5219 let http_client = create_test_http_client(clock);
5220 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
5221 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
5222 let dispatch_state = create_tracked_dispatch_state(client_order_id, instrument_id);
5223 dispatch_state.insert_accepted(client_order_id);
5224 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
5225 let account_id = AccountId::from("BINANCE-001");
5226
5227 let cancel_id = cancel_replace_cancel_id();
5228 dispatch_state.insert_cancel_replace(cancel_id.clone());
5229 dispatch_state.pending_requests.insert(
5230 "req-modify".to_string(),
5231 PendingRequest {
5232 client_order_id,
5233 venue_order_id: Some(VenueOrderId::from("12345678")),
5234 operation: PendingOperation::Modify,
5235 },
5236 );
5237 dispatch_state
5238 .cancel_replace_request_ids
5239 .insert("req-modify".to_string(), cancel_id.clone());
5240 dispatch_state.begin_replace(client_order_id, VenueOrderId::from("12345678"));
5241 let canceled = cancel_replace_cancel_report(&cancel_id);
5242
5243 dispatch_execution_report(
5244 &canceled,
5245 &emitter,
5246 &http_client,
5247 account_id,
5248 false,
5249 &dispatch_state,
5250 &seen_trade_ids,
5251 clock.get_time_ns(),
5252 );
5253 reject_cancel_replace(
5254 "req-modify",
5255 i32::try_from(code).unwrap(),
5256 &emitter,
5257 &http_client,
5258 &dispatch_state,
5259 &task_spawner,
5260 clock,
5261 );
5262
5263 assert!(
5264 rx.try_recv().is_err(),
5265 "ambiguous outcome must remain pending"
5266 );
5267
5268 let json = crate::common::testing::load_fixture_string(
5270 "spot/user_data_json/execution_report_new.json",
5271 );
5272 let mut new: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
5273 new.order_id = canceled.order_id + 1;
5274
5275 dispatch_execution_report(
5276 &new,
5277 &emitter,
5278 &http_client,
5279 account_id,
5280 false,
5281 &dispatch_state,
5282 &seen_trade_ids,
5283 clock.get_time_ns(),
5284 );
5285
5286 match rx.try_recv().expect("OrderUpdated expected") {
5287 ExecutionEvent::Order(OrderEventAny::Updated(event)) => {
5288 assert_eq!(event.client_order_id, client_order_id);
5289 assert_eq!(
5290 event.venue_order_id,
5291 Some(VenueOrderId::new(new.order_id.to_string()))
5292 );
5293 }
5294 other => panic!("Expected OrderUpdated, was {other:?}"),
5295 }
5296 assert!(rx.try_recv().is_err());
5297 }
5298
5299 #[rstest]
5300 #[case::as_expired(false, OrderStatus::Expired)]
5301 #[case::as_canceled(true, OrderStatus::Canceled)]
5302 fn test_normalize_spot_order_status_report_expired_respects_config(
5303 #[case] treat_expired_as_canceled: bool,
5304 #[case] expected: OrderStatus,
5305 ) {
5306 let clock = get_atomic_clock_realtime();
5307 let json = crate::common::testing::load_fixture_string(
5308 "spot/user_data_json/execution_report_expired.json",
5309 );
5310 let msg: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
5311 let mut report = parse_spot_exec_report_to_order_status(
5312 &msg,
5313 InstrumentId::from("ETHUSDT.BINANCE"),
5314 2,
5315 5,
5316 AccountId::from("BINANCE-001"),
5317 false,
5318 clock.get_time_ns(),
5319 )
5320 .unwrap();
5321 let mut reports = vec![report.clone()];
5322
5323 normalize_spot_order_status_report(&mut report, treat_expired_as_canceled);
5324 normalize_spot_order_status_reports(&mut reports, treat_expired_as_canceled);
5325
5326 assert_eq!(report.order_status, expected);
5327 assert_eq!(reports[0].order_status, expected);
5328 }
5329
5330 #[rstest]
5331 #[case::tracked(true)]
5332 #[case::external(false)]
5333 fn test_cancel_identifier_preserves_external_orders(#[case] tracked: bool) {
5334 let client_order_id = ClientOrderId::from("cancel-target");
5335 let cmd = CancelOrder::new(
5336 TraderId::from("TESTER-001"),
5337 None,
5338 StrategyId::from("TEST-STRATEGY"),
5339 InstrumentId::from("ETHUSDT.BINANCE"),
5340 client_order_id,
5341 Some(VenueOrderId::from("12345")),
5342 UUID4::new(),
5343 UnixNanos::default(),
5344 None,
5345 None,
5346 );
5347
5348 let params = build_cancel_order_params(&cmd, tracked);
5349
5350 assert_eq!(params.symbol, "ETHUSDT");
5351 assert_eq!(params.order_id, (!tracked).then_some(12345));
5352 assert_eq!(
5353 params.orig_client_order_id,
5354 tracked.then(|| encode_broker_id(&client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID))
5355 );
5356 assert_eq!(params.new_client_order_id, None);
5357 }
5358
5359 #[rstest]
5360 #[case::stream_first(false, false)]
5361 #[case::http_first(true, false)]
5362 #[case::replacement_rejected(false, true)]
5363 fn test_cancel_replace_preserves_logical_order(
5364 #[case] http_first: bool,
5365 #[case] rejected: bool,
5366 ) {
5367 let clock = get_atomic_clock_realtime();
5368 let (emitter, mut rx) = create_test_emitter(clock);
5369 let client = create_test_http_client(clock);
5370 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
5371 let state =
5372 create_tracked_dispatch_state(client_order_id, InstrumentId::from("ETHUSDT.BINANCE"));
5373 let seen = Arc::new(Mutex::new(FifoCache::new()));
5374 let json = crate::common::testing::load_fixture_string(
5375 "spot/user_data_json/execution_report_new.json",
5376 );
5377 let mut original: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
5378 original.client_order_id =
5379 encode_broker_id(&client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
5380 let dispatch = |report: &BinanceSpotExecutionReport| {
5381 dispatch_execution_report(
5382 report,
5383 &emitter,
5384 &client,
5385 AccountId::from("BINANCE-001"),
5386 false,
5387 &state,
5388 &seen,
5389 clock.get_time_ns(),
5390 );
5391 };
5392 dispatch(&original);
5393 let old_id = VenueOrderId::from(original.order_id.to_string());
5394 let new_id = VenueOrderId::from((original.order_id + 1).to_string());
5395 state.begin_replace(client_order_id, old_id);
5396 let mut replacement = original.clone();
5397 replacement.order_id += 1;
5398 replacement.price = "2501.00".to_string();
5399 replacement.original_qty = "2.00000".to_string();
5400
5401 if http_first {
5402 assert!(state.record_order_update(
5403 client_order_id,
5404 new_id,
5405 Quantity::from("2.00000"),
5406 Price::from("2501.00"),
5407 None
5408 ));
5409 }
5410 let mut canceled = original.clone();
5411 canceled.execution_type = BinanceSpotExecutionType::Canceled;
5412 canceled.original_client_order_id = Some(original.client_order_id);
5413 canceled.client_order_id = "venue-cancel-request".to_string();
5414 dispatch(&canceled);
5415
5416 if rejected {
5417 let event = state
5418 .reject_replace(client_order_id)
5419 .expect("Cancellation must be retained");
5420 assert_eq!(event.venue_order_id, Some(old_id));
5421 state.cleanup_terminal(client_order_id);
5422 emitter.send_order_event(OrderEventAny::Canceled(event));
5423 } else {
5424 dispatch(&replacement);
5425 dispatch(&replacement);
5426 dispatch(&canceled);
5427 assert!(state.order_identities.contains_key(&client_order_id));
5428 assert!(!state.record_order_update(
5429 client_order_id,
5430 new_id,
5431 Quantity::from("2.00000"),
5432 Price::from("2501.00"),
5433 None
5434 ));
5435 }
5436 assert!(matches!(
5437 rx.try_recv().unwrap(),
5438 ExecutionEvent::Order(OrderEventAny::Accepted(_))
5439 ));
5440
5441 if rejected {
5442 assert!(matches!(
5443 rx.try_recv().unwrap(),
5444 ExecutionEvent::Order(OrderEventAny::Canceled(_))
5445 ));
5446 } else if !http_first {
5447 let ExecutionEvent::Order(OrderEventAny::Updated(event)) = rx.try_recv().unwrap()
5448 else {
5449 panic!("Expected replacement update");
5450 };
5451 assert_eq!(event.client_order_id, client_order_id);
5452 assert_eq!(event.venue_order_id, Some(new_id));
5453 assert_eq!(event.quantity, Quantity::from("2.00000"));
5454 assert_eq!(event.price, Some(Price::from("2501.00")));
5455 }
5456 assert!(rx.try_recv().is_err());
5457 }
5458
5459 #[rstest]
5460 #[case::duplicate("duplicate")]
5461 #[case::quantity("quantity")]
5462 #[case::price("price")]
5463 #[case::trigger("trigger")]
5464 #[case::replacement("replacement")]
5465 fn test_dispatch_new_emits_only_changed_order_terms(#[case] change: &str) {
5466 let clock = get_atomic_clock_realtime();
5467 let (emitter, mut rx) = create_test_emitter(clock);
5468 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
5469 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
5470 let state = create_tracked_dispatch_state(client_order_id, instrument_id);
5471 let identity = state
5472 .order_identities
5473 .get(&client_order_id)
5474 .unwrap()
5475 .clone();
5476 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
5477 let json = crate::common::testing::load_fixture_string(
5478 "spot/user_data_json/execution_report_new.json",
5479 );
5480 let mut report: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
5481 let dispatch = |report: &BinanceSpotExecutionReport| {
5482 dispatch_tracked_execution_report(
5483 report,
5484 &emitter,
5485 AccountId::from("BINANCE-001"),
5486 false,
5487 &state,
5488 &seen_trade_ids,
5489 client_order_id,
5490 &identity,
5491 instrument_id,
5492 2,
5493 5,
5494 Currency::USDT(),
5495 clock.get_time_ns(),
5496 );
5497 };
5498 dispatch(&report);
5499 match change {
5500 "quantity" => report.original_qty = "2.00000".to_string(),
5501 "price" => report.price = "2501.00".to_string(),
5502 "trigger" => report.stop_price = "2499.00".to_string(),
5503 "replacement" => report.order_id += 1,
5504 "duplicate" => {}
5505 _ => unreachable!(),
5506 }
5507 dispatch(&report);
5508 dispatch(&report);
5509
5510 assert!(matches!(
5511 rx.try_recv().unwrap(),
5512 ExecutionEvent::Order(OrderEventAny::Accepted(_))
5513 ));
5514
5515 if change != "duplicate" {
5516 let ExecutionEvent::Order(OrderEventAny::Updated(updated)) = rx.try_recv().unwrap()
5517 else {
5518 panic!("expected changed order update");
5519 };
5520 assert_eq!(updated.client_order_id, client_order_id);
5521 assert_eq!(
5522 updated.venue_order_id,
5523 Some(VenueOrderId::from(report.order_id.to_string()))
5524 );
5525 assert_eq!(
5526 updated.quantity,
5527 Quantity::from(report.original_qty.as_str())
5528 );
5529 assert_eq!(updated.price, Some(Price::from(report.price.as_str())));
5530 assert_eq!(
5531 updated.trigger_price,
5532 (change == "trigger").then(|| Price::from("2499.00"))
5533 );
5534 }
5535 assert!(rx.try_recv().is_err());
5536 }
5537
5538 #[rstest]
5539 #[case::as_expired(false)]
5540 #[case::as_canceled(true)]
5541 fn test_dispatch_tracked_execution_report_expired_respects_config(
5542 #[case] treat_expired_as_canceled: bool,
5543 ) {
5544 let clock = get_atomic_clock_realtime();
5545 let (emitter, mut rx) = create_test_emitter(clock);
5546 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-0");
5547 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
5548 let dispatch_state = WsDispatchState::default();
5549 dispatch_state.insert_accepted(client_order_id);
5550 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
5551 let identity = OrderIdentity {
5552 instrument_id,
5553 strategy_id: StrategyId::from("TEST-STRATEGY"),
5554 order_side: OrderSide::Buy,
5555 order_type: OrderType::Limit,
5556 price: None,
5557 quantity: Quantity::from("1"),
5558 venue_position_id: None,
5559 };
5560
5561 let json = crate::common::testing::load_fixture_string(
5562 "spot/user_data_json/execution_report_expired.json",
5563 );
5564 let report: BinanceSpotExecutionReport = serde_json::from_str(&json).unwrap();
5565
5566 dispatch_tracked_execution_report(
5567 &report,
5568 &emitter,
5569 AccountId::from("BINANCE-001"),
5570 treat_expired_as_canceled,
5571 &dispatch_state,
5572 &seen_trade_ids,
5573 client_order_id,
5574 &identity,
5575 instrument_id,
5576 2,
5577 5,
5578 Currency::USDT(),
5579 clock.get_time_ns(),
5580 );
5581
5582 let event = rx.try_recv().expect("terminal order event expected");
5583 match (treat_expired_as_canceled, event) {
5584 (true, ExecutionEvent::Order(OrderEventAny::Canceled(event))) => {
5585 assert_eq!(event.client_order_id, client_order_id);
5586 }
5587 (false, ExecutionEvent::Order(OrderEventAny::Expired(event))) => {
5588 assert_eq!(event.client_order_id, client_order_id);
5589 }
5590 (_, other) => panic!("Expected terminal expired/canceled event, was {other:?}"),
5591 }
5592 assert!(rx.try_recv().is_err());
5593 }
5594
5595 #[rstest]
5596 fn test_dispatch_tracked_execution_report_rejected_gtx_sets_post_only() {
5597 let tasks = TaskGroup::new();
5598 let task_spawner = tasks.spawner().expect("task spawner");
5599 let clock = get_atomic_clock_realtime();
5600 let (emitter, mut rx) = create_test_emitter(clock);
5601 let http_client = create_test_http_client(clock);
5602 let client_order_id = ClientOrderId::from("O-20200101-000000-000-000-1");
5603 let dispatch_state =
5604 create_tracked_dispatch_state(client_order_id, InstrumentId::from("ETHUSDT.BINANCE"));
5605 let ws_authenticated = tokio::sync::Notify::new();
5606 let ws_user_data_subscribed = tokio::sync::Notify::new();
5607 let (ws_setup_error_tx, _ws_setup_error_rx) = tokio::sync::mpsc::unbounded_channel();
5608 let seen_trade_ids = Arc::new(Mutex::new(FifoCache::new()));
5609
5610 let encoded = encode_broker_id(&client_order_id, BINANCE_NAUTILUS_SPOT_BROKER_ID);
5611 let report_json = format!(
5612 r#"{{
5613 "e":"executionReport","E":1709654400000,"s":"ETHUSDT",
5614 "c":"{encoded}","S":"BUY","o":"LIMIT","f":"GTX",
5615 "q":"1.00000000","p":"2500.00000000","P":"0.00000000",
5616 "x":"REJECTED","X":"REJECTED","r":"NONE","i":12345678,
5617 "l":"0.00000000","z":"0.00000000","L":"0.00000000",
5618 "n":"0","N":null,"T":1709654400000,"t":-1,"w":false,"m":false,
5619 "O":1709654400000,"Z":"0.00000000","C":""
5620 }}"#,
5621 );
5622 let report: BinanceSpotExecutionReport = serde_json::from_str(&report_json).unwrap();
5623
5624 dispatch_ws_trading_message(
5625 BinanceSpotWsTradingMessage::ExecutionReport(Box::new(report)),
5626 &emitter,
5627 &http_client,
5628 AccountId::from("BINANCE-001"),
5629 false,
5630 clock,
5631 &dispatch_state,
5632 &ws_authenticated,
5633 &ws_user_data_subscribed,
5634 &ws_setup_error_tx,
5635 &seen_trade_ids,
5636 &task_spawner,
5637 );
5638
5639 match rx.try_recv().expect("OrderRejected event expected") {
5640 ExecutionEvent::Order(OrderEventAny::Rejected(event)) => {
5641 assert_eq!(event.client_order_id, client_order_id);
5642 assert_eq!(event.account_id, AccountId::from("BINANCE-001"));
5643 assert!(event.due_post_only);
5644 }
5645 other => panic!("Expected OrderRejected event, was {other:?}"),
5646 }
5647 }
5648
5649 fn test_execution_client(
5650 base_url_http: String,
5651 ) -> (BinanceSpotExecutionClient, Rc<RefCell<Cache>>) {
5652 let cache = Rc::new(RefCell::new(Cache::default()));
5653 let core = ExecutionClientCore::new(
5654 TraderId::from("TESTER-001"),
5655 *BINANCE_CLIENT_ID,
5656 *BINANCE_VENUE,
5657 OmsType::Hedging,
5658 AccountId::from("BINANCE-001"),
5659 AccountType::Cash,
5660 None,
5661 cache.clone(),
5662 );
5663 let config = BinanceExecutionClientConfig {
5664 base_url_http: Some(base_url_http),
5665 use_ws_trading: false,
5666 api_key: Some("test_api_key".into()),
5667 api_secret: Some("test_api_secret".into()),
5668 ..Default::default()
5669 };
5670
5671 (
5672 BinanceSpotExecutionClient::new(core, config).unwrap(),
5673 cache,
5674 )
5675 }
5676
5677 async fn wait_for_spawned_tasks(client: &BinanceSpotExecutionClient) {
5678 for _ in 0..40 {
5679 if client.pending_tasks.all_finished() {
5680 return;
5681 }
5682
5683 tokio::time::sleep(Duration::from_millis(25)).await;
5684 }
5685
5686 panic!("timed out waiting for spawned Binance Spot execution tasks");
5687 }
5688
5689 struct MockVenueHits {
5690 single_cancel: Arc<AtomicUsize>,
5691 batch: Arc<AtomicUsize>,
5692 cancel_all: Arc<AtomicUsize>,
5693 }
5694
5695 async fn start_cancel_reject_server() -> (String, MockVenueHits) {
5696 let single_hits = Arc::new(AtomicUsize::new(0));
5697 let single_hits_clone = single_hits.clone();
5698 let batch_hits = Arc::new(AtomicUsize::new(0));
5699 let batch_hits_clone = batch_hits.clone();
5700 let cancel_all_hits = Arc::new(AtomicUsize::new(0));
5701 let cancel_all_hits_clone = cancel_all_hits.clone();
5702 let app = axum::Router::new()
5703 .route(
5704 "/api/v3/order",
5705 axum::routing::delete(move || {
5706 let hits = single_hits_clone.clone();
5707 async move {
5708 hits.fetch_add(1, Ordering::Relaxed);
5709
5710 (
5711 axum::http::StatusCode::BAD_REQUEST,
5712 axum::Json(
5713 serde_json::json!({"code": -2011, "msg": "Unknown order sent"}),
5714 ),
5715 )
5716 }
5717 }),
5718 )
5719 .route(
5720 "/api/v3/batchOrders",
5721 axum::routing::delete(move || {
5722 let hits = batch_hits_clone.clone();
5723 async move {
5724 hits.fetch_add(1, Ordering::Relaxed);
5725
5726 axum::Json(serde_json::json!([]))
5727 }
5728 }),
5729 )
5730 .route(
5731 "/api/v3/openOrders",
5732 axum::routing::delete(move || {
5733 let hits = cancel_all_hits_clone.clone();
5734 async move {
5735 hits.fetch_add(1, Ordering::Relaxed);
5736
5737 axum::Json(serde_json::json!([]))
5738 }
5739 }),
5740 );
5741 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
5742 let addr = listener.local_addr().unwrap();
5743
5744 tokio::spawn(async move {
5745 axum::serve(listener, app.into_make_service())
5746 .await
5747 .unwrap();
5748 });
5749
5750 (
5751 format!("http://{addr}"),
5752 MockVenueHits {
5753 single_cancel: single_hits,
5754 batch: batch_hits,
5755 cancel_all: cancel_all_hits,
5756 },
5757 )
5758 }
5759
5760 #[rstest]
5761 #[case::buy(Some(OrderSide::Buy), vec![0, 2])]
5762 #[case::sell(Some(OrderSide::Sell), vec![1, 3])]
5763 #[tokio::test]
5764 async fn test_cancel_all_orders_filters_by_side_and_preserves_owners(
5765 #[case] order_side: Option<OrderSide>,
5766 #[case] expected_indices: Vec<usize>,
5767 ) {
5768 let (base_url, hits) = start_cancel_reject_server().await;
5769 let (mut client, cache) = test_execution_client(base_url);
5770 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
5771 client.emitter.set_sender(tx);
5772
5773 let instrument_id = InstrumentId::from("BTCUSDT.BINANCE");
5774 let other_instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
5775 let account_id = client.core.account_id;
5776 let mut orders = Vec::new();
5777
5778 for (index, (instrument, owner, side, open)) in [
5779 (instrument_id, "S-001", OrderSide::Buy, true),
5780 (instrument_id, "S-001", OrderSide::Sell, true),
5781 (instrument_id, "S-002", OrderSide::Buy, true),
5782 (instrument_id, "S-002", OrderSide::Sell, true),
5783 (other_instrument_id, "S-001", OrderSide::Buy, true),
5784 (other_instrument_id, "S-002", OrderSide::Sell, true),
5785 (instrument_id, "S-001", OrderSide::Buy, false),
5786 (instrument_id, "S-002", OrderSide::Sell, false),
5787 ]
5788 .into_iter()
5789 .enumerate()
5790 {
5791 let client_order_id = ClientOrderId::from(format!("O-CANCEL-ALL-{index}"));
5792 let mut builder = OrderTestBuilder::new(OrderType::Limit);
5793 let order = builder
5794 .instrument_id(instrument)
5795 .client_order_id(client_order_id)
5796 .strategy_id(StrategyId::from(owner))
5797 .side(side)
5798 .quantity(Quantity::from("1"))
5799 .price(Price::from("10000.00"))
5800 .build();
5801 let venue_order_id = VenueOrderId::from(format!("{}", 1000 + index));
5802 let accepted = TestOrderEventStubs::accepted(&order, account_id, venue_order_id);
5803 cache
5804 .borrow_mut()
5805 .add_order(order, None, Some(*BINANCE_CLIENT_ID), false)
5806 .unwrap();
5807 let order = cache.borrow_mut().update_order(&accepted).unwrap();
5808
5809 if !open {
5810 let canceled =
5811 TestOrderEventStubs::canceled(&order, account_id, Some(venue_order_id));
5812 cache.borrow_mut().update_order(&canceled).unwrap();
5813 }
5814
5815 orders.push(order);
5816 }
5817
5818 client
5819 .cancel_all_orders(CancelAllOrders::new(
5820 TraderId::from("TESTER-001"),
5821 Some(*BINANCE_CLIENT_ID),
5822 StrategyId::from("S-001"),
5823 instrument_id,
5824 order_side,
5825 UUID4::new(),
5826 UnixNanos::default(),
5827 None,
5828 None,
5829 ))
5830 .unwrap();
5831 wait_for_spawned_tasks(&client).await;
5832
5833 assert_eq!(hits.single_cancel.load(Ordering::Relaxed), 2);
5835 assert_eq!(hits.batch.load(Ordering::Relaxed), 0);
5836 assert_eq!(hits.cancel_all.load(Ordering::Relaxed), 0);
5837
5838 let mut actual = Vec::new();
5839
5840 for _ in &expected_indices {
5841 let event = rx.try_recv().expect("expected OrderCancelRejected event");
5842 match event {
5843 ExecutionEvent::Order(OrderEventAny::CancelRejected(rejected)) => {
5844 actual.push((rejected.client_order_id, rejected.strategy_id));
5845 }
5846 event => panic!("expected OrderCancelRejected, was {event:?}"),
5847 }
5848 }
5849
5850 let mut expected: Vec<_> = expected_indices
5851 .iter()
5852 .map(|&index| {
5853 let order = &orders[index];
5854 (order.client_order_id(), order.strategy_id())
5855 })
5856 .collect();
5857
5858 actual.sort();
5859 expected.sort();
5860
5861 assert_eq!(actual, expected);
5862 assert!(rx.try_recv().is_err());
5863 }
5864
5865 #[rstest]
5866 #[tokio::test]
5867 async fn test_cancel_all_orders_with_side_and_empty_cache_sends_nothing() {
5868 let (base_url, hits) = start_cancel_reject_server().await;
5869 let (mut client, _cache) = test_execution_client(base_url);
5870 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
5871 client.emitter.set_sender(tx);
5872
5873 client
5874 .cancel_all_orders(CancelAllOrders::new(
5875 TraderId::from("TESTER-001"),
5876 Some(*BINANCE_CLIENT_ID),
5877 StrategyId::from("S-001"),
5878 InstrumentId::from("BTCUSDT.BINANCE"),
5879 Some(OrderSide::Buy),
5880 UUID4::new(),
5881 UnixNanos::default(),
5882 None,
5883 None,
5884 ))
5885 .unwrap();
5886 wait_for_spawned_tasks(&client).await;
5887
5888 assert_eq!(hits.single_cancel.load(Ordering::Relaxed), 0);
5889 assert_eq!(hits.batch.load(Ordering::Relaxed), 0);
5890 assert_eq!(hits.cancel_all.load(Ordering::Relaxed), 0);
5891 assert!(rx.try_recv().is_err());
5892 assert!(client.dispatch_state.pending_requests.is_empty());
5893 }
5894}