1use std::{
19 future::Future,
20 sync::{
21 Arc,
22 atomic::{AtomicBool, Ordering},
23 },
24 time::Duration,
25};
26
27use ahash::AHashSet;
28use anyhow::Context;
29use async_trait::async_trait;
30use dashmap::DashMap;
31use nautilus_common::{
32 cache::fifo::FifoCache,
33 clients::ExecutionClient,
34 live::runner::get_exec_event_sender,
35 messages::execution::{
36 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
37 GenerateFillReportsBuilder, GenerateOrderStatusReport, GenerateOrderStatusReports,
38 GeneratePositionStatusReports, GeneratePositionStatusReportsBuilder, ModifyOrder,
39 PARAMS_CLOSE_POSITION, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
40 },
41};
42use nautilus_core::{
43 AtomicSet, Params, UUID4, UnixNanos,
44 datetime::{
45 NANOSECONDS_IN_DAY, NANOSECONDS_IN_MILLISECOND, NANOSECONDS_IN_SECOND,
46 checked_mins_to_nanos,
47 },
48 time::{AtomicTime, get_atomic_clock_realtime},
49};
50use nautilus_live::{
51 ExecutionClientCore, ExecutionEventEmitter, SocketControlFactory,
52 task::{TaskGroup, TaskGroupGuard, TaskJoinOutcome, TaskShutdownError, TaskSlot, finish_task},
53};
54use nautilus_model::{
55 accounts::AccountAny,
56 enums::{
57 AccountType, OmsType, OrderType, PositionSide, TimeInForce, TrailingOffsetType, TriggerType,
58 },
59 events::{
60 AccountState, OrderCancelRejected, OrderCanceled, OrderEventAny, OrderModifyRejected,
61 OrderRejected, OrderUpdated,
62 },
63 identifiers::{
64 AccountId, ClientId, ClientOrderId, InstrumentId, PositionId, Venue, VenueOrderId,
65 },
66 instruments::{Instrument, InstrumentAny},
67 orders::{Order, OrderAny},
68 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
69 types::{AccountBalance, Currency, MarginBalance, Money, Quantity},
70};
71use parking_lot::{Mutex, RwLock};
72use rust_decimal::Decimal;
73use tokio_util::sync::CancellationToken;
74
75use super::{
76 http::{
77 BinanceFuturesHttpError,
78 client::{
79 BinanceFuturesAlgoOrderQueryResult, BinanceFuturesHttpClient, BinanceFuturesInstrument,
80 is_algo_order_type,
81 },
82 models::{BatchOrderResult, BinanceFuturesAlgoOrder, BinancePositionRisk},
83 query::{
84 BatchCancelItem, BinanceAllOrdersParamsBuilder, BinanceOpenOrdersParamsBuilder,
85 BinanceOrderQueryParamsBuilder, BinancePositionRiskParamsBuilder,
86 BinanceSetLeverageParams, BinanceSetMarginTypeParams, BinanceUserTradesParamsBuilder,
87 },
88 },
89 websocket::{
90 streams::{
91 client::BinanceFuturesWebSocketClient,
92 dispatch::{
93 DispatchCtx, dispatch_user_stream_message, make_venue_position_id,
94 run_user_stream_dispatch, with_venue_position_id,
95 },
96 recovery::{
97 RecoveryCtx, WsBuildParams, build_and_connect_user_stream, run_recovery_driver,
98 },
99 },
100 trading::{client::BinanceFuturesWsTradingClient, dispatch::dispatch_ws_trading_message},
101 },
102};
103use crate::{
104 common::{
105 consts::{
106 BINANCE_FUTURES_DUAL_SIDE_SYNC_REJECT_CODE, BINANCE_FUTURES_USD_WS_API_TESTNET_URL,
107 BINANCE_FUTURES_USD_WS_API_URL, BINANCE_GTX_ORDER_REJECT_CODE,
108 BINANCE_NAUTILUS_FUTURES_BROKER_ID, BINANCE_STATUS_UNKNOWN_CODE,
109 BINANCE_UNEXPECTED_RESPONSE_CODE, BINANCE_VENUE, BINANCE_WS_HEARTBEAT_SECS,
110 },
111 credential::resolve_credentials,
112 dispatch::{OrderIdentity, PendingOperation, PendingRequest, WsDispatchState},
113 encoder::encode_broker_id,
114 enums::{
115 BinanceEnvironment, BinanceFuturesOrderType, BinancePositionSide, BinancePriceMatch,
116 BinanceProductType, BinanceSide, BinanceTimeInForce, BinanceWorkingType,
117 },
118 symbol::{format_binance_symbol, format_instrument_id},
119 urls::{get_usdm_ws_route_base_url, get_ws_private_base_url},
120 },
121 config::BinanceExecutionClientConfig,
122 futures::{
123 conversions::{
124 determine_position_side, normalize_futures_asset, reduce_only_param,
125 trailing_offset_to_callback_rate, trailing_offset_to_callback_rate_string,
126 },
127 http::{
128 client::order_type_to_binance_futures,
129 models::BinanceFuturesAccountInfo,
130 query::{
131 BatchOrderItem, BinanceCancelOrderParamsBuilder, BinanceModifyOrderParamsBuilder,
132 BinanceNewOrderParams,
133 },
134 },
135 },
136};
137
138const LISTEN_KEY_KEEPALIVE_SECS: u64 = 30 * 60;
140
141const MAX_KEEPALIVE_FAILURES: u32 = 1;
143
144const USER_TRADES_MAX_INTERVAL_MS: i64 = 7 * 24 * 60 * 60 * 1_000;
145
146const USER_TRADES_PAGE_LIMIT: u32 = 1_000;
147
148const USER_TRADES_COMPLETE_INTERVAL_NS: u64 = 88 * NANOSECONDS_IN_DAY;
152
153const USER_TRADES_MAX_TS_INIT_AGE_NS: u64 = NANOSECONDS_IN_DAY / 2;
155
156const BINANCE_GTD_MIN_LEAD_SECS: u64 = 600;
157
158const BINANCE_GTD_MAX_MILLIS: u64 = 253_402_300_799_000;
159
160pub const BINANCE_VENUE_ORDER_ID_IS_ALGO_ID_PARAM: &str = "venue_order_id_is_algo_id";
163
164#[derive(Debug)]
173pub struct BinanceFuturesExecutionClient {
174 core: ExecutionClientCore,
175 clock: &'static AtomicTime,
176 config: BinanceExecutionClientConfig,
177 emitter: ExecutionEventEmitter,
178 dispatch_state: Arc<WsDispatchState>,
179 product_type: BinanceProductType,
180 http_client: BinanceFuturesHttpClient,
181 ws_client: Arc<Mutex<Option<BinanceFuturesWebSocketClient>>>,
182 socket_factory: SocketControlFactory,
183 ws_trading_client: Option<BinanceFuturesWsTradingClient>,
184 listen_key: Arc<RwLock<Option<String>>>,
185 recovery_listen_key: Arc<RwLock<Option<String>>>,
186 cancellation_token: CancellationToken,
187 triggered_algo_order_ids: Arc<AtomicSet<ClientOrderId>>,
188 algo_client_order_ids: Arc<AtomicSet<ClientOrderId>>,
189 ws_task: Arc<tokio::sync::Mutex<TaskSlot<()>>>,
190 recovery_lock: Arc<tokio::sync::Mutex<()>>,
191 recovery_tx: Option<tokio::sync::mpsc::UnboundedSender<()>>,
192 session_tasks: TaskGroup,
193 pending_tasks: TaskGroup,
194 is_hedge_mode: AtomicBool,
195 shutdown_errors: Vec<String>,
196}
197
198impl BinanceFuturesExecutionClient {
199 pub fn new(
206 core: ExecutionClientCore,
207 config: BinanceExecutionClientConfig,
208 ) -> anyhow::Result<Self> {
209 config.validate()?;
210 let product_type = config.product_type;
211 match product_type {
212 BinanceProductType::UsdM | BinanceProductType::CoinM => {}
213 _ => {
214 anyhow::bail!(
215 "BinanceFuturesExecutionClient requires UsdM or CoinM product type, was {product_type:?}"
216 );
217 }
218 }
219
220 let (api_key, api_secret) = resolve_credentials(
221 config.api_key.clone(),
222 config.api_secret.clone(),
223 config.environment,
224 product_type,
225 )?;
226
227 let clock = get_atomic_clock_realtime();
228 let socket_factory = SocketControlFactory::new(core.client_id, Some(*BINANCE_VENUE));
229
230 let http_client = BinanceFuturesHttpClient::new(
231 product_type,
232 config.environment,
233 clock,
234 Some(api_key.clone()),
235 Some(api_secret.clone()),
236 config.base_url_http.clone(),
237 Some(config.recv_window_ms),
238 None, config.proxy_url.clone(),
240 config.treat_expired_as_canceled,
241 )
242 .context("failed to construct Binance Futures HTTP client")?;
243
244 let ws_trading_client = if config.use_ws_trading && product_type == BinanceProductType::UsdM
245 {
246 let ws_trading_url =
247 config
248 .base_url_ws_trading
249 .clone()
250 .or_else(|| match config.environment {
251 BinanceEnvironment::Testnet | BinanceEnvironment::Demo => {
252 Some(BINANCE_FUTURES_USD_WS_API_TESTNET_URL.to_string())
253 }
254 _ => Some(BINANCE_FUTURES_USD_WS_API_URL.to_string()),
255 });
256
257 Some(
258 BinanceFuturesWsTradingClient::new(
259 ws_trading_url,
260 api_key,
261 api_secret,
262 Some(BINANCE_WS_HEARTBEAT_SECS),
263 config.transport_backend,
264 )
265 .with_proxy(config.proxy_url.clone())
266 .with_recv_window(Some(config.recv_window_ms))
267 .with_socket_control(socket_factory.control("binance-futures-trading")),
268 )
269 } else {
270 None
271 };
272
273 let emitter = ExecutionEventEmitter::new(
274 clock,
275 core.trader_id,
276 core.account_id,
277 core.account_type,
278 core.base_currency,
279 );
280
281 let session_tasks = TaskGroup::new();
282 let pending_tasks = TaskGroup::new();
283
284 Ok(Self {
285 core,
286 clock,
287 config,
288 emitter,
289 dispatch_state: Arc::new(WsDispatchState::default()),
290 product_type,
291 http_client,
292 ws_client: Arc::new(Mutex::new(None)),
293 socket_factory,
294 ws_trading_client,
295 listen_key: Arc::new(RwLock::new(None)),
296 recovery_listen_key: Arc::new(RwLock::new(None)),
297 cancellation_token: CancellationToken::new(),
298 triggered_algo_order_ids: Arc::new(AtomicSet::new()),
299 algo_client_order_ids: Arc::new(AtomicSet::new()),
300 ws_task: Arc::new(tokio::sync::Mutex::new(TaskSlot::new())),
301 recovery_lock: Arc::new(tokio::sync::Mutex::new(())),
302 recovery_tx: None,
303 session_tasks,
304 pending_tasks,
305 is_hedge_mode: AtomicBool::new(false),
306 shutdown_errors: Vec::new(),
307 })
308 }
309
310 #[must_use]
312 pub fn is_hedge_mode(&self) -> bool {
313 self.is_hedge_mode.load(Ordering::Acquire)
314 }
315
316 fn resolve_algo_lookup(
317 &self,
318 client_order_id: Option<ClientOrderId>,
319 params: Option<&Params>,
320 ) -> BinanceFuturesAlgoLookup {
321 if params.and_then(|p| p.get_bool(BINANCE_VENUE_ORDER_ID_IS_ALGO_ID_PARAM)) == Some(true) {
322 return BinanceFuturesAlgoLookup::AlgoId;
323 }
324
325 let Some(client_order_id) = client_order_id else {
326 return BinanceFuturesAlgoLookup::ClientAlgoId;
327 };
328 let cache = self.core.cache();
329 match cache.order(&client_order_id) {
330 Some(order) if !is_algo_order_type(order.order_type()) => {
331 BinanceFuturesAlgoLookup::Skip
332 }
333 _ => BinanceFuturesAlgoLookup::ClientAlgoId,
334 }
335 }
336
337 #[doc(hidden)]
339 #[must_use]
340 pub fn instruments_cache(&self) -> Arc<DashMap<ustr::Ustr, BinanceFuturesInstrument>> {
341 self.http_client.instruments_cache()
342 }
343
344 fn create_account_state(&self, account_info: &BinanceFuturesAccountInfo) -> AccountState {
346 Self::create_account_state_from(
347 account_info,
348 self.core.account_id,
349 self.core.account_type,
350 self.config.bnfcr_currency,
351 self.clock,
352 )
353 }
354
355 fn create_account_state_from(
356 account_info: &BinanceFuturesAccountInfo,
357 account_id: AccountId,
358 account_type: AccountType,
359 bnfcr_currency: Currency,
360 clock: &'static AtomicTime,
361 ) -> AccountState {
362 let ts_now = clock.get_time_ns();
363
364 let balances: Vec<AccountBalance> = account_info
365 .assets
366 .iter()
367 .filter_map(|b| {
368 if b.wallet_balance.is_zero() {
369 return None;
370 }
371
372 let currency = normalize_futures_asset(b.asset, bnfcr_currency);
373 AccountBalance::from_total_and_free(b.wallet_balance, b.available_balance, currency)
374 .ok()
375 })
376 .collect();
377
378 let mut margins: Vec<MarginBalance> = Vec::new();
383
384 for asset in &account_info.assets {
385 let initial_dec = asset.initial_margin.unwrap_or_default();
386 let maint_dec = asset.maint_margin.unwrap_or_default();
387
388 if initial_dec.is_zero() && maint_dec.is_zero() {
389 continue;
390 }
391
392 let currency = normalize_futures_asset(asset.asset, bnfcr_currency);
393 let initial = Money::from_decimal(initial_dec, currency)
394 .unwrap_or_else(|_| Money::zero(currency));
395 let maintenance =
396 Money::from_decimal(maint_dec, currency).unwrap_or_else(|_| Money::zero(currency));
397 margins.push(MarginBalance::new(initial, maintenance, None));
398 }
399
400 let mut info = Params::new();
401 let mut push_decimal = |key: &str, val: Option<Decimal>| {
402 if let Some(decimal) = val {
403 info.insert(
404 key.to_string(),
405 serde_json::Value::from(decimal.to_string()),
406 );
407 }
408 };
409 push_decimal("total_wallet_balance", account_info.total_wallet_balance);
410 push_decimal("total_margin_balance", account_info.total_margin_balance);
411 push_decimal("total_initial_margin", account_info.total_initial_margin);
412 push_decimal("total_maint_margin", account_info.total_maint_margin);
413 push_decimal(
414 "total_unrealized_profit",
415 account_info.total_unrealized_profit,
416 );
417 push_decimal(
418 "total_cross_wallet_balance",
419 account_info.total_cross_wallet_balance,
420 );
421 push_decimal("total_cross_unpnl", account_info.total_cross_un_pnl);
422 push_decimal("available_balance", account_info.available_balance);
423 push_decimal("max_withdraw_amount", account_info.max_withdraw_amount);
424 let info = if info.is_empty() { None } else { Some(info) };
425
426 AccountState::new(
427 account_id,
428 account_type,
429 balances,
430 margins,
431 true, UUID4::new(),
433 ts_now,
434 ts_now,
435 None, )
437 .with_info(info)
438 }
439
440 async fn refresh_account_state(&self) -> anyhow::Result<AccountState> {
441 let account_info = match self.http_client.query_account().await {
442 Ok(info) => info,
443 Err(e) => {
444 log::error!("Binance Futures account state request failed: {e}");
445 anyhow::bail!("Binance Futures account state request failed: {e}");
446 }
447 };
448
449 Ok(self.create_account_state(&account_info))
450 }
451
452 fn update_account_state(&self) {
453 let http_client = self.http_client.clone();
454 let account_id = self.core.account_id;
455 let account_type = self.core.account_type;
456 let bnfcr_currency = self.config.bnfcr_currency;
457 let emitter = self.emitter.clone();
458 let clock = self.clock;
459
460 self.spawn_task("query_account", async move {
461 let account_info = http_client
462 .query_account()
463 .await
464 .context("Binance Futures account state request failed")?;
465 let account_state = Self::create_account_state_from(
466 &account_info,
467 account_id,
468 account_type,
469 bnfcr_currency,
470 clock,
471 );
472 let ts_now = clock.get_time_ns();
473 emitter.emit_account_state(
474 account_state.balances.clone(),
475 account_state.margins.clone(),
476 account_state.is_reported,
477 ts_now,
478 account_state.info,
479 );
480 Ok(())
481 });
482 }
483
484 async fn init_hedge_mode(&self) -> anyhow::Result<bool> {
485 let response = self.http_client.query_hedge_mode().await?;
486 Ok(response.dual_side_position)
487 }
488
489 fn ws_trading_active(&self) -> bool {
491 self.ws_trading_client
492 .as_ref()
493 .is_some_and(|c| c.is_active())
494 }
495
496 fn submit_order_internal(
497 &self,
498 cmd: &SubmitOrder,
499 lifetime: FuturesOrderLifetime,
500 position_side: Option<BinancePositionSide>,
501 venue_position_id: Option<PositionId>,
502 ) -> anyhow::Result<()> {
503 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
504
505 let emitter = self.emitter.clone();
506 let trader_id = self.core.trader_id;
507 let account_id = self.core.account_id;
508 let clock = self.clock;
509 let client_order_id = order.client_order_id();
510 let strategy_id = order.strategy_id();
511 let instrument_id = order.instrument_id();
512 let order_side = order.order_side();
513 let order_type = order.order_type();
514 let quantity = order.quantity();
515 let time_in_force = lifetime.time_in_force;
516 let good_till_date = lifetime.good_till_date;
517 let price = order.price();
518 let trigger_price = order.trigger_price();
519 let reduce_only = order.is_reduce_only();
520 let post_only = order.is_post_only();
521 let activation_price = order.activation_price();
522 let trailing_offset = order.trailing_offset();
523 let trigger_type = order.trigger_type();
524
525 let close_position = cmd
526 .params
527 .as_ref()
528 .and_then(|p| p.get_bool(PARAMS_CLOSE_POSITION))
529 .unwrap_or(false);
530
531 self.dispatch_state.order_identities.insert(
533 client_order_id,
534 OrderIdentity {
535 instrument_id,
536 strategy_id,
537 order_side,
538 order_type,
539 price,
540 quantity,
541 venue_position_id,
542 },
543 );
544
545 let use_algo_api = is_algo_order_type(order_type);
546
547 let price_match = cmd
548 .params
549 .as_ref()
550 .and_then(|p| p.get_str("price_match"))
551 .map(BinancePriceMatch::from_param)
552 .transpose()?;
553
554 let callback_rate = trailing_offset
555 .map(trailing_offset_to_callback_rate_string)
556 .transpose()?;
557
558 let working_type = match trigger_type {
559 Some(TriggerType::MarkPrice) => Some(BinanceWorkingType::MarkPrice),
560 Some(TriggerType::LastPrice | TriggerType::Default) => {
561 Some(BinanceWorkingType::ContractPrice)
562 }
563 _ => None,
564 };
565
566 if self.ws_trading_active() && !use_algo_api {
568 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
569 let dispatch_state = self.dispatch_state.clone();
570
571 let symbol = format_binance_symbol(&instrument_id);
572 let binance_side = BinanceSide::try_from(order_side)?;
573 let binance_order_type = order_type_to_binance_futures(order_type)?;
574 let binance_tif = if post_only {
575 BinanceTimeInForce::Gtx
576 } else {
577 BinanceTimeInForce::try_from(time_in_force)?
578 };
579
580 let requires_time_in_force = matches!(
581 order_type,
582 OrderType::Limit | OrderType::StopLimit | OrderType::LimitIfTouched
583 );
584
585 let client_id_str =
586 encode_broker_id(&client_order_id, BINANCE_NAUTILUS_FUTURES_BROKER_ID);
587
588 let params = BinanceNewOrderParams {
589 symbol,
590 side: binance_side,
591 order_type: binance_order_type,
592 time_in_force: if requires_time_in_force {
593 Some(binance_tif)
594 } else {
595 None
596 },
597 quantity: Some(quantity.to_string()),
598 price: if price_match.is_some() {
599 None
600 } else {
601 price.map(|p| p.to_string())
602 },
603 new_client_order_id: Some(client_id_str),
604 stop_price: trigger_price.map(|p| p.to_string()),
605 reduce_only: reduce_only_param(reduce_only, position_side),
606 position_side,
607 close_position: None,
608 activation_price: activation_price.map(|p| p.to_string()),
609 callback_rate,
610 working_type,
611 price_protect: None,
612 new_order_resp_type: None,
613 good_till_date,
614 recv_window: None,
615 price_match,
616 self_trade_prevention_mode: None,
617 };
618
619 let request_id = ws_client.next_request_id();
621 dispatch_state.pending_requests.insert(
622 request_id.clone(),
623 PendingRequest {
624 client_order_id,
625 venue_order_id: None,
626 operation: PendingOperation::Place,
627 },
628 );
629
630 self.spawn_task("submit_order_ws", async move {
631 if let Err(e) = ws_client
632 .place_order_with_id(request_id.clone(), params)
633 .await
634 {
635 dispatch_state.pending_requests.remove(&request_id);
636 log::error!("WS submit request failed for {client_order_id}: {e}");
637 anyhow::bail!("WS submit order failed: {e}");
638 }
639 Ok(())
640 });
641
642 return Ok(());
643 }
644
645 let http_client = self.http_client.clone();
646
647 self.spawn_task("submit_order", async move {
648 let result = if use_algo_api {
649 http_client
650 .submit_algo_order(
651 account_id,
652 instrument_id,
653 client_order_id,
654 order_side,
655 order_type,
656 quantity,
657 time_in_force,
658 price,
659 trigger_price,
660 reduce_only,
661 close_position,
662 position_side,
663 activation_price,
664 callback_rate,
665 working_type,
666 good_till_date,
667 )
668 .await
669 } else {
670 http_client
671 .submit_order(
672 account_id,
673 instrument_id,
674 client_order_id,
675 order_side,
676 order_type,
677 quantity,
678 time_in_force,
679 price,
680 trigger_price,
681 reduce_only,
682 post_only,
683 position_side,
684 price_match,
685 good_till_date,
686 )
687 .await
688 };
689
690 match result {
691 Ok(report) => {
692 log::debug!(
693 "Order submit accepted: client_order_id={}, venue_order_id={}",
694 client_order_id,
695 report.venue_order_id
696 );
697 }
698 Err(e) => {
699 if is_ambiguous_submit_error(&e) {
703 log::warn!(
704 "Ambiguous submit failure for {client_order_id}, awaiting reconciliation: {e}"
705 );
706 } else if is_structured_venue_rejection(&e) {
707 let due_post_only = classify_submit_order_error(&e);
708 let ts_now = clock.get_time_ns();
709
710 let rejected = OrderRejected::new(
711 trader_id,
712 strategy_id,
713 instrument_id,
714 client_order_id,
715 account_id,
716 format!("submit-order-error: {e}").into(),
717 UUID4::new(),
718 ts_now,
719 ts_now,
720 false,
721 due_post_only,
722 );
723
724 emitter.send_order_event(OrderEventAny::Rejected(rejected));
725 } else {
726 log::warn!(
727 "Ambiguous submit failure for {client_order_id}, awaiting reconciliation: {e}"
728 );
729 }
730
731 return Err(e);
732 }
733 }
734
735 Ok(())
736 });
737
738 Ok(())
739 }
740
741 fn cancel_order_internal(&self, cmd: &CancelOrder) {
742 let command = cmd.clone();
743
744 let is_algo = self
746 .core
747 .cache()
748 .order(&command.client_order_id)
749 .is_some_and(|order| is_algo_order_type(order.order_type()));
750 let promoted_venue_order_id = self
751 .dispatch_state
752 .promoted_algo_order_id(&command.client_order_id);
753 let use_algo_cancel = should_use_algo_cancel(
754 is_algo,
755 self.triggered_algo_order_ids
756 .contains(&command.client_order_id),
757 promoted_venue_order_id.is_some(),
758 );
759
760 let emitter = self.emitter.clone();
761 let trader_id = self.core.trader_id;
762 let account_id = self.core.account_id;
763 let clock = self.clock;
764 let instrument_id = command.instrument_id;
765 let venue_order_id =
766 cancel_venue_order_id(is_algo, command.venue_order_id, promoted_venue_order_id);
767 let client_order_id = command.client_order_id;
768
769 if self.ws_trading_active() && !use_algo_cancel {
771 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
772 let dispatch_state = self.dispatch_state.clone();
773
774 let mut cancel_builder = BinanceCancelOrderParamsBuilder::default();
775 cancel_builder.symbol(format_binance_symbol(&instrument_id));
776
777 if let Some(venue_id) = venue_order_id {
778 match venue_id.inner().parse::<i64>() {
779 Ok(order_id) => {
780 cancel_builder.order_id(order_id);
781 }
782 Err(e) => {
783 log::warn!(
784 "Unable to parse venue_order_id {venue_id} for cancel {client_order_id}, canceling by client_order_id: {e}"
785 );
786 }
787 }
788 }
789
790 cancel_builder.orig_client_order_id(encode_broker_id(
791 &client_order_id,
792 BINANCE_NAUTILUS_FUTURES_BROKER_ID,
793 ));
794
795 let params = cancel_builder.build().unwrap();
796
797 let request_id = ws_client.next_request_id();
799 dispatch_state.pending_requests.insert(
800 request_id.clone(),
801 PendingRequest {
802 client_order_id,
803 venue_order_id,
804 operation: PendingOperation::Cancel,
805 },
806 );
807
808 self.spawn_task("cancel_order_ws", async move {
809 if let Err(e) = ws_client
810 .cancel_order_with_id(request_id.clone(), params)
811 .await
812 {
813 dispatch_state.pending_requests.remove(&request_id);
814 log::error!("WS cancel request failed for {client_order_id}: {e}");
815 anyhow::bail!("WS cancel order failed: {e}");
816 }
817 Ok(())
818 });
819
820 return;
821 }
822
823 let http_client = self.http_client.clone();
824
825 self.spawn_task("cancel_order", async move {
826 let result = if use_algo_cancel {
827 match http_client.cancel_algo_order(client_order_id).await {
830 Ok(()) => Ok(()),
831 Err(algo_err) => {
832 log::debug!("Algo cancel failed, trying regular cancel: {algo_err}");
833 http_client
834 .cancel_order(instrument_id, venue_order_id, Some(client_order_id))
835 .await
836 .map(|_| ())
837 }
838 }
839 } else {
840 http_client
841 .cancel_order(instrument_id, venue_order_id, Some(client_order_id))
842 .await
843 .map(|_| ())
844 };
845
846 match result {
847 Ok(()) => {
848 log::debug!("Cancel request accepted: client_order_id={client_order_id}");
849 }
850 Err(e) => {
851 if is_structured_venue_rejection(&e) {
852 let ts_now = clock.get_time_ns();
853
854 let rejected = OrderCancelRejected::new(
855 trader_id,
856 command.strategy_id,
857 command.instrument_id,
858 client_order_id,
859 format!("cancel-order-error: {e}").into(),
860 UUID4::new(),
861 ts_now,
862 ts_now,
863 false,
864 command.venue_order_id,
865 Some(account_id),
866 );
867
868 emitter.send_order_event(OrderEventAny::CancelRejected(rejected));
869 } else if is_local_command_failure(&e) {
870 log::warn!(
871 "Cancel command failed local validation for {client_order_id}: {e}"
872 );
873 } else {
874 log::warn!(
875 "Ambiguous cancel failure for {client_order_id}, awaiting reconciliation: {e}"
876 );
877 }
878
879 return Err(e);
880 }
881 }
882
883 Ok(())
884 });
885 }
886
887 fn spawn_task<F>(&self, description: &'static str, fut: F)
888 where
889 F: Future<Output = anyhow::Result<()>> + Send + 'static,
890 {
891 crate::common::execution::spawn_task(&self.pending_tasks, description, fut);
892 }
893
894 fn abort_pending_tasks(&self) {
895 crate::common::execution::abort_pending_tasks(&self.pending_tasks);
896 }
897
898 fn abort_session_tasks(&self) {
899 self.session_tasks.begin_shutdown();
900 }
901
902 async fn await_pending_tasks(&self) -> anyhow::Result<()> {
903 crate::common::execution::await_pending_tasks(&self.pending_tasks).await
904 }
905
906 async fn await_session_tasks(&self) -> anyhow::Result<()> {
907 self.finish_session_tasks().await.map_err(|e| {
908 anyhow::anyhow!("Failed to terminate Binance Futures session tasks: {e}")
909 })?;
910 Ok(())
911 }
912
913 async fn finish_session_tasks(&self) -> Result<(), TaskShutdownError> {
914 self.session_tasks.begin_shutdown();
915 self.session_tasks
916 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
917 .await?;
918 Ok(())
919 }
920
921 async fn await_dispatch_task(&self) -> anyhow::Result<()> {
922 let _recovery_guard = self.recovery_lock.lock().await;
923 let mut task_slot = self.ws_task.lock().await;
924 let Some(outcome) = finish_task(
925 &mut task_slot,
926 Duration::from_secs(1),
927 Duration::from_secs(2),
928 )
929 .await
930 else {
931 return Ok(());
932 };
933
934 match outcome {
935 TaskJoinOutcome::Completed(()) | TaskJoinOutcome::Aborted => Ok(()),
936 TaskJoinOutcome::Failed(e) => {
937 Err(anyhow::anyhow!("Binance Futures dispatch task failed: {e}"))
938 }
939 TaskJoinOutcome::Incomplete => Err(anyhow::anyhow!(
940 "Binance Futures dispatch task did not stop after abort"
941 )),
942 }
943 }
944
945 async fn close_listen_key_slot(
946 &self,
947 slot: &RwLock<Option<String>>,
948 context: &str,
949 ) -> anyhow::Result<()> {
950 let key = slot.read().clone();
951 let Some(key) = key else {
952 return Ok(());
953 };
954
955 self.http_client
956 .close_listen_key(&key)
957 .await
958 .with_context(|| context.to_string())?;
959 let mut owned = slot.write();
960 if owned.as_deref() == Some(key.as_str()) {
961 *owned = None;
962 }
963 Ok(())
964 }
965
966 fn get_instrument_precision(&self, instrument_id: InstrumentId) -> anyhow::Result<(u8, u8)> {
968 self.http_client
969 .instrument_reconciliation(&instrument_id)
970 .map(|instrument| (instrument.price_precision(), instrument.size_precision()))
971 .ok_or_else(|| {
972 anyhow::anyhow!(
973 "Binance Futures instrument {instrument_id} is not loaded for reconciliation"
974 )
975 })
976 }
977
978 fn is_instrument_out_of_scope(&self, instrument_id: InstrumentId) -> bool {
979 let provider = &self.config.instrument_provider;
980 !provider.load_all
981 && provider.load_ids.as_ref().is_some_and(|load_ids| {
982 load_ids
983 .iter()
984 .all(|raw_id| InstrumentId::from(raw_id.as_str()) != instrument_id)
985 })
986 }
987
988 fn create_position_report(
990 &self,
991 position: &BinancePositionRisk,
992 instrument_id: InstrumentId,
993 size_precision: u8,
994 ) -> anyhow::Result<PositionStatusReport> {
995 let position_amount: Decimal = position
996 .position_amt
997 .parse()
998 .context("invalid position_amt")?;
999
1000 if position_amount.is_zero() {
1001 anyhow::bail!("Position is flat");
1002 }
1003
1004 let entry_price: Decimal = position
1005 .entry_price
1006 .parse()
1007 .context("invalid entry_price")?;
1008
1009 let position_side = if position_amount > Decimal::ZERO {
1010 PositionSide::Long
1011 } else {
1012 PositionSide::Short
1013 };
1014
1015 if self.config.use_position_ids {
1016 match position.position_side {
1017 Some(BinancePositionSide::Long) => anyhow::ensure!(
1018 position_side == PositionSide::Long,
1019 "position_side LONG conflicts with negative position_amt"
1020 ),
1021 Some(BinancePositionSide::Short) => anyhow::ensure!(
1022 position_side == PositionSide::Short,
1023 "position_side SHORT conflicts with positive position_amt"
1024 ),
1025 _ => {}
1026 }
1027 }
1028
1029 let venue_position_id = make_venue_position_id(
1030 self.config.use_position_ids,
1031 instrument_id,
1032 position.position_side,
1033 )?;
1034
1035 if let Some(venue_position_id) = venue_position_id {
1036 self.ensure_cached_position_id_compatible(
1037 instrument_id,
1038 position_side,
1039 venue_position_id,
1040 )?;
1041 }
1042
1043 let ts_now = self.clock.get_time_ns();
1044
1045 Ok(PositionStatusReport::new(
1046 self.core.account_id,
1047 instrument_id,
1048 position_side,
1049 Quantity::from_decimal_dp(position_amount.abs(), size_precision)?,
1050 ts_now,
1051 ts_now,
1052 Some(UUID4::new()),
1053 venue_position_id,
1054 Some(entry_price),
1055 ))
1056 }
1057
1058 fn ensure_cached_position_id_compatible(
1059 &self,
1060 instrument_id: InstrumentId,
1061 position_side: PositionSide,
1062 venue_position_id: PositionId,
1063 ) -> anyhow::Result<()> {
1064 let cache = self.core.cache();
1065 let mut incompatible_ids: Vec<_> = cache
1066 .positions_open(
1067 Some(&BINANCE_VENUE),
1068 Some(&instrument_id),
1069 None,
1070 Some(&self.core.account_id),
1071 Some(position_side),
1072 )
1073 .into_iter()
1074 .filter(|position| position.id != venue_position_id)
1075 .map(|position| position.id.to_string())
1076 .collect();
1077 incompatible_ids.sort_unstable();
1078
1079 anyhow::ensure!(
1080 incompatible_ids.is_empty(),
1081 "incompatible cached {position_side:?} position IDs for {instrument_id}: {}; expected {venue_position_id}",
1082 incompatible_ids.join(", "),
1083 );
1084 Ok(())
1085 }
1086
1087 async fn generate_open_order_status_reports(
1088 &self,
1089 instrument_id: Option<InstrumentId>,
1090 ts_init: UnixNanos,
1091 ) -> anyhow::Result<Vec<OpenOrderStatusReport>> {
1092 if let Some(instrument_id) = instrument_id
1093 && self
1094 .http_client
1095 .instrument_reconciliation(&instrument_id)
1096 .is_none()
1097 {
1098 if self.is_instrument_out_of_scope(instrument_id) {
1099 log::debug!(
1100 "Dropping out-of-scope Binance Futures order request for instrument {instrument_id}"
1101 );
1102 return Ok(Vec::new());
1103 }
1104
1105 anyhow::bail!(
1106 "Binance Futures open order request has unresolved instrument {instrument_id}"
1107 );
1108 }
1109
1110 let symbol = instrument_id.map(|id| format_binance_symbol(&id));
1111 let mut builder = BinanceOpenOrdersParamsBuilder::default();
1112
1113 if let Some(symbol) = symbol {
1114 builder.symbol(symbol);
1115 }
1116 let params = builder.build().map_err(|e| anyhow::anyhow!("{e}"))?;
1117
1118 let (orders, algo_orders) = tokio::try_join!(
1119 self.http_client.query_open_orders(¶ms),
1120 self.http_client.query_open_algo_orders(instrument_id),
1121 )?;
1122 let mut reports = Vec::with_capacity(orders.len() + algo_orders.len());
1123
1124 for order in orders {
1125 let instrument_id = instrument_id
1126 .unwrap_or_else(|| format_instrument_id(&order.symbol, self.product_type));
1127 let Some(instrument) = self.http_client.instrument_reconciliation(&instrument_id)
1128 else {
1129 if self.is_instrument_out_of_scope(instrument_id) {
1130 log::debug!(
1131 "Dropping out-of-scope Binance Futures open order for instrument {instrument_id}"
1132 );
1133 continue;
1134 }
1135 anyhow::bail!(
1136 "Binance Futures open order has unresolved instrument {instrument_id}"
1137 );
1138 };
1139
1140 if let Ok(report) = order.to_order_status_report(
1141 self.core.account_id,
1142 instrument.id(),
1143 instrument.price_precision(),
1144 instrument.size_precision(),
1145 self.config.treat_expired_as_canceled,
1146 ts_init,
1147 ) {
1148 let venue_position_id = make_venue_position_id(
1149 self.config.use_position_ids,
1150 instrument.id(),
1151 order.position_side,
1152 )?;
1153 reports.push(OpenOrderStatusReport {
1154 report: with_venue_position_id(report, venue_position_id),
1155 quantity_free_close_position_side: None,
1156 });
1157 }
1158 }
1159
1160 for algo_order in algo_orders {
1161 let instrument_id = instrument_id
1162 .unwrap_or_else(|| format_instrument_id(&algo_order.symbol, self.product_type));
1163 let Some(instrument) = self.http_client.instrument_reconciliation(&instrument_id)
1164 else {
1165 if self.is_instrument_out_of_scope(instrument_id) {
1166 log::debug!(
1167 "Dropping out-of-scope Binance Futures open algo order for instrument {instrument_id}"
1168 );
1169 continue;
1170 }
1171 anyhow::bail!(
1172 "Binance Futures open algo order has unresolved instrument {instrument_id}"
1173 );
1174 };
1175
1176 if let Ok(report) = algo_order.to_order_status_report(
1177 self.core.account_id,
1178 instrument.id(),
1179 instrument.price_precision(),
1180 instrument.size_precision(),
1181 ts_init,
1182 ) {
1183 let venue_position_id = make_venue_position_id(
1184 self.config.use_position_ids,
1185 instrument.id(),
1186 algo_order.position_side,
1187 )?;
1188 reports.push(OpenOrderStatusReport {
1189 report: with_venue_position_id(report, venue_position_id),
1190 quantity_free_close_position_side: quantity_free_close_position_side(
1191 &algo_order,
1192 ),
1193 });
1194 }
1195 }
1196
1197 Ok(reports)
1198 }
1199
1200 async fn apply_futures_config(&self) -> anyhow::Result<()> {
1201 if let Some(ref leverages) = self.config.futures_leverages {
1202 for (symbol, leverage) in leverages {
1203 let params = BinanceSetLeverageParams {
1204 symbol: symbol.clone(),
1205 leverage: *leverage,
1206 recv_window: None,
1207 };
1208 match self.http_client.set_leverage(¶ms).await {
1211 Ok(response) => {
1212 log::info!("Set leverage {} {}X", response.symbol, response.leverage);
1213 }
1214 Err(BinanceFuturesHttpError::BinanceError { code, message }) => {
1215 log::warn!(
1216 "Unable to set leverage for {symbol} to {leverage}x: [{code}] {message}; skipping (leverage init is best-effort)"
1217 );
1218 }
1219 Err(e) => {
1220 return Err(e).context(format!("failed to set leverage for {symbol}"));
1221 }
1222 }
1223 }
1224 }
1225
1226 if let Some(ref margin_types) = self.config.futures_margin_types {
1227 for (symbol, margin_type) in margin_types {
1228 let params = BinanceSetMarginTypeParams {
1229 symbol: symbol.clone(),
1230 margin_type: *margin_type,
1231 recv_window: None,
1232 };
1233
1234 match self.http_client.set_margin_type(¶ms).await {
1235 Ok(_) => {
1236 log::info!("Set {symbol} margin type to {margin_type:?}");
1237 }
1238 Err(e) => {
1239 let err_str = format!("{e}");
1240 if err_str.contains("-4046") {
1241 log::debug!("{symbol} margin type already {margin_type:?}");
1242 } else {
1243 return Err(e)
1244 .context(format!("failed to set margin type for {symbol}"));
1245 }
1246 }
1247 }
1248 }
1249 }
1250
1251 Ok(())
1252 }
1253}
1254
1255fn quantity_free_close_position_side(order: &BinanceFuturesAlgoOrder) -> Option<PositionSide> {
1256 let quantity_free = match order.quantity.as_deref() {
1257 None => true,
1258 Some(quantity) => quantity
1259 .parse::<Decimal>()
1260 .is_ok_and(|quantity| quantity.is_zero()),
1261 };
1262
1263 if order.close_position != Some(true) || !quantity_free {
1264 return None;
1265 }
1266
1267 match (order.side, order.position_side) {
1268 (BinanceSide::Sell, Some(BinancePositionSide::Long | BinancePositionSide::Both) | None) => {
1269 Some(PositionSide::Long)
1270 }
1271 (BinanceSide::Buy, Some(BinancePositionSide::Short | BinancePositionSide::Both) | None) => {
1272 Some(PositionSide::Short)
1273 }
1274 _ => None,
1275 }
1276}
1277
1278fn restore_close_position_quantities(
1279 order_reports: &mut [OpenOrderStatusReport],
1280 position_reports: &[PositionStatusReport],
1281) {
1282 for order_report in order_reports {
1286 let Some(position_side) = order_report.quantity_free_close_position_side else {
1287 continue;
1288 };
1289
1290 if order_report.report.quantity.is_positive() {
1291 continue;
1292 }
1293
1294 let mut matching_positions = position_reports.iter().filter(|position| {
1295 position.account_id == order_report.report.account_id
1296 && position.instrument_id == order_report.report.instrument_id
1297 && position.position_side == position_side
1298 });
1299 let Some(position) = matching_positions.next() else {
1300 continue;
1301 };
1302
1303 if matching_positions.next().is_some() {
1304 log::warn!(
1305 "Cannot restore close-position quantity for order {}: multiple {:?} position reports for {}",
1306 order_report.report.venue_order_id,
1307 position_side,
1308 order_report.report.instrument_id,
1309 );
1310 continue;
1311 }
1312
1313 order_report.report.quantity = position.quantity;
1314 }
1315}
1316
1317fn resolve_order_position_identity(
1318 is_hedge_mode: bool,
1319 use_position_ids: bool,
1320 order: &OrderAny,
1321 close_position: bool,
1322) -> anyhow::Result<(Option<BinancePositionSide>, Option<PositionId>)> {
1323 let position_side = determine_position_side(
1326 is_hedge_mode,
1327 order.order_side(),
1328 order.is_reduce_only() || close_position,
1329 );
1330 let venue_position_id = make_venue_position_id(
1331 use_position_ids,
1332 order.instrument_id(),
1333 Some(position_side.unwrap_or(BinancePositionSide::Both)),
1334 )?;
1335 Ok((position_side, venue_position_id))
1336}
1337
1338fn validate_submit_position_id(
1339 submitted_position_id: Option<PositionId>,
1340 venue_position_id: Option<PositionId>,
1341) -> anyhow::Result<()> {
1342 if let (Some(submitted_position_id), Some(venue_position_id)) =
1343 (submitted_position_id, venue_position_id)
1344 {
1345 anyhow::ensure!(
1346 submitted_position_id == venue_position_id,
1347 "submitted position ID {submitted_position_id} conflicts with canonical Binance Futures venue position ID {venue_position_id}; omit position_id while use_position_ids=true, or set use_position_ids=false for virtual hedging",
1348 );
1349 }
1350 Ok(())
1351}
1352
1353fn build_futures_order_list_batch(
1354 orders: &[OrderAny],
1355 is_hedge_mode: bool,
1356 close_position: bool,
1357 price_match: Option<BinancePriceMatch>,
1358 product_type: BinanceProductType,
1359 use_gtd: bool,
1360 ts_now: UnixNanos,
1361) -> Result<Vec<BatchOrderItem>, String> {
1362 if orders.len() > 5 {
1363 return Err(format!(
1364 "Binance Futures batch order submission supports at most 5 orders, was {}",
1365 orders.len()
1366 ));
1367 }
1368
1369 if close_position {
1370 return Err(
1371 "`close_position` is not supported for Binance Futures batch order submission"
1372 .to_string(),
1373 );
1374 }
1375
1376 if orders.iter().any(is_grouped_order) {
1377 return Err(
1378 "Binance Futures linked order-list contingencies require adapter-level OCO-on-fill state"
1379 .to_string(),
1380 );
1381 }
1382
1383 if let Some(order) = orders
1384 .iter()
1385 .find(|order| is_algo_order_type(order.order_type()))
1386 {
1387 return Err(format!(
1388 "Binance Futures batch order submission does not support conditional order type {:?}",
1389 order.order_type()
1390 ));
1391 }
1392
1393 let mut batch_items = Vec::with_capacity(orders.len());
1394 for order in orders {
1395 if price_match.is_some() && order.is_post_only() {
1396 return Err("price_match cannot be combined with post-only orders".to_string());
1397 }
1398
1399 if price_match.is_some() && order.order_type() != OrderType::Limit {
1400 return Err(format!(
1401 "price_match is not supported for order type {:?}",
1402 order.order_type()
1403 ));
1404 }
1405
1406 let binance_side = BinanceSide::try_from(order.order_side()).map_err(|e| e.to_string())?;
1407 let binance_order_type =
1408 order_type_to_binance_futures(order.order_type()).map_err(|e| e.to_string())?;
1409 let lifetime = determine_futures_order_lifetime(
1410 product_type,
1411 order.order_type(),
1412 order.time_in_force(),
1413 order.expire_time(),
1414 order.is_post_only(),
1415 use_gtd,
1416 ts_now,
1417 )
1418 .map_err(|e| e.to_string())?;
1419 let binance_tif = if order.is_post_only() {
1420 BinanceTimeInForce::Gtx
1421 } else {
1422 BinanceTimeInForce::try_from(lifetime.time_in_force).map_err(|e| e.to_string())?
1423 };
1424 let position_side =
1425 determine_position_side(is_hedge_mode, order.order_side(), order.is_reduce_only());
1426 let requires_time_in_force = matches!(order.order_type(), OrderType::Limit);
1427
1428 batch_items.push(BatchOrderItem {
1429 symbol: format_binance_symbol(&order.instrument_id()),
1430 side: binance_side_wire(binance_side).to_string(),
1431 order_type: binance_futures_order_type_wire(binance_order_type).to_string(),
1432 time_in_force: if requires_time_in_force {
1433 Some(binance_time_in_force_wire(binance_tif).to_string())
1434 } else {
1435 None
1436 },
1437 quantity: Some(order.quantity().to_string()),
1438 price: if price_match.is_some() {
1439 None
1440 } else {
1441 order.price().map(|price| price.to_string())
1442 },
1443 reduce_only: reduce_only_param(order.is_reduce_only(), position_side),
1444 new_client_order_id: Some(encode_broker_id(
1445 &order.client_order_id(),
1446 BINANCE_NAUTILUS_FUTURES_BROKER_ID,
1447 )),
1448 stop_price: None,
1449 position_side: position_side.map(|side| binance_position_side_wire(side).to_string()),
1450 activation_price: None,
1451 callback_rate: None,
1452 working_type: None,
1453 price_protect: None,
1454 close_position: None,
1455 good_till_date: lifetime.good_till_date,
1456 price_match: price_match
1457 .and_then(binance_price_match_wire)
1458 .map(str::to_string),
1459 self_trade_prevention_mode: None,
1460 });
1461 }
1462
1463 Ok(batch_items)
1464}
1465
1466fn determine_futures_order_lifetime(
1467 product_type: BinanceProductType,
1468 order_type: OrderType,
1469 time_in_force: TimeInForce,
1470 expire_time: Option<UnixNanos>,
1471 post_only: bool,
1472 use_gtd: bool,
1473 ts_now: UnixNanos,
1474) -> anyhow::Result<FuturesOrderLifetime> {
1475 if time_in_force != TimeInForce::Gtd {
1476 return Ok(FuturesOrderLifetime {
1477 time_in_force,
1478 good_till_date: None,
1479 });
1480 }
1481
1482 if !use_gtd {
1483 log::warn!(
1484 "Binance Futures GTD submitted as GTC because use_gtd=false. Enable manage_gtd_expiry on the submitting strategy"
1485 );
1486 return Ok(FuturesOrderLifetime {
1487 time_in_force: TimeInForce::Gtc,
1488 good_till_date: None,
1489 });
1490 }
1491
1492 anyhow::ensure!(
1493 matches!(
1494 order_type,
1495 OrderType::Limit | OrderType::StopLimit | OrderType::LimitIfTouched
1496 ),
1497 "Binance Futures does not support GTD for order type {order_type:?}"
1498 );
1499 anyhow::ensure!(!post_only, "Binance Futures GTD cannot be post-only");
1500
1501 anyhow::ensure!(
1502 product_type == BinanceProductType::UsdM,
1503 "Binance {product_type:?} Futures does not support native GTD"
1504 );
1505
1506 let expire_time = expire_time.context("Binance Futures GTD requires an expire_time")?;
1507 let expire_ns = expire_time.as_u64();
1508 anyhow::ensure!(
1509 expire_ns.is_multiple_of(NANOSECONDS_IN_SECOND),
1510 "Binance Futures goodTillDate requires whole-second precision"
1511 );
1512
1513 let minimum_ns = ts_now
1514 .as_u64()
1515 .checked_add(BINANCE_GTD_MIN_LEAD_SECS * NANOSECONDS_IN_SECOND)
1516 .context("Binance Futures GTD minimum timestamp overflow")?;
1517 anyhow::ensure!(
1518 expire_ns > minimum_ns,
1519 "Binance Futures goodTillDate must be strictly greater than current time plus {BINANCE_GTD_MIN_LEAD_SECS} seconds"
1520 );
1521
1522 let good_till_date = expire_ns / NANOSECONDS_IN_MILLISECOND;
1523 anyhow::ensure!(
1524 good_till_date < BINANCE_GTD_MAX_MILLIS,
1525 "Binance Futures goodTillDate must be smaller than {BINANCE_GTD_MAX_MILLIS}"
1526 );
1527
1528 Ok(FuturesOrderLifetime {
1529 time_in_force,
1530 good_till_date: Some(i64::try_from(good_till_date)?),
1531 })
1532}
1533
1534fn is_grouped_order(order: &OrderAny) -> bool {
1535 order.contingency_type().is_some()
1536 || order
1537 .linked_order_ids()
1538 .is_some_and(|linked_order_ids| !linked_order_ids.is_empty())
1539}
1540
1541fn binance_side_wire(side: BinanceSide) -> &'static str {
1542 match side {
1543 BinanceSide::Buy => "BUY",
1544 BinanceSide::Sell => "SELL",
1545 }
1546}
1547
1548fn binance_futures_order_type_wire(order_type: BinanceFuturesOrderType) -> &'static str {
1549 match order_type {
1550 BinanceFuturesOrderType::Limit => "LIMIT",
1551 BinanceFuturesOrderType::Market => "MARKET",
1552 BinanceFuturesOrderType::Stop => "STOP",
1553 BinanceFuturesOrderType::StopMarket => "STOP_MARKET",
1554 BinanceFuturesOrderType::TakeProfit => "TAKE_PROFIT",
1555 BinanceFuturesOrderType::TakeProfitMarket => "TAKE_PROFIT_MARKET",
1556 BinanceFuturesOrderType::TrailingStopMarket => "TRAILING_STOP_MARKET",
1557 BinanceFuturesOrderType::Liquidation => "LIQUIDATION",
1558 BinanceFuturesOrderType::Adl => "ADL",
1559 BinanceFuturesOrderType::Unknown => "UNKNOWN",
1560 }
1561}
1562
1563fn binance_time_in_force_wire(time_in_force: BinanceTimeInForce) -> &'static str {
1564 match time_in_force {
1565 BinanceTimeInForce::Gtc => "GTC",
1566 BinanceTimeInForce::Ioc => "IOC",
1567 BinanceTimeInForce::Fok => "FOK",
1568 BinanceTimeInForce::Gtx => "GTX",
1569 BinanceTimeInForce::Gtd => "GTD",
1570 BinanceTimeInForce::Rpi => "RPI",
1571 BinanceTimeInForce::Unknown => "UNKNOWN",
1572 }
1573}
1574
1575fn binance_position_side_wire(position_side: BinancePositionSide) -> &'static str {
1576 match position_side {
1577 BinancePositionSide::Both => "BOTH",
1578 BinancePositionSide::Long => "LONG",
1579 BinancePositionSide::Short => "SHORT",
1580 BinancePositionSide::Unknown => "UNKNOWN",
1581 }
1582}
1583
1584fn binance_price_match_wire(price_match: BinancePriceMatch) -> Option<&'static str> {
1585 match price_match {
1586 BinancePriceMatch::None | BinancePriceMatch::Unknown => None,
1587 BinancePriceMatch::Opponent => Some("OPPONENT"),
1588 BinancePriceMatch::Opponent5 => Some("OPPONENT_5"),
1589 BinancePriceMatch::Opponent10 => Some("OPPONENT_10"),
1590 BinancePriceMatch::Opponent20 => Some("OPPONENT_20"),
1591 BinancePriceMatch::Queue => Some("QUEUE"),
1592 BinancePriceMatch::Queue5 => Some("QUEUE_5"),
1593 BinancePriceMatch::Queue10 => Some("QUEUE_10"),
1594 BinancePriceMatch::Queue20 => Some("QUEUE_20"),
1595 }
1596}
1597
1598pub(crate) fn classify_submit_order_error(err: &anyhow::Error) -> bool {
1604 let venue_code = err
1605 .downcast_ref::<BinanceFuturesHttpError>()
1606 .and_then(|be| match be {
1607 BinanceFuturesHttpError::BinanceError { code, .. } => Some(*code),
1608 _ => None,
1609 });
1610
1611 if venue_code == Some(BINANCE_FUTURES_DUAL_SIDE_SYNC_REJECT_CODE) {
1612 log::warn!(
1613 "Order rejected by Binance Futures with code -4531 \
1614 (UM/CM dualSidePosition sync); confirm Portfolio Margin hedge mode \
1615 matches the order positionSide before resubmitting"
1616 );
1617 }
1618 venue_code == Some(BINANCE_GTX_ORDER_REJECT_CODE)
1619}
1620
1621fn is_structured_venue_rejection(err: &anyhow::Error) -> bool {
1622 err.downcast_ref::<BinanceFuturesHttpError>()
1623 .is_some_and(|be| matches!(be, BinanceFuturesHttpError::BinanceError { .. }))
1624}
1625
1626fn is_ambiguous_submit_error(err: &anyhow::Error) -> bool {
1627 err.downcast_ref::<BinanceFuturesHttpError>()
1628 .is_some_and(|be| {
1629 matches!(
1630 be,
1631 BinanceFuturesHttpError::BinanceError {
1632 code: BINANCE_UNEXPECTED_RESPONSE_CODE | BINANCE_STATUS_UNKNOWN_CODE,
1633 ..
1634 }
1635 )
1636 })
1637}
1638
1639fn is_local_command_failure(err: &anyhow::Error) -> bool {
1640 err.downcast_ref::<BinanceFuturesHttpError>()
1641 .is_some_and(is_local_http_command_failure)
1642}
1643
1644fn is_local_http_command_failure(err: &BinanceFuturesHttpError) -> bool {
1645 matches!(
1646 err,
1647 BinanceFuturesHttpError::MissingCredentials | BinanceFuturesHttpError::ValidationError(_)
1648 )
1649}
1650
1651#[async_trait(?Send)]
1652impl ExecutionClient for BinanceFuturesExecutionClient {
1653 fn is_connected(&self) -> bool {
1654 self.core.is_connected()
1655 }
1656
1657 fn client_id(&self) -> ClientId {
1658 self.core.client_id
1659 }
1660
1661 fn account_id(&self) -> AccountId {
1662 self.core.account_id
1663 }
1664
1665 fn venue(&self) -> Venue {
1666 *BINANCE_VENUE
1667 }
1668
1669 fn oms_type(&self) -> OmsType {
1670 self.core.oms_type
1671 }
1672
1673 fn get_account(&self) -> Option<AccountAny> {
1674 self.core.cache().account_owned(&self.core.account_id)
1675 }
1676
1677 async fn connect(&mut self) -> anyhow::Result<()> {
1678 if self.core.is_connected() && self.session_tasks.is_open() && self.pending_tasks.is_open()
1679 {
1680 return Ok(());
1681 }
1682
1683 if !self.pending_tasks.is_open() || !self.session_tasks.is_open() {
1684 self.disconnect().await?;
1685 }
1686
1687 if !self.pending_tasks.is_open() {
1688 self.await_pending_tasks().await?;
1689 self.pending_tasks.start_generation().map_err(|e| {
1690 anyhow::anyhow!("Failed to start Binance Futures task generation: {e}")
1691 })?;
1692 }
1693
1694 if !self.session_tasks.is_open() {
1695 self.await_session_tasks().await?;
1696 self.await_dispatch_task().await?;
1697 self.session_tasks.start_generation().map_err(|e| {
1698 anyhow::anyhow!("Failed to start Binance Futures session generation: {e}")
1699 })?;
1700 }
1701
1702 self.cancellation_token = CancellationToken::new();
1703 let cancellation_token = self.cancellation_token.clone();
1704 let ws_client = Arc::clone(&self.ws_client);
1705 let ws_trading_client = self.ws_trading_client.clone();
1706 let setup_guard =
1707 TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
1708 cancellation_token.cancel();
1709
1710 if let Some(client) = ws_client.lock().as_ref() {
1711 client.begin_shutdown();
1712 }
1713
1714 if let Some(client) = ws_trading_client {
1715 client.begin_shutdown();
1716 }
1717 });
1718
1719 let connect_result: anyhow::Result<()> = async {
1720 let is_hedge_mode = self
1722 .init_hedge_mode()
1723 .await
1724 .context("failed to query hedge mode")?;
1725 self.is_hedge_mode.store(is_hedge_mode, Ordering::Release);
1726 log::info!("Hedge mode (dual side position): {is_hedge_mode}");
1727 if is_hedge_mode != (self.core.oms_type == OmsType::Hedging) {
1728 log::warn!(
1729 "Binance Futures account position mode does not match the configured OMS type: \
1730 dual_side_position={is_hedge_mode}, oms_type={:?}; set oms_type to match the account position mode",
1731 self.core.oms_type,
1732 );
1733 }
1734
1735 let instruments = self
1736 .http_client
1737 .request_instruments_with_config(&self.config.instrument_provider)
1738 .await
1739 .context("failed to request Binance Futures instruments")?;
1740
1741 if instruments.is_empty() {
1742 log::warn!("No instruments returned for Binance Futures");
1743 } else {
1744 log::debug!("Loaded {} Futures instruments", instruments.len());
1745 }
1746 self.core.set_instruments_initialized();
1747
1748 self.apply_futures_config()
1750 .await
1751 .context("failed to apply futures config")?;
1752
1753 log::debug!("Creating listen key for user data stream...");
1755 let listen_key_response = self
1756 .http_client
1757 .create_listen_key()
1758 .await
1759 .context("failed to create listen key")?;
1760 let listen_key = listen_key_response.listen_key;
1761 log::debug!("Listen key created successfully");
1762
1763 {
1764 let mut key_guard = self.listen_key.write();
1765 *key_guard = Some(listen_key.clone());
1766 }
1767
1768 let (api_key, api_secret) = resolve_credentials(
1769 self.config.api_key.clone(),
1770 self.config.api_secret.clone(),
1771 self.config.environment,
1772 self.product_type,
1773 )?;
1774
1775 let private_base_url = self.config.base_url_ws.clone().map_or_else(
1776 || get_ws_private_base_url(self.product_type, self.config.environment).to_string(),
1777 |url| {
1778 if self.product_type == BinanceProductType::UsdM
1779 && self.config.environment == BinanceEnvironment::Live
1780 {
1781 get_usdm_ws_route_base_url(&url, "private")
1782 } else {
1783 url
1784 }
1785 },
1786 );
1787
1788 let (recovery_tx, recovery_rx) = tokio::sync::mpsc::unbounded_channel::<()>();
1789 self.recovery_tx = Some(recovery_tx.clone());
1790
1791 let seen_trade_ids: Arc<Mutex<FifoCache<(ustr::Ustr, i64), 10_000>>> =
1792 Arc::new(Mutex::new(FifoCache::new()));
1793
1794 let dispatch_ctx = Arc::new(DispatchCtx {
1795 emitter: self.emitter.clone(),
1796 http_client: self.http_client.clone(),
1797 account_id: self.core.account_id,
1798 product_type: self.product_type,
1799 clock: self.clock,
1800 dispatch_state: self.dispatch_state.clone(),
1801 triggered_algo_ids: self.triggered_algo_order_ids.clone(),
1802 algo_client_ids: self.algo_client_order_ids.clone(),
1803 use_position_ids: self.config.use_position_ids,
1804 default_taker_fee: self.config.default_taker_fee,
1805 bnfcr_currency: self.config.bnfcr_currency,
1806 treat_expired_as_canceled: self.config.treat_expired_as_canceled,
1807 use_trade_lite: self.config.use_trade_lite,
1808 seen_trade_ids,
1809 cancellation_token: self.cancellation_token.clone(),
1810 });
1811
1812 let ws_build_params = WsBuildParams {
1813 product_type: self.product_type,
1814 environment: self.config.environment,
1815 api_key: api_key.clone(),
1816 api_secret: api_secret.clone(),
1817 private_base_url: private_base_url.clone(),
1818 transport_backend: self.config.transport_backend,
1819 proxy_url: self.config.proxy_url.clone(),
1820 socket_factory: self.socket_factory.clone(),
1821 };
1822
1823 let ws_client = build_and_connect_user_stream(&ws_build_params, &listen_key).await?;
1824 let stream = ws_client.stream();
1825 *self.ws_client.lock() = Some(ws_client);
1826
1827 self.ws_task
1828 .lock()
1829 .await
1830 .spawn(run_user_stream_dispatch(
1831 stream,
1832 dispatch_ctx.clone(),
1833 recovery_tx.clone(),
1834 dispatch_user_stream_message,
1835 ))
1836 .map_err(|e| anyhow::anyhow!("failed to start user stream dispatch task: {e}"))?;
1837
1838 {
1840 let http_client = self.http_client.clone();
1841 let listen_key_ref = self.listen_key.clone();
1842 let cancel = self.cancellation_token.clone();
1843 let recovery_tx = recovery_tx.clone();
1844
1845 self.session_tasks.spawn(async move {
1846 let mut interval =
1847 tokio::time::interval(Duration::from_secs(LISTEN_KEY_KEEPALIVE_SECS));
1848 let mut consecutive_failures: u32 = 0;
1849
1850 loop {
1851 tokio::select! {
1852 _ = interval.tick() => {
1853 let key = {
1854 let guard = listen_key_ref.read();
1855 guard.clone()
1856 };
1857
1858 if let Some(ref key) = key {
1859 match http_client.keepalive_listen_key(key).await {
1860 Ok(()) => {
1861 log::debug!("Listen key keepalive sent successfully");
1862 consecutive_failures = 0;
1863 }
1864 Err(e) => {
1865 consecutive_failures += 1;
1866 log::warn!(
1867 "Listen key keepalive failed ({consecutive_failures}/{MAX_KEEPALIVE_FAILURES}): {e}",
1868 );
1869
1870 if consecutive_failures >= MAX_KEEPALIVE_FAILURES
1871 && recovery_tx.send(()).is_err()
1872 {
1873 log::warn!(
1874 "Recovery channel closed, keepalive exiting",
1875 );
1876 break;
1877 }
1878 }
1879 }
1880 }
1881 }
1882 () = cancel.cancelled() => {
1883 log::debug!("Listen key keepalive task cancelled");
1884 break;
1885 }
1886 }
1887 }
1888 })?;
1889 }
1890
1891 {
1893 let recovery_ctx = RecoveryCtx {
1894 http_client: self.http_client.clone(),
1895 listen_key: self.listen_key.clone(),
1896 recovery_listen_key: self.recovery_listen_key.clone(),
1897 ws_client: self.ws_client.clone(),
1898 ws_task: self.ws_task.clone(),
1899 recovery_lock: self.recovery_lock.clone(),
1900 ws_build_params,
1901 dispatch_ctx,
1902 recovery_tx: recovery_tx.clone(),
1903 };
1904 let cancel = self.cancellation_token.clone();
1905
1906 self.session_tasks.spawn(async move {
1907 run_recovery_driver(
1908 recovery_ctx,
1909 recovery_rx,
1910 cancel,
1911 dispatch_user_stream_message,
1912 )
1913 .await;
1914 })?;
1915 }
1916
1917 let account_state = self
1919 .refresh_account_state()
1920 .await
1921 .context("failed to request Binance Futures account state")?;
1922
1923 if !account_state.balances.is_empty() {
1924 log::debug!(
1925 "Received account state with {} balance(s) and {} margin(s)",
1926 account_state.balances.len(),
1927 account_state.margins.len()
1928 );
1929 }
1930
1931 self.emitter.send_account_state(account_state);
1932
1933 crate::common::execution::await_account_registered(&self.core, self.core.account_id, 30.0)
1934 .await?;
1935
1936 if let Some(ref mut ws_trading) = self.ws_trading_client {
1938 match ws_trading.connect().await {
1939 Ok(()) => {
1940 log::debug!("Connected to Binance Futures WS trading API");
1941
1942 let ws_trading_clone = ws_trading.clone();
1943 let emitter = self.emitter.clone();
1944 let account_id = self.core.account_id;
1945 let clock = self.clock;
1946 let dispatch_state = self.dispatch_state.clone();
1947
1948 self.session_tasks.spawn(async move {
1949 while let Some(msg) = ws_trading_clone.recv().await {
1950 dispatch_ws_trading_message(
1951 msg,
1952 &emitter,
1953 account_id,
1954 clock,
1955 &dispatch_state,
1956 );
1957 }
1958 })?;
1959 }
1960 Err(e) => {
1961 log::error!(
1962 "Failed to connect WS trading API: {e}. \
1963 Order operations will use HTTP fallback"
1964 );
1965 }
1966 }
1967 }
1968
1969 let refresh_secs = self.config.instrument_refresh_interval_secs;
1970 if refresh_secs > 0 {
1971 let http_client = self.http_client.clone();
1972 let provider = self.config.instrument_provider.clone();
1973
1974 self.session_tasks.spawn(async move {
1975 let mut interval = tokio::time::interval(Duration::from_secs(refresh_secs));
1976 interval.tick().await;
1977
1978 loop {
1979 interval.tick().await;
1980
1981 match http_client.request_instruments_with_config(&provider).await {
1982 Ok(instruments) => log::debug!(
1983 "Refreshed Binance Futures execution instruments: count={}",
1984 instruments.len()
1985 ),
1986 Err(e) => {
1987 log::warn!("Binance Futures execution instrument refresh failed: {e}");
1988 }
1989 }
1990 }
1991 })?;
1992 }
1993
1994 Ok(())
1995 }
1996 .await;
1997
1998 if let Err(e) = connect_result {
1999 self.recovery_tx.take();
2000 self.cancellation_token.cancel();
2001 self.abort_session_tasks();
2002 self.abort_pending_tasks();
2003
2004 if let Some(client) = self.ws_client.lock().as_ref() {
2005 client.begin_shutdown();
2006 }
2007
2008 if let Some(client) = self.ws_trading_client.as_ref() {
2009 client.begin_shutdown();
2010 }
2011
2012 if let Some(ref mut ws_trading) = self.ws_trading_client
2013 && let Err(e) = ws_trading.disconnect().await
2014 {
2015 self.shutdown_errors.push(format!(
2016 "Binance Futures trading WebSocket shutdown failed: {e}"
2017 ));
2018 }
2019
2020 let session_drained = match self.finish_session_tasks().await {
2021 Ok(()) => true,
2022 Err(e) => {
2023 let drained = matches!(e, TaskShutdownError::Join(_));
2024 self.shutdown_errors.push(format!(
2025 "Failed to terminate Binance Futures session tasks: {e}"
2026 ));
2027 drained
2028 }
2029 };
2030
2031 if session_drained && let Err(e) = self.await_dispatch_task().await {
2032 self.shutdown_errors.push(e.to_string());
2033 }
2034
2035 let ws_client = self.ws_client.lock().clone();
2036 if let Some(mut ws_client) = ws_client {
2037 match ws_client.close().await {
2038 Ok(()) => *self.ws_client.lock() = None,
2039 Err(e) => self.shutdown_errors.push(format!(
2040 "Binance Futures stream close after failed startup failed: {e}"
2041 )),
2042 }
2043 }
2044
2045 if let Err(e) = self
2046 .close_listen_key_slot(
2047 &self.listen_key,
2048 "failed to close listen key after failed startup",
2049 )
2050 .await
2051 {
2052 self.shutdown_errors.push(e.to_string());
2053 }
2054
2055 if let Err(e) = self
2056 .close_listen_key_slot(
2057 &self.recovery_listen_key,
2058 "failed to close recovery listen key after failed startup",
2059 )
2060 .await
2061 {
2062 self.shutdown_errors.push(e.to_string());
2063 }
2064
2065 if let Err(e) = self.await_pending_tasks().await {
2066 self.shutdown_errors.push(e.to_string());
2067 }
2068
2069 if self.shutdown_errors.is_empty() {
2070 return Err(e);
2071 }
2072 let shutdown_errors = std::mem::take(&mut self.shutdown_errors);
2073 return Err(e.context(format!(
2074 "Binance Futures startup teardown failed: {}",
2075 shutdown_errors.join("; ")
2076 )));
2077 }
2078
2079 setup_guard.disarm();
2080 self.core.set_connected();
2081 log::info!("Connected: client_id={}", self.core.client_id);
2082 Ok(())
2083 }
2084
2085 async fn disconnect(&mut self) -> anyhow::Result<()> {
2086 self.recovery_tx.take();
2088
2089 self.cancellation_token.cancel();
2091 self.abort_session_tasks();
2092 self.abort_pending_tasks();
2093
2094 if let Some(client) = self.ws_client.lock().as_ref() {
2095 client.begin_shutdown();
2096 }
2097
2098 if let Some(client) = self.ws_trading_client.as_ref() {
2099 client.begin_shutdown();
2100 }
2101
2102 if let Some(ref mut ws_trading) = self.ws_trading_client
2103 && let Err(e) = ws_trading.disconnect().await
2104 {
2105 self.shutdown_errors.push(format!(
2106 "Binance Futures trading WebSocket shutdown failed: {e}"
2107 ));
2108 }
2109
2110 let session_drained = match self.finish_session_tasks().await {
2111 Ok(()) => true,
2112 Err(e) => {
2113 let drained = matches!(e, TaskShutdownError::Join(_));
2114 self.shutdown_errors.push(format!(
2115 "Failed to terminate Binance Futures session tasks: {e}"
2116 ));
2117 drained
2118 }
2119 };
2120
2121 if session_drained && let Err(e) = self.await_dispatch_task().await {
2122 self.shutdown_errors.push(e.to_string());
2123 }
2124
2125 let ws_client = self.ws_client.lock().clone();
2127 if let Some(mut ws_client) = ws_client {
2128 match ws_client.close().await {
2129 Ok(()) => *self.ws_client.lock() = None,
2130 Err(e) => {
2131 self.shutdown_errors
2132 .push(format!("Binance Futures stream close failed: {e}"));
2133 }
2134 }
2135 }
2136
2137 if let Err(e) = self
2139 .close_listen_key_slot(&self.listen_key, "failed to close listen key")
2140 .await
2141 {
2142 self.shutdown_errors.push(e.to_string());
2143 }
2144
2145 if let Err(e) = self
2146 .close_listen_key_slot(
2147 &self.recovery_listen_key,
2148 "failed to close recovery listen key",
2149 )
2150 .await
2151 {
2152 self.shutdown_errors.push(e.to_string());
2153 }
2154
2155 if let Err(e) = self.await_pending_tasks().await {
2156 self.shutdown_errors.push(e.to_string());
2157 }
2158
2159 self.core.set_disconnected();
2160
2161 if !self.shutdown_errors.is_empty() {
2162 let shutdown_errors = std::mem::take(&mut self.shutdown_errors);
2163 anyhow::bail!(
2164 "Binance Futures shutdown failed: {}",
2165 shutdown_errors.join("; ")
2166 );
2167 }
2168 log::info!("Disconnected: client_id={}", self.core.client_id);
2169 Ok(())
2170 }
2171
2172 async fn generate_order_status_report(
2173 &self,
2174 cmd: &GenerateOrderStatusReport,
2175 ) -> anyhow::Result<Option<OrderStatusReport>> {
2176 let Some(instrument_id) = cmd.instrument_id else {
2177 log::warn!("generate_order_status_report requires instrument_id: {cmd:?}");
2178 return Ok(None);
2179 };
2180 let Some(instrument) = self.http_client.instrument_reconciliation(&instrument_id) else {
2181 if self.is_instrument_out_of_scope(instrument_id) {
2182 log::debug!(
2183 "Dropping out-of-scope historical Binance Futures order for instrument {instrument_id}"
2184 );
2185 } else {
2186 log::warn!(
2187 "Dropping historical Binance Futures order for unresolved instrument {instrument_id}"
2188 );
2189 }
2190 return Ok(None);
2191 };
2192
2193 let symbol = format_binance_symbol(&instrument_id);
2194 let order_id = cmd
2195 .venue_order_id
2196 .as_ref()
2197 .map(|id| {
2198 id.inner()
2199 .parse::<i64>()
2200 .context("failed to parse venue_order_id as numeric")
2201 })
2202 .transpose()?;
2203 let orig_client_order_id = cmd
2204 .client_order_id
2205 .map(|id| encode_broker_id(&id, BINANCE_NAUTILUS_FUTURES_BROKER_ID));
2206
2207 let mut builder = BinanceOrderQueryParamsBuilder::default();
2208 builder.symbol(symbol);
2209
2210 if let Some(oid) = order_id {
2211 builder.order_id(oid);
2212 }
2213
2214 if let Some(ref coid) = orig_client_order_id {
2215 builder.orig_client_order_id(coid.clone());
2216 }
2217 let params = builder.build().map_err(|e| anyhow::anyhow!("{e}"))?;
2218
2219 let price_precision = instrument.price_precision();
2220 let size_precision = instrument.size_precision();
2221 let ts_init = self.clock.get_time_ns();
2222 let algo_lookup = self.resolve_algo_lookup(cmd.client_order_id, cmd.params.as_ref());
2223
2224 if algo_lookup == BinanceFuturesAlgoLookup::AlgoId {
2225 let algo_order = self
2226 .http_client
2227 .query_algo_order_with_history(
2228 instrument_id,
2229 cmd.client_order_id,
2230 cmd.venue_order_id,
2231 )
2232 .await?;
2233
2234 return match algo_order {
2235 Some(result) => Ok(Some(create_algo_order_status_report(
2236 &result,
2237 self.core.account_id,
2238 instrument_id,
2239 price_precision,
2240 size_precision,
2241 self.config.treat_expired_as_canceled,
2242 self.config.use_position_ids,
2243 ts_init,
2244 )?)),
2245 None => {
2246 log::debug!("Algo order query returned no matching order");
2247 Ok(None)
2248 }
2249 };
2250 }
2251
2252 match self.http_client.query_order(¶ms).await {
2253 Ok(order) => {
2254 let report = order.to_order_status_report(
2255 self.core.account_id,
2256 instrument_id,
2257 price_precision,
2258 size_precision,
2259 self.config.treat_expired_as_canceled,
2260 ts_init,
2261 )?;
2262 let venue_position_id = make_venue_position_id(
2263 self.config.use_position_ids,
2264 instrument_id,
2265 order.position_side,
2266 )?;
2267 Ok(Some(with_venue_position_id(report, venue_position_id)))
2268 }
2269 Err(BinanceFuturesHttpError::BinanceError { code: -2013, .. }) => {
2270 if algo_lookup == BinanceFuturesAlgoLookup::Skip {
2271 log::debug!("Skipping Algo Service fallback for known regular order");
2272 return Ok(None);
2273 }
2274
2275 let algo_venue_order_id = if algo_lookup == BinanceFuturesAlgoLookup::AlgoId {
2280 cmd.venue_order_id
2281 } else {
2282 None
2283 };
2284 let algo_order = self
2285 .http_client
2286 .query_algo_order_with_history(
2287 instrument_id,
2288 cmd.client_order_id,
2289 algo_venue_order_id,
2290 )
2291 .await?;
2292
2293 match algo_order {
2294 Some(result) => Ok(Some(create_algo_order_status_report(
2295 &result,
2296 self.core.account_id,
2297 instrument_id,
2298 price_precision,
2299 size_precision,
2300 self.config.treat_expired_as_canceled,
2301 self.config.use_position_ids,
2302 ts_init,
2303 )?)),
2304 None => {
2305 log::debug!("Algo order query returned no matching order");
2306 Ok(None)
2307 }
2308 }
2309 }
2310 Err(e) => Err(e.into()),
2311 }
2312 }
2313
2314 async fn generate_order_status_reports(
2315 &self,
2316 cmd: &GenerateOrderStatusReports,
2317 ) -> anyhow::Result<Vec<OrderStatusReport>> {
2318 let ts_init = self.clock.get_time_ns();
2319
2320 if cmd.open_only {
2321 let mut reports = self
2322 .generate_open_order_status_reports(cmd.instrument_id, ts_init)
2323 .await?;
2324
2325 if reports.iter().any(|report| {
2326 report.quantity_free_close_position_side.is_some()
2327 && !report.report.quantity.is_positive()
2328 }) {
2329 let position_cmd = GeneratePositionStatusReportsBuilder::default()
2330 .ts_init(ts_init)
2331 .instrument_id(cmd.instrument_id)
2332 .build()
2333 .map_err(|e| anyhow::anyhow!("{e}"))?;
2334 let position_reports = self.generate_position_status_reports(&position_cmd).await?;
2335 restore_close_position_quantities(&mut reports, &position_reports);
2336 }
2337
2338 return Ok(reports.into_iter().map(|report| report.report).collect());
2339 }
2340
2341 let mut reports = Vec::new();
2342
2343 if let Some(instrument_id) = cmd.instrument_id
2344 && self
2345 .http_client
2346 .instrument_reconciliation(&instrument_id)
2347 .is_none()
2348 {
2349 if self.is_instrument_out_of_scope(instrument_id) {
2350 log::debug!(
2351 "Dropping out-of-scope Binance Futures order request for instrument {instrument_id}"
2352 );
2353 return Ok(reports);
2354 }
2355
2356 log::warn!(
2357 "Dropping historical Binance Futures orders for unresolved instrument {instrument_id}"
2358 );
2359 return Ok(reports);
2360 }
2361
2362 if let Some(instrument_id) = cmd.instrument_id {
2363 let Some(instrument) = self.http_client.instrument_reconciliation(&instrument_id)
2364 else {
2365 if self.is_instrument_out_of_scope(instrument_id) {
2366 log::debug!(
2367 "Dropping out-of-scope historical Binance Futures orders for instrument {instrument_id}"
2368 );
2369 } else {
2370 log::warn!(
2371 "Dropping historical Binance Futures orders for unresolved instrument {instrument_id}"
2372 );
2373 }
2374 return Ok(reports);
2375 };
2376 let symbol = format_binance_symbol(&instrument_id);
2377 let start_time = cmd
2378 .start
2379 .map(|t| t.as_i64() / NANOSECONDS_IN_MILLISECOND as i64);
2380 let end_time = cmd
2381 .end
2382 .map(|t| t.as_i64() / NANOSECONDS_IN_MILLISECOND as i64);
2383
2384 let mut builder = BinanceAllOrdersParamsBuilder::default();
2385 builder.symbol(symbol);
2386
2387 if let Some(st) = start_time {
2388 builder.start_time(st);
2389 }
2390
2391 if let Some(et) = end_time {
2392 builder.end_time(et);
2393 }
2394 let params = builder.build().map_err(|e| anyhow::anyhow!("{e}"))?;
2395
2396 let orders = self.http_client.query_all_orders(¶ms).await?;
2397
2398 for order in orders {
2399 if let Ok(report) = order.to_order_status_report(
2400 self.core.account_id,
2401 instrument.id(),
2402 instrument.price_precision(),
2403 instrument.size_precision(),
2404 self.config.treat_expired_as_canceled,
2405 ts_init,
2406 ) {
2407 let venue_position_id = make_venue_position_id(
2408 self.config.use_position_ids,
2409 instrument.id(),
2410 order.position_side,
2411 )?;
2412 reports.push(with_venue_position_id(report, venue_position_id));
2413 }
2414 }
2415 }
2416
2417 Ok(reports)
2418 }
2419
2420 async fn generate_fill_reports(
2421 &self,
2422 cmd: GenerateFillReports,
2423 ) -> anyhow::Result<Vec<FillReport>> {
2424 let Some(instrument_id) = cmd.instrument_id else {
2425 log::warn!("generate_fill_reports requires instrument_id for Binance Futures");
2426 return Ok(Vec::new());
2427 };
2428 let Some(instrument) = self.http_client.instrument_reconciliation(&instrument_id) else {
2429 if self.is_instrument_out_of_scope(instrument_id) {
2430 log::debug!(
2431 "Dropping out-of-scope historical Binance Futures fills for instrument {instrument_id}"
2432 );
2433 } else {
2434 log::warn!(
2435 "Dropping historical Binance Futures fills for unresolved instrument {instrument_id}"
2436 );
2437 }
2438 return Ok(Vec::new());
2439 };
2440
2441 let symbol = format_binance_symbol(&instrument_id);
2442 let mut trades = Vec::new();
2443 let mut seen_trade_ids = AHashSet::new();
2444 let requested_end_time = cmd
2445 .end
2446 .map(|end| end.as_i64() / NANOSECONDS_IN_MILLISECOND as i64);
2447
2448 if let Some(start) = cmd.start {
2449 let query_start_time = start.as_i64() / NANOSECONDS_IN_MILLISECOND as i64;
2450 let query_end_time = requested_end_time.unwrap_or_else(|| {
2451 self.clock.get_time_ns().as_i64() / NANOSECONDS_IN_MILLISECOND as i64
2452 });
2453 anyhow::ensure!(
2454 query_start_time <= query_end_time,
2455 "fill report start time must not exceed end time"
2456 );
2457 let complete_start = user_trades_complete_start(cmd.ts_init, self.clock.get_time_ns());
2458 anyhow::ensure!(
2459 start >= complete_start,
2460 "Binance Futures fill report range is incomplete: start {start} precedes complete-history boundary {complete_start}"
2461 );
2462 let mut window_start = query_start_time;
2463
2464 loop {
2465 let window_end = window_start
2466 .saturating_add(USER_TRADES_MAX_INTERVAL_MS)
2467 .min(query_end_time);
2468 let mut from_id = None;
2469
2470 loop {
2471 let mut builder = BinanceUserTradesParamsBuilder::default();
2472 builder.symbol(symbol.clone());
2473 builder.limit(USER_TRADES_PAGE_LIMIT);
2474
2475 if let Some(cursor) = from_id {
2476 builder.from_id(cursor);
2477 } else {
2478 builder.start_time(window_start);
2479 builder.end_time(window_end);
2480 }
2481 let params = builder.build().map_err(|e| anyhow::anyhow!("{e}"))?;
2482 let page = self.http_client.query_user_trades(¶ms).await?;
2483
2484 if page.is_empty() {
2485 break;
2486 }
2487
2488 let page_len = page.len();
2489 let max_trade_id = page.iter().map(|trade| trade.id).max().unwrap();
2490 let passed_window_end = page.iter().any(|trade| trade.time > window_end);
2491
2492 trades.extend(page.into_iter().filter(|trade| {
2493 trade.time >= window_start
2494 && trade.time <= window_end
2495 && seen_trade_ids.insert(trade.id)
2496 }));
2497
2498 if page_len < USER_TRADES_PAGE_LIMIT as usize || passed_window_end {
2499 break;
2500 }
2501
2502 let next_from_id = max_trade_id
2503 .checked_add(1)
2504 .context("Binance user trade ID overflow during pagination")?;
2505 anyhow::ensure!(
2506 from_id.is_none_or(|cursor| next_from_id > cursor),
2507 "Binance user-trades pagination made no progress"
2508 );
2509 from_id = Some(next_from_id);
2510 }
2511
2512 if window_end >= query_end_time {
2513 break;
2514 }
2515 window_start = window_end.saturating_add(1);
2516 }
2517 } else {
2518 anyhow::ensure!(
2519 cmd.end.is_none(),
2520 "Binance Futures fill report end time requires start time for a complete range"
2521 );
2522 let mut from_id = 0;
2523
2524 loop {
2525 let params = BinanceUserTradesParamsBuilder::default()
2526 .symbol(symbol.clone())
2527 .from_id(from_id)
2528 .limit(USER_TRADES_PAGE_LIMIT)
2529 .build()
2530 .map_err(|e| anyhow::anyhow!("{e}"))?;
2531 let page = self.http_client.query_user_trades(¶ms).await?;
2532
2533 if page.is_empty() {
2534 break;
2535 }
2536
2537 let page_len = page.len();
2538 let max_trade_id = page.iter().map(|trade| trade.id).max().unwrap();
2539 let passed_end = requested_end_time
2540 .is_some_and(|end_time| page.iter().any(|trade| trade.time > end_time));
2541
2542 trades.extend(page.into_iter().filter(|trade| {
2543 requested_end_time.is_none_or(|end_time| trade.time <= end_time)
2544 && seen_trade_ids.insert(trade.id)
2545 }));
2546
2547 if page_len < USER_TRADES_PAGE_LIMIT as usize || passed_end {
2548 break;
2549 }
2550
2551 let next_from_id = max_trade_id
2552 .checked_add(1)
2553 .context("Binance user trade ID overflow during pagination")?;
2554 anyhow::ensure!(
2555 next_from_id > from_id,
2556 "Binance user-trades pagination made no progress"
2557 );
2558 from_id = next_from_id;
2559 }
2560 }
2561
2562 trades.sort_unstable_by_key(|trade| (trade.time, trade.id));
2563 let ts_init = self.clock.get_time_ns();
2564
2565 let mut reports = Vec::new();
2566
2567 for trade in trades {
2568 let venue_position_id = make_venue_position_id(
2569 self.config.use_position_ids,
2570 instrument.id(),
2571 trade.position_side,
2572 )?;
2573 let mut report = trade.to_fill_report(
2574 self.core.account_id,
2575 instrument.id(),
2576 instrument.price_precision(),
2577 instrument.size_precision(),
2578 self.config.bnfcr_currency,
2579 ts_init,
2580 )?;
2581 report.venue_position_id = venue_position_id;
2582 reports.push(report);
2583 }
2584
2585 Ok(reports)
2586 }
2587
2588 async fn generate_position_status_reports(
2589 &self,
2590 cmd: &GeneratePositionStatusReports,
2591 ) -> anyhow::Result<Vec<PositionStatusReport>> {
2592 if let Some(instrument_id) = cmd.instrument_id
2593 && self
2594 .http_client
2595 .instrument_reconciliation(&instrument_id)
2596 .is_none()
2597 {
2598 if self.is_instrument_out_of_scope(instrument_id) {
2599 log::debug!(
2600 "Dropping out-of-scope Binance Futures position request for instrument {instrument_id}"
2601 );
2602 return Ok(Vec::new());
2603 }
2604 anyhow::bail!(
2605 "Binance Futures position request has unresolved instrument {instrument_id}"
2606 );
2607 }
2608 let symbol = cmd.instrument_id.map(|id| format_binance_symbol(&id));
2609
2610 let mut builder = BinancePositionRiskParamsBuilder::default();
2611
2612 if let Some(s) = symbol {
2613 builder.symbol(s);
2614 }
2615 let params = builder.build().map_err(|e| anyhow::anyhow!("{e}"))?;
2616
2617 let positions = self.http_client.query_positions(¶ms).await?;
2618
2619 let mut reports = Vec::new();
2620 let mut position_reports_failed = 0usize;
2621
2622 for position in positions {
2623 let instrument_id = format_instrument_id(&position.symbol, self.product_type);
2624
2625 if self.is_instrument_out_of_scope(instrument_id) {
2626 log::debug!(
2627 "Dropping out-of-scope Binance Futures position for instrument {instrument_id}"
2628 );
2629 continue;
2630 }
2631
2632 let position_amt = match position.position_amt.parse::<Decimal>() {
2633 Ok(value) => value,
2634 Err(e) => {
2635 log::warn!(
2636 "Failed to parse Futures position_amt for symbol={}: {e}",
2637 position.symbol
2638 );
2639 position_reports_failed += 1;
2640 continue;
2641 }
2642 };
2643
2644 if position_amt.is_zero() {
2645 if self.config.use_position_ids {
2646 let position_side = match position.position_side {
2647 Some(BinancePositionSide::Long) => PositionSide::Long,
2648 Some(BinancePositionSide::Short) => PositionSide::Short,
2649 _ => continue,
2650 };
2651
2652 let venue_position_id =
2653 make_venue_position_id(true, instrument_id, position.position_side)?
2654 .expect("hedge position sides always produce an ID");
2655
2656 if let Err(e) = self.ensure_cached_position_id_compatible(
2657 instrument_id,
2658 position_side,
2659 venue_position_id,
2660 ) {
2661 log::warn!(
2662 "Failed to create Futures position report for symbol={}: {e}",
2663 position.symbol
2664 );
2665 position_reports_failed += 1;
2666 }
2667 }
2668 continue;
2669 }
2670
2671 let Some(instrument) = self.http_client.instrument_reconciliation(&instrument_id)
2672 else {
2673 log::warn!(
2674 "Failed to create Futures position report for symbol={}: instrument {instrument_id} is unresolved",
2675 position.symbol
2676 );
2677 position_reports_failed += 1;
2678 continue;
2679 };
2680
2681 match self.create_position_report(
2682 &position,
2683 instrument.id(),
2684 instrument.size_precision(),
2685 ) {
2686 Ok(report) => reports.push(report),
2687 Err(e) => {
2688 log::warn!(
2689 "Failed to create Futures position report for symbol={}: {e}",
2690 position.symbol
2691 );
2692 position_reports_failed += 1;
2693 }
2694 }
2695 }
2696
2697 anyhow::ensure!(
2698 position_reports_failed == 0,
2699 "Failed to process {position_reports_failed} Binance Futures position reports",
2700 );
2701
2702 Ok(reports)
2703 }
2704
2705 async fn generate_mass_status(
2706 &self,
2707 lookback_mins: Option<u64>,
2708 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
2709 log::info!("Generating ExecutionMassStatus (lookback_mins={lookback_mins:?})");
2710
2711 let ts_now = self.clock.get_time_ns();
2712
2713 let requested_start = if let Some(mins) = lookback_mins {
2714 let lookback_ns = checked_mins_to_nanos(mins)
2715 .context("lookback minutes exceed the nanosecond range")?;
2716 Some(UnixNanos::from(ts_now.as_u64().saturating_sub(lookback_ns)))
2717 } else {
2718 None
2719 };
2720 let complete_start = user_trades_complete_start(ts_now, ts_now);
2721 let (report_start, fill_start, mut reports_complete) = match requested_start {
2722 Some(start) if start < complete_start => (complete_start, Some(complete_start), false),
2723 Some(start) => (start, Some(start), true),
2724 None => (complete_start, None, false),
2725 };
2726
2727 let position_cmd = GeneratePositionStatusReportsBuilder::default()
2728 .ts_init(ts_now)
2729 .start(Some(report_start))
2730 .build()
2731 .map_err(|e| anyhow::anyhow!("{e}"))?;
2732
2733 let (mut open_order_reports, position_reports) = tokio::try_join!(
2734 self.generate_open_order_status_reports(None, ts_now),
2735 self.generate_position_status_reports(&position_cmd),
2736 )?;
2737 restore_close_position_quantities(&mut open_order_reports, &position_reports);
2738 let order_reports: Vec<_> = open_order_reports
2739 .into_iter()
2740 .map(|report| report.report)
2741 .collect();
2742
2743 let mut instrument_ids: Vec<_> = order_reports
2744 .iter()
2745 .map(|report| report.instrument_id)
2746 .chain(position_reports.iter().map(|report| report.instrument_id))
2747 .collect();
2748 {
2749 let cache = self.core.cache();
2750 instrument_ids.extend(
2751 cache
2752 .orders_open(
2753 Some(&BINANCE_VENUE),
2754 None,
2755 None,
2756 Some(&self.core.account_id),
2757 None,
2758 )
2759 .into_iter()
2760 .map(|order| order.instrument_id()),
2761 );
2762 instrument_ids.extend(
2763 cache
2764 .orders_inflight(
2765 Some(&BINANCE_VENUE),
2766 None,
2767 None,
2768 Some(&self.core.account_id),
2769 None,
2770 )
2771 .into_iter()
2772 .map(|order| order.instrument_id()),
2773 );
2774 instrument_ids.extend(
2775 cache
2776 .positions_open(
2777 Some(&BINANCE_VENUE),
2778 None,
2779 None,
2780 Some(&self.core.account_id),
2781 None,
2782 )
2783 .into_iter()
2784 .map(|position| position.instrument_id),
2785 );
2786 instrument_ids.retain(|instrument_id| {
2787 cache.instrument(instrument_id).is_none_or(|instrument| {
2788 is_instrument_for_product(instrument, self.product_type)
2789 })
2790 });
2791 }
2792 instrument_ids.sort_unstable();
2793 instrument_ids.dedup();
2794
2795 let mut fill_reports = Vec::new();
2796
2797 for instrument_id in instrument_ids {
2798 if self
2799 .http_client
2800 .instrument_reconciliation(&instrument_id)
2801 .is_none()
2802 {
2803 if self.is_instrument_out_of_scope(instrument_id) {
2804 log::debug!(
2805 "Dropping out-of-scope historical Binance Futures fills for instrument {instrument_id}"
2806 );
2807 } else {
2808 log::warn!(
2809 "Dropping historical Binance Futures fills for unresolved instrument {instrument_id}"
2810 );
2811 reports_complete = false;
2812 }
2813 continue;
2814 }
2815 let fill_cmd = GenerateFillReportsBuilder::default()
2816 .ts_init(ts_now)
2817 .instrument_id(Some(instrument_id))
2818 .start(fill_start)
2819 .build()
2820 .map_err(|e| anyhow::anyhow!("{e}"))?;
2821 fill_reports.extend(
2822 self.generate_fill_reports(fill_cmd)
2823 .await?
2824 .into_iter()
2825 .filter(|report| report.ts_event >= report_start),
2826 );
2827 }
2828
2829 log::info!("Received {} OrderStatusReports", order_reports.len());
2830 log::info!("Received {} FillReports", fill_reports.len());
2831 log::info!("Received {} PositionReports", position_reports.len());
2832
2833 let mut mass_status = ExecutionMassStatus::new(
2834 self.core.client_id,
2835 self.core.account_id,
2836 *BINANCE_VENUE,
2837 ts_now,
2838 None,
2839 );
2840
2841 mass_status.add_order_reports(order_reports);
2842 mass_status.add_fill_reports(fill_reports);
2843 mass_status.add_position_reports(position_reports);
2844 mass_status.set_report_window(Some(report_start), reports_complete);
2845
2846 Ok(Some(mass_status))
2847 }
2848
2849 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
2850 self.update_account_state();
2851 Ok(())
2852 }
2853
2854 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
2855 log::debug!("query_order: client_order_id={}", cmd.client_order_id);
2856
2857 let algo_lookup = self.resolve_algo_lookup(Some(cmd.client_order_id), cmd.params.as_ref());
2858 let http_client = self.http_client.clone();
2859 let command = cmd;
2860 let emitter = self.emitter.clone();
2861 let account_id = self.core.account_id;
2862 let clock = self.clock;
2863
2864 let symbol = format_binance_symbol(&command.instrument_id);
2865 let order_id = command
2866 .venue_order_id
2867 .map(|id| {
2868 id.inner()
2869 .parse::<i64>()
2870 .map_err(|e| anyhow::anyhow!("failed to parse venue_order_id: {e}"))
2871 })
2872 .transpose()?;
2873 let orig_client_order_id = Some(encode_broker_id(
2874 &command.client_order_id,
2875 BINANCE_NAUTILUS_FUTURES_BROKER_ID,
2876 ));
2877 let (price_precision, size_precision) =
2878 self.get_instrument_precision(command.instrument_id)?;
2879 let treat_expired_as_canceled = self.config.treat_expired_as_canceled;
2880 let use_position_ids = self.config.use_position_ids;
2881
2882 self.spawn_task("query_order", async move {
2883 if algo_lookup == BinanceFuturesAlgoLookup::AlgoId {
2884 match http_client
2885 .query_algo_order_with_history(
2886 command.instrument_id,
2887 Some(command.client_order_id),
2888 command.venue_order_id,
2889 )
2890 .await
2891 {
2892 Ok(Some(result)) => {
2893 let report = create_algo_order_status_report(
2894 &result,
2895 account_id,
2896 command.instrument_id,
2897 price_precision,
2898 size_precision,
2899 treat_expired_as_canceled,
2900 use_position_ids,
2901 clock.get_time_ns(),
2902 )?;
2903 emitter.send_order_status_report(report);
2904 }
2905 Ok(None) => log::warn!("Algo order query returned no matching order"),
2906 Err(e) => log::warn!("Failed to query algo order status: {e}"),
2907 }
2908
2909 return Ok(());
2910 }
2911
2912 let mut builder = BinanceOrderQueryParamsBuilder::default();
2913 builder.symbol(symbol.clone());
2914
2915 if let Some(oid) = order_id {
2916 builder.order_id(oid);
2917 }
2918
2919 if let Some(coid) = orig_client_order_id {
2920 builder.orig_client_order_id(coid);
2921 }
2922 let params = builder
2923 .build()
2924 .map_err(|e| anyhow::anyhow!("failed to build order query params: {e}"))?;
2925
2926 let result = http_client.query_order(¶ms).await;
2927
2928 match result {
2929 Ok(order) => {
2930 let ts_init = clock.get_time_ns();
2931 let report = order.to_order_status_report(
2932 account_id,
2933 command.instrument_id,
2934 price_precision,
2935 size_precision,
2936 treat_expired_as_canceled,
2937 ts_init,
2938 )?;
2939 let venue_position_id = make_venue_position_id(
2940 use_position_ids,
2941 command.instrument_id,
2942 order.position_side,
2943 )?;
2944
2945 emitter.send_order_status_report(with_venue_position_id(
2946 report,
2947 venue_position_id,
2948 ));
2949 }
2950 Err(BinanceFuturesHttpError::BinanceError { code: -2013, .. }) => {
2951 if algo_lookup == BinanceFuturesAlgoLookup::Skip {
2952 log::debug!("Skipping Algo Service fallback for known regular order");
2953 return Ok(());
2954 }
2955
2956 let algo_venue_order_id = if algo_lookup == BinanceFuturesAlgoLookup::AlgoId {
2963 command.venue_order_id
2964 } else {
2965 None
2966 };
2967
2968 match http_client
2969 .query_algo_order_with_history(
2970 command.instrument_id,
2971 Some(command.client_order_id),
2972 algo_venue_order_id,
2973 )
2974 .await
2975 {
2976 Ok(Some(result)) => {
2977 let report = create_algo_order_status_report(
2978 &result,
2979 account_id,
2980 command.instrument_id,
2981 price_precision,
2982 size_precision,
2983 treat_expired_as_canceled,
2984 use_position_ids,
2985 clock.get_time_ns(),
2986 )?;
2987 emitter.send_order_status_report(report);
2988 }
2989 Ok(None) => log::warn!("Algo order query returned no matching order"),
2990 Err(e) => log::warn!("Algo order query also failed: {e}"),
2991 }
2992 }
2993 Err(e) => log::warn!("Failed to query order status: {e}"),
2994 }
2995
2996 Ok(())
2997 });
2998
2999 Ok(())
3000 }
3001
3002 fn generate_account_state(
3003 &self,
3004 balances: Vec<AccountBalance>,
3005 margins: Vec<MarginBalance>,
3006 reported: bool,
3007 ts_event: UnixNanos,
3008 info: Option<Params>,
3009 ) -> anyhow::Result<()> {
3010 self.emitter
3011 .emit_account_state(balances, margins, reported, ts_event, info);
3012 Ok(())
3013 }
3014
3015 fn start(&mut self) -> anyhow::Result<()> {
3016 if self.core.is_started() {
3017 return Ok(());
3018 }
3019
3020 self.emitter.set_sender(get_exec_event_sender());
3021 self.core.set_started();
3022
3023 log::info!(
3024 "Started: client_id={}, account_id={}, account_type={:?}, environment={:?}",
3025 self.core.client_id,
3026 self.core.account_id,
3027 self.core.account_type,
3028 self.config.environment,
3029 );
3030 Ok(())
3031 }
3032
3033 fn stop(&mut self) -> anyhow::Result<()> {
3034 if self.core.is_stopped() {
3035 return Ok(());
3036 }
3037
3038 self.cancellation_token.cancel();
3039 self.abort_session_tasks();
3040
3041 if let Some(client) = self.ws_client.lock().as_ref() {
3042 client.begin_shutdown();
3043 }
3044
3045 if let Some(client) = self.ws_trading_client.as_ref() {
3046 client.begin_shutdown();
3047 }
3048
3049 if let Ok(mut task_slot) = self.ws_task.try_lock() {
3050 task_slot.abort();
3051 }
3052
3053 self.recovery_tx.take();
3054
3055 self.abort_pending_tasks();
3056 self.core.set_stopped();
3057 self.core.set_disconnected();
3058 log::info!("Stopped: client_id={}", self.core.client_id);
3059 Ok(())
3060 }
3061
3062 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
3063 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
3064
3065 if order.is_closed() {
3066 let client_order_id = order.client_order_id();
3067 log::warn!("Cannot submit closed order {client_order_id}");
3068 return Ok(());
3069 }
3070
3071 if let Some(offset_type) = order.trailing_offset_type() {
3074 if offset_type != TrailingOffsetType::BasisPoints {
3075 anyhow::bail!(
3076 "Binance only supports TrailingOffsetType::BasisPoints, received {offset_type:?}"
3077 );
3078 }
3079
3080 if let Some(offset) = order.trailing_offset() {
3081 trailing_offset_to_callback_rate(offset)?;
3082 }
3083 }
3084
3085 let close_position = cmd
3086 .params
3087 .as_ref()
3088 .and_then(|p| p.get_bool(PARAMS_CLOSE_POSITION))
3089 .unwrap_or(false);
3090
3091 if close_position {
3092 let order_type = order.order_type();
3093
3094 if !matches!(
3095 order_type,
3096 OrderType::StopMarket | OrderType::MarketIfTouched
3097 ) {
3098 anyhow::bail!(
3099 "`close_position` is not supported for order type {order_type:?} on Binance"
3100 );
3101 }
3102
3103 if order.is_reduce_only() {
3104 anyhow::bail!("`close_position` cannot be combined with `reduce_only` on Binance");
3105 }
3106 }
3107
3108 if let Some(pm_str) = cmd.params.as_ref().and_then(|p| p.get_str("price_match")) {
3109 BinancePriceMatch::from_param(pm_str)?;
3110 let order_type = order.order_type();
3111 anyhow::ensure!(
3112 !order.is_post_only(),
3113 "price_match cannot be combined with post-only orders"
3114 );
3115 anyhow::ensure!(
3116 order_type == OrderType::Limit,
3117 "price_match is not supported for order type {order_type:?}"
3118 );
3119 }
3120
3121 let lifetime = determine_futures_order_lifetime(
3122 self.product_type,
3123 order.order_type(),
3124 order.time_in_force(),
3125 order.expire_time(),
3126 order.is_post_only(),
3127 self.config.use_gtd,
3128 self.clock.get_time_ns(),
3129 )?;
3130
3131 let (position_side, venue_position_id) = resolve_order_position_identity(
3132 self.is_hedge_mode(),
3133 self.config.use_position_ids,
3134 &order,
3135 close_position,
3136 )?;
3137 validate_submit_position_id(cmd.position_id, venue_position_id)?;
3138
3139 log::debug!("OrderSubmitted client_order_id={}", order.client_order_id());
3140 self.emitter.emit_order_submitted(&order);
3141
3142 self.submit_order_internal(&cmd, lifetime, position_side, venue_position_id)
3143 }
3144
3145 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
3146 if cmd.order_list.client_order_ids.is_empty() {
3147 log::debug!("submit_order_list called with empty order list");
3148 return Ok(());
3149 }
3150
3151 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
3152
3153 if let Some(order) = orders.iter().find(|order| order.is_closed()) {
3154 let reason = format!("Cannot submit closed order {}", order.client_order_id());
3155 for order in &orders {
3156 self.emitter.emit_order_denied(order, &reason);
3157 }
3158 return Ok(());
3159 }
3160
3161 let close_position = cmd
3162 .params
3163 .as_ref()
3164 .and_then(|p| p.get_bool(PARAMS_CLOSE_POSITION))
3165 .unwrap_or(false);
3166 let price_match = match cmd
3167 .params
3168 .as_ref()
3169 .and_then(|p| p.get_str("price_match"))
3170 .map(BinancePriceMatch::from_param)
3171 .transpose()
3172 {
3173 Ok(price_match) => price_match,
3174 Err(e) => {
3175 for order in &orders {
3176 self.emitter.emit_order_denied(order, &e.to_string());
3177 }
3178 return Ok(());
3179 }
3180 };
3181
3182 let batch_items = match build_futures_order_list_batch(
3183 &orders,
3184 self.is_hedge_mode(),
3185 close_position,
3186 price_match,
3187 self.product_type,
3188 self.config.use_gtd,
3189 self.clock.get_time_ns(),
3190 ) {
3191 Ok(batch_items) => batch_items,
3192 Err(reason) => {
3193 for order in &orders {
3194 self.emitter.emit_order_denied(order, &reason);
3195 }
3196 return Ok(());
3197 }
3198 };
3199
3200 let venue_position_ids = orders
3201 .iter()
3202 .map(|order| {
3203 resolve_order_position_identity(
3204 self.is_hedge_mode(),
3205 self.config.use_position_ids,
3206 order,
3207 close_position,
3208 )
3209 .map(|(_, venue_position_id)| venue_position_id)
3210 })
3211 .collect::<anyhow::Result<Vec<_>>>()?;
3212
3213 for venue_position_id in &venue_position_ids {
3214 validate_submit_position_id(cmd.position_id, *venue_position_id)?;
3215 }
3216
3217 for (order, venue_position_id) in orders.iter().zip(venue_position_ids) {
3218 self.dispatch_state.order_identities.insert(
3219 order.client_order_id(),
3220 OrderIdentity {
3221 instrument_id: order.instrument_id(),
3222 strategy_id: order.strategy_id(),
3223 order_side: order.order_side(),
3224 order_type: order.order_type(),
3225 price: order.price(),
3226 quantity: order.quantity(),
3227 venue_position_id,
3228 },
3229 );
3230 self.emitter.emit_order_submitted(order);
3231 }
3232
3233 let http_client = self.http_client.clone();
3234 let emitter = self.emitter.clone();
3235 let trader_id = self.core.trader_id;
3236 let account_id = self.core.account_id;
3237 let clock = self.clock;
3238
3239 self.spawn_task("submit_order_list", async move {
3240 match http_client.submit_order_list(&batch_items).await {
3241 Ok(results) => {
3242 for (order, result) in orders.iter().zip(results.iter()) {
3243 match result {
3244 BatchOrderResult::Success(response) => {
3245 log::debug!(
3246 "Order-list leg submit accepted: client_order_id={}, venue_order_id={}",
3247 order.client_order_id(),
3248 response.order_id
3249 );
3250 }
3251 BatchOrderResult::Error(error) => {
3252 let ts_now = clock.get_time_ns();
3253 let rejected = OrderRejected::new(
3254 trader_id,
3255 order.strategy_id(),
3256 order.instrument_id(),
3257 order.client_order_id(),
3258 account_id,
3259 format!(
3260 "submit-order-list-error: code={}, msg={}",
3261 error.code, error.msg
3262 )
3263 .into(),
3264 UUID4::new(),
3265 ts_now,
3266 ts_now,
3267 false,
3268 false,
3269 );
3270 emitter.send_order_event(OrderEventAny::Rejected(rejected));
3271 }
3272 }
3273 }
3274 }
3275 Err(e) => {
3276 let e = anyhow::Error::new(e);
3277
3278 if is_ambiguous_submit_error(&e) {
3279 log::warn!(
3280 "Ambiguous order-list submit failure, awaiting reconciliation: {e}"
3281 );
3282 } else if is_structured_venue_rejection(&e) {
3283 let ts_now = clock.get_time_ns();
3284 let due_post_only = classify_submit_order_error(&e);
3285
3286 for order in &orders {
3287 let rejected = OrderRejected::new(
3288 trader_id,
3289 order.strategy_id(),
3290 order.instrument_id(),
3291 order.client_order_id(),
3292 account_id,
3293 format!("submit-order-list-error: {e}").into(),
3294 UUID4::new(),
3295 ts_now,
3296 ts_now,
3297 false,
3298 due_post_only,
3299 );
3300 emitter.send_order_event(OrderEventAny::Rejected(rejected));
3301 }
3302 } else if is_local_command_failure(&e) {
3303 log::warn!(
3304 "Order-list submit command failed local validation for {} orders: {e}",
3305 orders.len()
3306 );
3307 } else {
3308 log::warn!(
3309 "Ambiguous order-list submit failure, awaiting reconciliation: {e}"
3310 );
3311 }
3312
3313 return Err(e);
3314 }
3315 }
3316 Ok(())
3317 });
3318
3319 Ok(())
3320 }
3321
3322 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
3323 let order = {
3324 let cache = self.core.cache();
3325 cache.order(&cmd.client_order_id).map(|o| o.clone())
3326 };
3327
3328 let Some(order) = order else {
3329 log::warn!(
3330 "Cannot modify order {}: not found in cache",
3331 cmd.client_order_id
3332 );
3333 let ts_init = self.clock.get_time_ns();
3334
3335 let rejected = OrderModifyRejected::new(
3336 self.core.trader_id,
3337 cmd.strategy_id,
3338 cmd.instrument_id,
3339 cmd.client_order_id,
3340 "Order not found in cache for modify".into(),
3341 UUID4::new(),
3342 ts_init, ts_init,
3344 false,
3345 cmd.venue_order_id,
3346 Some(self.core.account_id),
3347 );
3348
3349 self.emitter
3350 .send_order_event(OrderEventAny::ModifyRejected(rejected));
3351 return Ok(());
3352 };
3353
3354 let http_client = self.http_client.clone();
3355 let emitter = self.emitter.clone();
3356 let trader_id = self.core.trader_id;
3357 let account_id = self.core.account_id;
3358 let instrument_id = cmd.instrument_id;
3359 let venue_order_id = cmd.venue_order_id;
3360 let client_order_id = Some(cmd.client_order_id);
3361 let order_side = order.order_side();
3362 let quantity = cmd.quantity.unwrap_or_else(|| order.quantity());
3363 let price = cmd.price.or_else(|| order.price());
3364
3365 let Some(price) = price else {
3366 log::warn!(
3367 "Cannot modify order {}: price required",
3368 cmd.client_order_id
3369 );
3370 let ts_init = self.clock.get_time_ns();
3371
3372 let rejected = OrderModifyRejected::new(
3373 self.core.trader_id,
3374 cmd.strategy_id,
3375 cmd.instrument_id,
3376 cmd.client_order_id,
3377 "Price required for order modification".into(),
3378 UUID4::new(),
3379 ts_init, ts_init,
3381 false,
3382 cmd.venue_order_id,
3383 Some(self.core.account_id),
3384 );
3385
3386 self.emitter
3387 .send_order_event(OrderEventAny::ModifyRejected(rejected));
3388 return Ok(());
3389 };
3390 let command = cmd;
3391 let clock = self.clock;
3392
3393 if self.ws_trading_active() {
3394 let ws_client = self.ws_trading_client.as_ref().unwrap().clone();
3395 let dispatch_state = self.dispatch_state.clone();
3396
3397 let binance_side = BinanceSide::try_from(order_side)?;
3398 let orig_client_order_id =
3399 client_order_id.map(|id| encode_broker_id(&id, BINANCE_NAUTILUS_FUTURES_BROKER_ID));
3400
3401 let mut modify_builder = BinanceModifyOrderParamsBuilder::default();
3402 modify_builder
3403 .symbol(format_binance_symbol(&instrument_id))
3404 .side(binance_side)
3405 .quantity(quantity.to_string())
3406 .price(price.to_string());
3407
3408 if let Some(venue_id) = venue_order_id {
3409 let order_id: i64 = venue_id
3410 .inner()
3411 .parse()
3412 .context("failed to parse venue_order_id as numeric")?;
3413 modify_builder.order_id(order_id);
3414 }
3415
3416 if let Some(client_id) = orig_client_order_id {
3417 modify_builder.orig_client_order_id(client_id);
3418 }
3419
3420 let params = modify_builder
3421 .build()
3422 .context("failed to build modify params")?;
3423
3424 let request_id = ws_client.next_request_id();
3426 dispatch_state.pending_requests.insert(
3427 request_id.clone(),
3428 PendingRequest {
3429 client_order_id: command.client_order_id,
3430 venue_order_id,
3431 operation: PendingOperation::Modify,
3432 },
3433 );
3434
3435 self.spawn_task("modify_order_ws", async move {
3436 if let Err(e) = ws_client
3437 .modify_order_with_id(request_id.clone(), params)
3438 .await
3439 {
3440 dispatch_state.pending_requests.remove(&request_id);
3441 log::error!(
3442 "WS modify request failed for {}: {e}",
3443 command.client_order_id
3444 );
3445 anyhow::bail!("WS modify order failed: {e}");
3446 }
3447 Ok(())
3448 });
3449
3450 return Ok(());
3451 }
3452
3453 self.spawn_task("modify_order", async move {
3454 let result = http_client
3455 .modify_order(
3456 account_id,
3457 instrument_id,
3458 venue_order_id,
3459 client_order_id,
3460 order_side,
3461 quantity,
3462 price,
3463 )
3464 .await;
3465
3466 match result {
3467 Ok(report) => {
3468 let ts_now = clock.get_time_ns();
3469 let updated_event = OrderUpdated::new(
3470 trader_id,
3471 command.strategy_id,
3472 command.instrument_id,
3473 command.client_order_id,
3474 quantity,
3475 UUID4::new(),
3476 ts_now,
3477 ts_now,
3478 false,
3479 Some(report.venue_order_id),
3480 Some(account_id),
3481 Some(price),
3482 None,
3483 None,
3484 false, );
3486
3487 emitter.send_order_event(OrderEventAny::Updated(updated_event));
3488 }
3489 Err(e) => {
3490 if is_structured_venue_rejection(&e) {
3491 let ts_now = clock.get_time_ns();
3492
3493 let rejected = OrderModifyRejected::new(
3494 trader_id,
3495 command.strategy_id,
3496 command.instrument_id,
3497 command.client_order_id,
3498 format!("modify-order-failed: {e}").into(),
3499 UUID4::new(),
3500 ts_now,
3501 ts_now,
3502 false,
3503 command.venue_order_id,
3504 Some(account_id),
3505 );
3506
3507 emitter.send_order_event(OrderEventAny::ModifyRejected(rejected));
3508 } else {
3509 log::warn!(
3510 "Ambiguous modify failure for {}, awaiting reconciliation: {e}",
3511 command.client_order_id
3512 );
3513 }
3514
3515 anyhow::bail!("Modify order failed: {e}");
3516 }
3517 }
3518
3519 Ok(())
3520 });
3521
3522 Ok(())
3523 }
3524
3525 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
3526 self.cancel_order_internal(&cmd);
3527 Ok(())
3528 }
3529
3530 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
3531 let http_client = self.http_client.clone();
3532 let instrument_id = cmd.instrument_id;
3533
3534 self.spawn_task("cancel_all_orders", async move {
3537 match http_client.cancel_all_orders(instrument_id).await {
3538 Ok(_) => {
3539 log::debug!("Cancel all regular orders request accepted for {instrument_id}");
3540 }
3541 Err(e) => {
3542 log::error!("Failed to cancel all regular orders for {instrument_id}: {e}");
3543 }
3544 }
3545
3546 match http_client.cancel_all_algo_orders(instrument_id).await {
3547 Ok(()) => {
3548 log::debug!("Cancel all algo orders request accepted for {instrument_id}");
3549 }
3550 Err(e) => {
3551 log::error!("Failed to cancel all algo orders for {instrument_id}: {e}");
3552 }
3553 }
3554
3555 Ok(())
3556 });
3557
3558 Ok(())
3559 }
3560
3561 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
3562 const BATCH_SIZE: usize = 10;
3563
3564 if cmd.cancels.is_empty() {
3565 return Ok(());
3566 }
3567
3568 let http_client = self.http_client.clone();
3569 let command = cmd;
3570
3571 let emitter = self.emitter.clone();
3572 let trader_id = self.core.trader_id;
3573 let account_id = self.core.account_id;
3574 let clock = self.clock;
3575
3576 self.spawn_task("batch_cancel_orders", async move {
3577 let symbol = format_binance_symbol(&command.instrument_id);
3578
3579 for chunk in command.cancels.chunks(BATCH_SIZE) {
3580 let mut order_id_batch = Vec::new();
3581 let mut client_order_id_batch = Vec::new();
3582
3583 for cancel in chunk {
3584 if let Some(venue_order_id) = cancel.venue_order_id {
3585 let order_id = venue_order_id.inner().parse::<i64>().unwrap_or(0);
3586 if order_id != 0 {
3587 order_id_batch.push((
3588 BatchCancelItem::by_order_id(symbol.clone(), order_id),
3589 cancel.clone(),
3590 ));
3591 continue;
3592 }
3593 }
3594
3595 client_order_id_batch.push((
3596 BatchCancelItem::by_client_order_id(
3597 symbol.clone(),
3598 encode_broker_id(
3599 &cancel.client_order_id,
3600 BINANCE_NAUTILUS_FUTURES_BROKER_ID,
3601 ),
3602 ),
3603 cancel.clone(),
3604 ));
3605 }
3606
3607 for batch in [order_id_batch, client_order_id_batch] {
3608 if batch.is_empty() {
3609 continue;
3610 }
3611
3612 let batch_len = batch.len();
3613 let (batch_items, batch_cancels): (Vec<_>, Vec<_>) =
3614 batch.into_iter().unzip();
3615
3616 match http_client.batch_cancel_orders(&batch_items).await {
3617 Ok(results) => {
3618 for (cancel, result) in batch_cancels.iter().zip(results.iter()) {
3619 match result {
3620 BatchOrderResult::Success(response) => {
3621 let venue_order_id =
3622 VenueOrderId::new(response.order_id.to_string());
3623 let canceled_event = OrderCanceled::new(
3624 trader_id,
3625 cancel.strategy_id,
3626 cancel.instrument_id,
3627 cancel.client_order_id,
3628 UUID4::new(),
3629 cancel.ts_init,
3630 clock.get_time_ns(),
3631 false,
3632 Some(venue_order_id),
3633 Some(account_id),
3634 );
3635
3636 emitter.send_order_event(OrderEventAny::Canceled(
3637 canceled_event,
3638 ));
3639 }
3640 BatchOrderResult::Error(error) => {
3641 let rejected = OrderCancelRejected::new(
3642 trader_id,
3643 cancel.strategy_id,
3644 cancel.instrument_id,
3645 cancel.client_order_id,
3646 format!(
3647 "batch-cancel-error: code={}, msg={}",
3648 error.code, error.msg
3649 )
3650 .into(),
3651 UUID4::new(),
3652 clock.get_time_ns(),
3653 cancel.ts_init,
3654 false,
3655 cancel.venue_order_id,
3656 Some(account_id),
3657 );
3658
3659 emitter.send_order_event(OrderEventAny::CancelRejected(
3660 rejected,
3661 ));
3662 }
3663 }
3664 }
3665 }
3666 Err(e) => {
3667 if is_local_http_command_failure(&e) {
3668 log::warn!(
3669 "Batch cancel command failed local validation for {batch_len} orders: {e}",
3670 );
3671 } else {
3672 log::warn!(
3673 "Ambiguous batch cancel request failure for {batch_len} orders, awaiting reconciliation: {e}",
3674 );
3675 }
3676 return Err(e.into());
3677 }
3678 }
3679 }
3680 }
3681
3682 Ok(())
3683 });
3684
3685 Ok(())
3686 }
3687}
3688
3689#[derive(Clone, Copy, Debug, Eq, PartialEq)]
3690enum BinanceFuturesAlgoLookup {
3691 Skip,
3692 AlgoId,
3693 ClientAlgoId,
3694}
3695
3696#[derive(Clone, Copy, Debug, Eq, PartialEq)]
3697struct FuturesOrderLifetime {
3698 time_in_force: TimeInForce,
3699 good_till_date: Option<i64>,
3700}
3701
3702struct OpenOrderStatusReport {
3703 report: OrderStatusReport,
3704 quantity_free_close_position_side: Option<PositionSide>,
3705}
3706
3707#[allow(
3708 clippy::too_many_arguments,
3709 reason = "the converter receives one cohesive set of report conversion inputs"
3710)]
3711fn create_algo_order_status_report(
3712 result: &BinanceFuturesAlgoOrderQueryResult,
3713 account_id: AccountId,
3714 instrument_id: InstrumentId,
3715 price_precision: u8,
3716 size_precision: u8,
3717 treat_expired_as_canceled: bool,
3718 use_position_ids: bool,
3719 ts_init: UnixNanos,
3720) -> anyhow::Result<OrderStatusReport> {
3721 let position_side = result
3722 .actual
3723 .as_ref()
3724 .and_then(|actual| actual.position_side)
3725 .or(result.algo.position_side);
3726 let venue_position_id = make_venue_position_id(use_position_ids, instrument_id, position_side)?;
3727
3728 if let Some(actual) = result.actual.as_ref() {
3729 match result.algo.to_order_status_report_with_actual(
3730 actual,
3731 account_id,
3732 instrument_id,
3733 price_precision,
3734 size_precision,
3735 treat_expired_as_canceled,
3736 ts_init,
3737 ) {
3738 Ok(report) => return Ok(with_venue_position_id(report, venue_position_id)),
3739 Err(e) => {
3740 log::warn!(
3741 "Failed to convert matching-engine enrichment for algo order {}: {e}; falling back to Algo Service report",
3742 result.algo.algo_id
3743 );
3744 }
3745 }
3746 }
3747
3748 let report = result.algo.to_order_status_report(
3749 account_id,
3750 instrument_id,
3751 price_precision,
3752 size_precision,
3753 ts_init,
3754 )?;
3755 Ok(with_venue_position_id(report, venue_position_id))
3756}
3757
3758fn user_trades_complete_start(ts_init: UnixNanos, ts_now: UnixNanos) -> UnixNanos {
3759 let oldest_reference = ts_now.saturating_sub_ns(USER_TRADES_MAX_TS_INIT_AGE_NS);
3760 ts_init
3761 .max(oldest_reference)
3762 .saturating_sub_ns(USER_TRADES_COMPLETE_INTERVAL_NS)
3763}
3764
3765fn should_use_algo_cancel(is_algo: bool, is_triggered: bool, has_promoted_id: bool) -> bool {
3766 is_algo && !is_triggered && !has_promoted_id
3767}
3768
3769fn cancel_venue_order_id(
3770 is_algo: bool,
3771 venue_order_id: Option<VenueOrderId>,
3772 promoted_algo_order_id: Option<VenueOrderId>,
3773) -> Option<VenueOrderId> {
3774 if is_algo {
3775 promoted_algo_order_id
3776 } else {
3777 venue_order_id
3778 }
3779}
3780
3781fn is_instrument_for_product(instrument: &InstrumentAny, product_type: BinanceProductType) -> bool {
3782 match product_type {
3783 BinanceProductType::UsdM => {
3784 matches!(
3785 instrument,
3786 InstrumentAny::CryptoFuture(_)
3787 | InstrumentAny::CryptoPerpetual(_)
3788 | InstrumentAny::PerpetualContract(_)
3789 ) && !instrument.is_inverse()
3790 }
3791 BinanceProductType::CoinM => {
3792 matches!(
3793 instrument,
3794 InstrumentAny::CryptoFuture(_) | InstrumentAny::CryptoPerpetual(_)
3795 ) && instrument.is_inverse()
3796 }
3797 _ => false,
3798 }
3799}
3800
3801#[cfg(test)]
3802mod tests {
3803 use nautilus_model::{
3804 instruments::stubs::{
3805 crypto_future_btcusdt, crypto_perpetual_ethusdt, currency_pair_btcusdt,
3806 perpetual_contract_eurusd, xbtusd_bitmex,
3807 },
3808 types::Price,
3809 };
3810 use rstest::rstest;
3811
3812 use super::*;
3813 use crate::common::testing::load_fixture_string;
3814
3815 fn http_error(code: i64) -> anyhow::Error {
3816 anyhow::Error::new(BinanceFuturesHttpError::BinanceError {
3817 code,
3818 message: format!("test error {code}"),
3819 })
3820 }
3821
3822 #[rstest]
3823 #[case(None, Some("BTCUSDT-PERP.BINANCE-LONG"))]
3824 #[case(Some("BTCUSDT-PERP.BINANCE-LONG"), Some("BTCUSDT-PERP.BINANCE-LONG"))]
3825 #[case(Some("P-VIRTUAL-LONG"), None)]
3826 fn test_validate_submit_position_id_accepts_compatible_identity(
3827 #[case] submitted_position_id: Option<&str>,
3828 #[case] venue_position_id: Option<&str>,
3829 ) {
3830 validate_submit_position_id(
3831 submitted_position_id.map(PositionId::from),
3832 venue_position_id.map(PositionId::from),
3833 )
3834 .unwrap();
3835 }
3836
3837 #[rstest]
3838 fn test_validate_submit_position_id_rejects_custom_venue_identity() {
3839 let error = validate_submit_position_id(
3840 Some(PositionId::from("P-VIRTUAL-LONG")),
3841 Some(PositionId::from("BTCUSDT-PERP.BINANCE-LONG")),
3842 )
3843 .unwrap_err();
3844
3845 assert_eq!(
3846 error.to_string(),
3847 "submitted position ID P-VIRTUAL-LONG conflicts with canonical Binance Futures venue position ID BTCUSDT-PERP.BINANCE-LONG; omit position_id while use_position_ids=true, or set use_position_ids=false for virtual hedging",
3848 );
3849 }
3850
3851 #[rstest]
3852 fn test_create_account_state_preserves_info_decimal_values() {
3853 let json = load_fixture_string("futures/http_json/account_info_v2.json");
3854 let mut account_info: BinanceFuturesAccountInfo = serde_json::from_str(&json).unwrap();
3855 account_info.total_wallet_balance = Some("1.0000000000000001".parse().unwrap());
3856 account_info.total_margin_balance = Some("2.0000000000000002".parse().unwrap());
3857 account_info.total_initial_margin = Some("3.0000000000000003".parse().unwrap());
3858 account_info.total_maint_margin = Some("4.0000000000000004".parse().unwrap());
3859 account_info.total_unrealized_profit = Some("5.0000000000000005".parse().unwrap());
3860 account_info.total_cross_wallet_balance = Some("6.0000000000000006".parse().unwrap());
3861 account_info.total_cross_un_pnl = Some("7.0000000000000007".parse().unwrap());
3862 account_info.available_balance = Some("8.0000000000000008".parse().unwrap());
3863 account_info.max_withdraw_amount = Some("9.0000000000000009".parse().unwrap());
3864
3865 let state = BinanceFuturesExecutionClient::create_account_state_from(
3866 &account_info,
3867 AccountId::from("BINANCE-001"),
3868 AccountType::Margin,
3869 Currency::USDT(),
3870 get_atomic_clock_realtime(),
3871 );
3872
3873 let info = state.info.as_ref().unwrap();
3874 assert_eq!(info.len(), 9);
3875 assert_eq!(
3876 info.get_str("total_wallet_balance"),
3877 Some("1.0000000000000001")
3878 );
3879 assert_eq!(
3880 info.get_str("total_margin_balance"),
3881 Some("2.0000000000000002")
3882 );
3883 assert_eq!(
3884 info.get_str("total_initial_margin"),
3885 Some("3.0000000000000003")
3886 );
3887 assert_eq!(
3888 info.get_str("total_maint_margin"),
3889 Some("4.0000000000000004")
3890 );
3891 assert_eq!(
3892 info.get_str("total_unrealized_profit"),
3893 Some("5.0000000000000005")
3894 );
3895 assert_eq!(
3896 info.get_str("total_cross_wallet_balance"),
3897 Some("6.0000000000000006")
3898 );
3899 assert_eq!(
3900 info.get_str("total_cross_unpnl"),
3901 Some("7.0000000000000007")
3902 );
3903 assert_eq!(
3904 info.get_str("available_balance"),
3905 Some("8.0000000000000008")
3906 );
3907 assert_eq!(
3908 info.get_str("max_withdraw_amount"),
3909 Some("9.0000000000000009")
3910 );
3911 }
3912
3913 #[rstest]
3914 #[case::regular(false, false, false, false)]
3915 #[case::untriggered_algo(true, false, false, true)]
3916 #[case::triggered_algo(true, true, false, false)]
3917 #[case::promoted_algo(true, false, true, false)]
3918 fn test_should_use_algo_cancel(
3919 #[case] is_algo: bool,
3920 #[case] is_triggered: bool,
3921 #[case] has_promoted_id: bool,
3922 #[case] expected: bool,
3923 ) {
3924 assert_eq!(
3925 should_use_algo_cancel(is_algo, is_triggered, has_promoted_id),
3926 expected
3927 );
3928 }
3929
3930 #[rstest]
3931 #[case::not_close(BinanceSide::Sell, Some(BinancePositionSide::Long), false, None, None)]
3932 #[case::positive_quantity(
3933 BinanceSide::Sell,
3934 Some(BinancePositionSide::Long),
3935 true,
3936 Some("0.001"),
3937 None
3938 )]
3939 #[case::hedge_long(
3940 BinanceSide::Sell,
3941 Some(BinancePositionSide::Long),
3942 true,
3943 None,
3944 Some(PositionSide::Long)
3945 )]
3946 #[case::hedge_long_zero(
3947 BinanceSide::Sell,
3948 Some(BinancePositionSide::Long),
3949 true,
3950 Some("0.0000"),
3951 Some(PositionSide::Long)
3952 )]
3953 #[case::hedge_short(
3954 BinanceSide::Buy,
3955 Some(BinancePositionSide::Short),
3956 true,
3957 None,
3958 Some(PositionSide::Short)
3959 )]
3960 #[case::invalid_hedge_long(BinanceSide::Buy, Some(BinancePositionSide::Long), true, None, None)]
3961 #[case::invalid_hedge_short(
3962 BinanceSide::Sell,
3963 Some(BinancePositionSide::Short),
3964 true,
3965 None,
3966 None
3967 )]
3968 #[case::one_way_long(
3969 BinanceSide::Sell,
3970 Some(BinancePositionSide::Both),
3971 true,
3972 None,
3973 Some(PositionSide::Long)
3974 )]
3975 #[case::one_way_short(
3976 BinanceSide::Buy,
3977 Some(BinancePositionSide::Both),
3978 true,
3979 None,
3980 Some(PositionSide::Short)
3981 )]
3982 #[case::missing_side(BinanceSide::Buy, None, true, None, Some(PositionSide::Short))]
3983 #[case::unknown_side(
3984 BinanceSide::Sell,
3985 Some(BinancePositionSide::Unknown),
3986 true,
3987 None,
3988 None
3989 )]
3990 fn test_quantity_free_close_position_side(
3991 #[case] side: BinanceSide,
3992 #[case] position_side: Option<BinancePositionSide>,
3993 #[case] close_position: bool,
3994 #[case] quantity: Option<&str>,
3995 #[case] expected: Option<PositionSide>,
3996 ) {
3997 let json = load_fixture_string("futures/http_json/open_algo_orders.json");
3998 let mut orders: Vec<BinanceFuturesAlgoOrder> = serde_json::from_str(&json).unwrap();
3999 let mut order = orders.remove(0);
4000 order.side = side;
4001 order.position_side = position_side;
4002 order.close_position = Some(close_position);
4003 order.quantity = quantity.map(str::to_string);
4004
4005 assert_eq!(quantity_free_close_position_side(&order), expected);
4006 }
4007
4008 #[rstest]
4009 #[case::regular(
4010 false,
4011 Some(VenueOrderId::from("8886774")),
4012 None,
4013 Some(VenueOrderId::from("8886774"))
4014 )]
4015 #[case::unpromoted_algo(true, Some(VenueOrderId::from("2148719")), None, None)]
4016 #[case::promoted(
4017 true,
4018 Some(VenueOrderId::from("2148719")),
4019 Some(VenueOrderId::from("22542179")),
4020 Some(VenueOrderId::from("22542179"))
4021 )]
4022 fn test_cancel_venue_order_id(
4023 #[case] is_algo: bool,
4024 #[case] venue_order_id: Option<VenueOrderId>,
4025 #[case] promoted_algo_order_id: Option<VenueOrderId>,
4026 #[case] expected: Option<VenueOrderId>,
4027 ) {
4028 assert_eq!(
4029 cancel_venue_order_id(is_algo, venue_order_id, promoted_algo_order_id),
4030 expected
4031 );
4032 }
4033
4034 #[rstest]
4035 fn test_classify_submit_order_error_gtx_is_post_only() {
4036 let err = http_error(BINANCE_GTX_ORDER_REJECT_CODE);
4037 assert!(classify_submit_order_error(&err));
4038 }
4039
4040 #[rstest]
4041 fn test_futures_order_lifetime_encodes_valid_usdm_gtd() {
4042 let ts_now = UnixNanos::from_seconds(1_700_000_000);
4043 let expire_time = UnixNanos::from_seconds(1_700_000_601);
4044
4045 let lifetime = determine_futures_order_lifetime(
4046 BinanceProductType::UsdM,
4047 OrderType::Limit,
4048 TimeInForce::Gtd,
4049 Some(expire_time),
4050 false,
4051 true,
4052 ts_now,
4053 )
4054 .unwrap();
4055
4056 assert_eq!(lifetime.time_in_force, TimeInForce::Gtd);
4057 assert_eq!(lifetime.good_till_date, Some(1_700_000_601_000));
4058 }
4059
4060 #[rstest]
4061 #[case::minimum(
4062 BinanceProductType::UsdM,
4063 OrderType::Limit,
4064 UnixNanos::from_seconds(1_700_000_600),
4065 false,
4066 "strictly greater"
4067 )]
4068 #[case::subsecond(
4069 BinanceProductType::UsdM,
4070 OrderType::Limit,
4071 UnixNanos::from(1_700_000_601_000_000_001),
4072 false,
4073 "whole-second precision"
4074 )]
4075 #[case::coin_m(
4076 BinanceProductType::CoinM,
4077 OrderType::Limit,
4078 UnixNanos::from_seconds(1_700_000_601),
4079 false,
4080 "does not support native GTD"
4081 )]
4082 #[case::market(
4083 BinanceProductType::UsdM,
4084 OrderType::Market,
4085 UnixNanos::from_seconds(1_700_000_601),
4086 false,
4087 "does not support GTD for order type Market"
4088 )]
4089 #[case::post_only(
4090 BinanceProductType::UsdM,
4091 OrderType::Limit,
4092 UnixNanos::from_seconds(1_700_000_601),
4093 true,
4094 "cannot be post-only"
4095 )]
4096 fn test_futures_order_lifetime_rejects_invalid_gtd(
4097 #[case] product_type: BinanceProductType,
4098 #[case] order_type: OrderType,
4099 #[case] expire_time: UnixNanos,
4100 #[case] post_only: bool,
4101 #[case] expected: &str,
4102 ) {
4103 let error = determine_futures_order_lifetime(
4104 product_type,
4105 order_type,
4106 TimeInForce::Gtd,
4107 Some(expire_time),
4108 post_only,
4109 true,
4110 UnixNanos::from_seconds(1_700_000_000),
4111 )
4112 .unwrap_err();
4113
4114 assert!(error.to_string().contains(expected));
4115 }
4116
4117 #[rstest]
4118 #[case::coin_m_limit(BinanceProductType::CoinM, OrderType::Limit, false)]
4119 #[case::usd_m_market(BinanceProductType::UsdM, OrderType::Market, false)]
4120 #[case::usd_m_post_only(BinanceProductType::UsdM, OrderType::Limit, true)]
4121 fn test_futures_order_lifetime_maps_locally_managed_gtd_to_gtc(
4122 #[case] product_type: BinanceProductType,
4123 #[case] order_type: OrderType,
4124 #[case] post_only: bool,
4125 ) {
4126 let lifetime = determine_futures_order_lifetime(
4127 product_type,
4128 order_type,
4129 TimeInForce::Gtd,
4130 Some(UnixNanos::from_seconds(1_700_000_601)),
4131 post_only,
4132 false,
4133 UnixNanos::from_seconds(1_700_000_000),
4134 )
4135 .unwrap();
4136
4137 assert_eq!(lifetime.time_in_force, TimeInForce::Gtc);
4138 assert_eq!(lifetime.good_till_date, None);
4139 }
4140
4141 #[rstest]
4142 fn test_futures_order_lifetime_requires_expire_time() {
4143 let error = determine_futures_order_lifetime(
4144 BinanceProductType::UsdM,
4145 OrderType::Limit,
4146 TimeInForce::Gtd,
4147 None,
4148 false,
4149 true,
4150 UnixNanos::from_seconds(1_700_000_000),
4151 )
4152 .unwrap_err();
4153
4154 assert!(error.to_string().contains("requires an expire_time"));
4155 }
4156
4157 #[rstest]
4158 fn test_futures_order_lifetime_maximum_exceeds_unix_nanos_range() {
4159 let maximum_unix_nanos_millis = u64::MAX / NANOSECONDS_IN_MILLISECOND;
4160
4161 assert!(maximum_unix_nanos_millis < BINANCE_GTD_MAX_MILLIS);
4162 }
4163
4164 #[rstest]
4165 fn test_classify_submit_order_error_dual_side_sync_is_not_post_only() {
4166 let err = http_error(BINANCE_FUTURES_DUAL_SIDE_SYNC_REJECT_CODE);
4169 assert!(!classify_submit_order_error(&err));
4170 }
4171
4172 #[rstest]
4173 fn test_classify_submit_order_error_other_venue_code_is_not_post_only() {
4174 let err = http_error(-2010);
4175 assert!(!classify_submit_order_error(&err));
4176 }
4177
4178 #[rstest]
4179 fn test_classify_submit_order_error_non_binance_error_is_not_post_only() {
4180 let err = anyhow::anyhow!("network failure");
4181 assert!(!classify_submit_order_error(&err));
4182 }
4183
4184 #[rstest]
4185 #[case(BINANCE_UNEXPECTED_RESPONSE_CODE)]
4186 #[case(BINANCE_STATUS_UNKNOWN_CODE)]
4187 fn test_unknown_status_submit_error_is_ambiguous(#[case] code: i64) {
4188 let err = http_error(code);
4189 assert!(is_ambiguous_submit_error(&err));
4190 assert!(is_structured_venue_rejection(&err));
4191 }
4192
4193 #[rstest]
4194 fn test_other_structured_submit_error_is_not_ambiguous() {
4195 let err = http_error(BINANCE_GTX_ORDER_REJECT_CODE);
4196 assert!(!is_ambiguous_submit_error(&err));
4197 assert!(is_structured_venue_rejection(&err));
4198 }
4199
4200 #[rstest]
4201 fn test_instrument_product_matching_distinguishes_futures_products_from_spot() {
4202 let usdm = InstrumentAny::CryptoPerpetual(crypto_perpetual_ethusdt());
4203 let generic_perpetual = InstrumentAny::PerpetualContract(perpetual_contract_eurusd());
4204 let coinm = InstrumentAny::CryptoPerpetual(xbtusd_bitmex());
4205 let delivery =
4206 || crypto_future_btcusdt(2, 6, Price::from("0.01"), Quantity::from("0.000001"));
4207 let usdm_delivery = InstrumentAny::CryptoFuture(delivery());
4208 let mut coinm_delivery = delivery();
4209 coinm_delivery.is_inverse = true;
4210 let coinm_delivery = InstrumentAny::CryptoFuture(coinm_delivery);
4211 let spot = InstrumentAny::CurrencyPair(currency_pair_btcusdt());
4212
4213 assert!(is_instrument_for_product(&usdm, BinanceProductType::UsdM));
4214 assert!(!is_instrument_for_product(&usdm, BinanceProductType::CoinM));
4215 assert!(is_instrument_for_product(
4216 &generic_perpetual,
4217 BinanceProductType::UsdM
4218 ));
4219 assert!(!is_instrument_for_product(
4220 &generic_perpetual,
4221 BinanceProductType::CoinM
4222 ));
4223 assert!(is_instrument_for_product(&coinm, BinanceProductType::CoinM));
4224 assert!(!is_instrument_for_product(&coinm, BinanceProductType::UsdM));
4225 assert!(is_instrument_for_product(
4226 &usdm_delivery,
4227 BinanceProductType::UsdM
4228 ));
4229 assert!(!is_instrument_for_product(
4230 &usdm_delivery,
4231 BinanceProductType::CoinM
4232 ));
4233 assert!(is_instrument_for_product(
4234 &coinm_delivery,
4235 BinanceProductType::CoinM
4236 ));
4237 assert!(!is_instrument_for_product(
4238 &coinm_delivery,
4239 BinanceProductType::UsdM
4240 ));
4241 assert!(!is_instrument_for_product(&spot, BinanceProductType::UsdM));
4242 assert!(!is_instrument_for_product(&spot, BinanceProductType::CoinM));
4243 }
4244}