Skip to main content

nautilus_binance/futures/
execution.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live execution client implementation for the Binance Futures adapter.
17
18use 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
138/// Listen key keepalive interval (30 minutes).
139const LISTEN_KEY_KEEPALIVE_SECS: u64 = 30 * 60;
140
141/// Consecutive keepalive failures before a listenKey rotation is triggered.
142const 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
148// Binance does not define whether its three-month retention is calendar-based or fixed-duration
149// An 88-day interval leaves at least one day inside either interpretation
150// https://developers.binance.com/en/docs/products/derivatives-trading-usds-futures/change-log
151const USER_TRADES_COMPLETE_INTERVAL_NS: u64 = 88 * NANOSECONDS_IN_DAY;
152
153// Keep half the one-day retention safety margin when a command waits before execution
154const 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
160/// Query parameter declaring that the command's venue order ID is a Binance Algo Service
161/// `algoId`, rather than a regular matching-engine `orderId` or triggered `actualOrderId`.
162pub const BINANCE_VENUE_ORDER_ID_IS_ALGO_ID_PARAM: &str = "venue_order_id_is_algo_id";
163
164/// Live execution client for Binance Futures trading.
165///
166/// Implements the [`ExecutionClient`] trait for order management on Binance
167/// USD-M and COIN-M Futures markets. Uses HTTP API for order operations and
168/// WebSocket for real-time order updates via user data stream.
169///
170/// Uses a two-tier architecture with an execution handler that maintains
171/// pending order maps for correlating WebSocket updates with order context.
172#[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    /// Creates a new [`BinanceFuturesExecutionClient`].
200    ///
201    /// # Errors
202    ///
203    /// Returns an error if the HTTP client fails to initialize, credentials are
204    /// missing, or the product type is not a futures type (UsdM or CoinM).
205    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, // timeout_secs
239            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    /// Returns whether the account is in hedge mode (dual side position).
311    #[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    /// Returns a clone of the HTTP client's instruments cache Arc.
338    #[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    /// Converts Binance futures account info to Nautilus account state.
345    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        // Emit account-wide (cross-margin) margin balances per collateral asset.
379        // Binance reports per-asset `initialMargin` / `maintMargin` which covers both
380        // USDT-M (typically USDT, or USDT+BNB under multi-assets mode) and COIN-M
381        // (one entry per base coin, e.g. BTC / ETH).
382        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, // reported
432            UUID4::new(),
433            ts_now,
434            ts_now,
435            None, // base currency
436        )
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    /// Returns whether the WS trading client is connected and active.
490    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        // Register identity for tracked/external dispatch routing
532        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        // Non-algo orders can route through WS trading API when active
567        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            // Pre-register before sending to avoid response racing the insert
620            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                    // Keep order registered - if HTTP failed due to timeout but order
700                    // reached Binance, WebSocket updates will still arrive. The order
701                    // will be cleaned up via WebSocket rejection or reconciliation.
702                    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        // Non-triggered algo orders use algo cancel endpoint, triggered use regular
745        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        // Non-algo cancels can route through WS trading API when active
770        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            // Pre-register before sending to avoid response racing the insert
798            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                // Try algo cancel first; if it fails, the order may have been triggered
828                // before this session started, so fall back to regular cancel
829                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    /// Returns the (price_precision, size_precision) for an instrument.
967    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    /// Creates a position status report from Binance position risk data.
989    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(&params),
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                // Best-effort: a venue reject is non-fatal, but transport errors
1209                // propagate so an unhealthy connection still surfaces.
1210                match self.http_client.set_leverage(&params).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(&params).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    // Binance accepts close-all orders without quantity and reports the omitted value as either
1283    // absent or zero. Core reconciliation needs a positive quantity to materialize the order, so
1284    // the current position leg is the narrowest valid source.
1285    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    // `close_position` retires an entire hedge leg, so it carries close intent on its own
1324    // and cannot be combined with `reduce_only`, which otherwise selects the closing side.
1325    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
1598/// Classifies a submit-order error for the rejection event.
1599///
1600/// Returns `true` when the venue indicated a post-only (GTX) rejection. Logs
1601/// a hint when the error matches the UM/CM `dualSidePosition` sync rejection
1602/// (`-4531`), which is an account/setup mismatch rather than a routing fault.
1603pub(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        // Check hedge mode
1721        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        // Apply configured leverage and margin types
1749        self.apply_futures_config()
1750            .await
1751            .context("failed to apply futures config")?;
1752
1753        // Create listen key for user data stream
1754        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        // Start listen key keepalive task
1839        {
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        // Start listen key recovery driver task
1892        {
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        // Request initial account state
1918        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        // Connect WS trading client (primary order transport for USD-M)
1937        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        // Drop the recovery tx so the driver exits its recv loop
2087        self.recovery_tx.take();
2088
2089        // Cancel all background tasks
2090        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        // Close WebSocket
2126        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        // Close listen key
2138        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(&params).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                // A conditional order may expose its Algo Service `algoId` before triggering and
2276                // its matching-engine `actualOrderId` afterwards. Only an explicit ID-kind hint
2277                // makes a venue ID safe for Algo Service lookup; cached conditional orders retain
2278                // their existing clientAlgoId fallback.
2279                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(&params).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(&params).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(&params).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(&params).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(&params).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                    // Untriggered algo orders (STOP_MARKET, STOP_LIMIT, MIT, LIT,
2957                    // TRAILING_STOP_MARKET) live in the Binance Futures Algo Service
2958                    // and return -2013 on the regular order endpoint. Mirror the
2959                    // reconciliation path (`generate_order_status_report`) so live
2960                    // inflight checks resolve them instead of exhausting retries into
2961                    // `OrderRejected(reason='INFLIGHT_TIMEOUT')`.
2962                    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        // Validate before submission (Initialized -> Denied is valid,
3072        // but Submitted -> Denied is not, so validate before emitting OrderSubmitted)
3073        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, // no venue timestamp, rejected locally
3343                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, // no venue timestamp, rejected locally
3380                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            // Pre-register before sending to avoid response racing the insert
3425            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, // is_quote_quantity
3485                    );
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        // USD-M Futures WS Trading API does not expose an openOrders.cancelAll
3535        // method, so regular and algo cancel-all both go through HTTP.
3536        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        // -4531 is a hedge-mode/account-setup issue, not a post-only rejection.
4167        // Make sure the new classifier branch does not mark it as post-only.
4168        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}