Skip to main content

nautilus_bybit/
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 Bybit adapter.
17
18use std::{
19    future::Future,
20    sync::Arc,
21    time::{Duration, Instant},
22};
23
24use ahash::AHashMap;
25use anyhow::Context;
26use async_trait::async_trait;
27use futures_util::{StreamExt, pin_mut};
28use nautilus_common::{
29    clients::ExecutionClient,
30    live::runner::get_exec_event_sender,
31    messages::execution::{
32        BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
33        GenerateFillReportsBuilder, GenerateOrderStatusReport, GenerateOrderStatusReports,
34        GenerateOrderStatusReportsBuilder, GeneratePositionStatusReports, ModifyOrder,
35        QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
36    },
37};
38use nautilus_core::{
39    DurationNanos, Params, UUID4, UnixNanos,
40    env::get_or_env_var,
41    string::secret::SecretString,
42    time::{AtomicTime, get_atomic_clock_realtime},
43};
44use nautilus_live::{
45    ExecutionClientCore, ExecutionEventEmitter, SocketControl,
46    execution::failure::CommandFailure,
47    task::{TaskGroup, TaskGroupGuard},
48};
49use nautilus_model::{
50    accounts::AccountAny,
51    enums::{OmsType, OrderType, TimeInForce},
52    events::OrderDeniedReason,
53    identifiers::{AccountId, ClientId, ClientOrderId, InstrumentId, Venue},
54    instruments::{Instrument, InstrumentAny},
55    orders::{Order, OrderAny},
56    reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
57    types::{AccountBalance, MarginBalance, Price},
58};
59use ustr::Ustr;
60
61use crate::{
62    common::{
63        consts::BYBIT_VENUE,
64        credential::credential_env_vars,
65        enums::{
66            BybitAccountType, BybitEnvironment, BybitOrderSide, BybitOrderSmpType, BybitOrderType,
67            BybitPositionIdx, BybitPositionMode, BybitProductType, BybitTimeInForce, BybitTpSlMode,
68            resolve_trigger_type,
69        },
70        parse::{
71            BybitTpSlParams, bybit_rejection_due_post_only, extract_raw_symbol, get_price_str,
72            make_hedge_venue_position_id, nanos_to_millis, parse_bybit_tp_sl_params,
73            resolve_position_idx as resolve_bybit_position_idx, spot_leverage, spot_market_unit,
74            trigger_direction,
75        },
76        rate_limit::batch_send_limit,
77        symbol::BybitSymbol,
78    },
79    config::BybitExecutionClientConfig,
80    http::{
81        client::BybitHttpClient,
82        error::{
83            BybitCancelOrderError, BybitHttpError, BybitModifyOrderError, BybitSubmitOrderError,
84            is_bybit_ambiguous_order_error_code,
85        },
86    },
87    websocket::{
88        client::BybitWebSocketClient,
89        dispatch::{
90            OrderIdentity, OrderStateSnapshot, PendingOperation, WsDispatchState,
91            dispatch_ws_message,
92        },
93        messages::{BybitWsAmendOrderParams, BybitWsCancelOrderParams, BybitWsPlaceOrderParams},
94    },
95};
96
97/// Live execution client for Bybit.
98#[derive(Debug)]
99pub struct BybitExecutionClient {
100    core: ExecutionClientCore,
101    clock: &'static AtomicTime,
102    config: BybitExecutionClientConfig,
103    emitter: ExecutionEventEmitter,
104    http_client: BybitHttpClient,
105    ws_private: BybitWebSocketClient,
106    ws_trade: BybitWebSocketClient,
107    session_tasks: TaskGroup,
108    pending_tasks: TaskGroup,
109    shutdown_errors: Vec<String>,
110    instruments_cache: Arc<AHashMap<Ustr, InstrumentAny>>,
111    dispatch_state: Arc<WsDispatchState>,
112}
113
114impl BybitExecutionClient {
115    /// Creates a new [`BybitExecutionClient`].
116    ///
117    /// # Errors
118    ///
119    /// Returns an error if the client fails to initialize.
120    pub fn new(
121        core: ExecutionClientCore,
122        config: BybitExecutionClientConfig,
123    ) -> anyhow::Result<Self> {
124        let (key_var, secret_var) = credential_env_vars(config.environment);
125        let api_key = get_or_env_var(
126            config.api_key.clone().map(SecretString::into_inner),
127            key_var,
128        )?;
129        let api_secret = get_or_env_var(
130            config.api_secret.clone().map(SecretString::into_inner),
131            secret_var,
132        )?;
133        let proxy_url = config
134            .proxy_url
135            .as_ref()
136            .map(|value| value.expose_secret().to_owned());
137
138        let http_client = BybitHttpClient::with_credentials(
139            api_key.clone(),
140            api_secret.clone(),
141            Some(config.http_base_url()),
142            config.http_timeout_secs,
143            config.max_retries,
144            config.retry_delay_initial_ms,
145            config.retry_delay_max_ms,
146            config.recv_window_ms,
147            proxy_url.clone(),
148        )?;
149        http_client.set_use_spot_position_reports(config.use_spot_position_reports);
150
151        let mut ws_private = BybitWebSocketClient::new_private(
152            config.environment,
153            Some(api_key.clone()),
154            Some(api_secret.clone()),
155            Some(config.ws_private_url()),
156            config.heartbeat_interval_secs,
157            config.transport_backend,
158            proxy_url.clone(),
159        )
160        .with_socket_control(SocketControl::new(
161            core.client_id,
162            Some(*BYBIT_VENUE),
163            "bybit-user-streams",
164        ));
165
166        if let Some(secs) = config.auth_timeout_secs {
167            ws_private.set_auth_wait_timeout(Duration::from_secs(secs));
168        }
169
170        let mut ws_trade = BybitWebSocketClient::new_trade(
171            config.environment,
172            Some(api_key),
173            Some(api_secret),
174            Some(config.ws_trade_url()),
175            config.heartbeat_interval_secs,
176            config.transport_backend,
177            proxy_url,
178        )
179        .with_socket_control(SocketControl::new(
180            core.client_id,
181            Some(*BYBIT_VENUE),
182            "bybit-trading",
183        ));
184
185        if let Some(secs) = config.auth_timeout_secs {
186            ws_trade.set_auth_wait_timeout(Duration::from_secs(secs));
187        }
188        ws_trade.set_recv_window_ms(config.recv_window_ms);
189
190        let clock = get_atomic_clock_realtime();
191        let emitter = ExecutionEventEmitter::new(
192            clock,
193            core.trader_id,
194            core.account_id,
195            core.account_type,
196            None,
197        );
198
199        let session_tasks = TaskGroup::new();
200        let pending_tasks = TaskGroup::new();
201
202        Ok(Self {
203            core,
204            clock,
205            config,
206            emitter,
207            http_client,
208            ws_private,
209            ws_trade,
210            session_tasks,
211            pending_tasks,
212            shutdown_errors: Vec::new(),
213            instruments_cache: Arc::new(AHashMap::new()),
214            dispatch_state: Arc::new(WsDispatchState::default()),
215        })
216    }
217
218    fn product_types(&self) -> Vec<BybitProductType> {
219        if self.config.product_types.is_empty() {
220            vec![BybitProductType::Linear]
221        } else {
222            self.config.product_types.clone()
223        }
224    }
225
226    fn update_account_state(&self) {
227        let http_client = self.http_client.clone();
228        let account_id = self.core.account_id;
229        let emitter = self.emitter.clone();
230
231        self.spawn_task("query_account", async move {
232            let account_state = http_client
233                .request_account_state(BybitAccountType::Unified, account_id)
234                .await
235                .context("failed to request Bybit account state")?;
236            emitter.send_account_state(account_state);
237            Ok(())
238        });
239    }
240
241    fn spawn_task<F>(&self, description: &'static str, fut: F)
242    where
243        F: Future<Output = anyhow::Result<()>> + Send + 'static,
244    {
245        let future = async move {
246            if let Err(e) = fut.await {
247                log::warn!("{description} failed: {e:?}");
248            }
249        };
250
251        if let Err(e) = self.pending_tasks.spawn(future) {
252            log::warn!("Skipping Bybit {description} after shutdown began: {e}");
253        }
254    }
255
256    fn abort_pending_tasks(&self) {
257        self.pending_tasks.begin_shutdown();
258    }
259
260    fn abort_session_tasks(&self) {
261        self.session_tasks.begin_shutdown();
262    }
263
264    async fn await_pending_tasks(&self) -> anyhow::Result<()> {
265        self.pending_tasks.begin_shutdown();
266        self.pending_tasks
267            .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
268            .await
269            .map_err(|e| anyhow::anyhow!("Failed to terminate Bybit execution tasks: {e}"))?;
270        Ok(())
271    }
272
273    async fn await_session_tasks(&self) -> anyhow::Result<()> {
274        self.session_tasks.begin_shutdown();
275        self.session_tasks
276            .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
277            .await
278            .map_err(|e| {
279                anyhow::anyhow!("Failed to terminate Bybit execution session tasks: {e}")
280            })?;
281        Ok(())
282    }
283
284    async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
285        self.abort_session_tasks();
286        self.abort_pending_tasks();
287        self.http_client.cancel_all_requests();
288        self.ws_private.begin_shutdown();
289        self.ws_trade.begin_shutdown();
290
291        if let Err(e) = self.ws_private.close().await {
292            self.shutdown_errors
293                .push(format!("private WebSocket shutdown failed: {e}"));
294        }
295
296        if let Err(e) = self.ws_trade.close().await {
297            self.shutdown_errors
298                .push(format!("trade WebSocket shutdown failed: {e}"));
299        }
300        self.dispatch_state.clear_repay_sender();
301
302        if let Err(e) = self.await_session_tasks().await {
303            self.shutdown_errors.push(e.to_string());
304        }
305
306        if let Err(e) = self.await_pending_tasks().await {
307            self.shutdown_errors.push(e.to_string());
308        }
309        self.core.set_disconnected();
310
311        if self.shutdown_errors.is_empty() {
312            Ok(())
313        } else {
314            let errors = std::mem::take(&mut self.shutdown_errors);
315            anyhow::bail!("Bybit execution shutdown failed: {}", errors.join("; "))
316        }
317    }
318
319    /// Polls the cache until the account is registered or timeout is reached.
320    async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
321        let account_id = self.core.account_id;
322
323        if self.core.cache().account(&account_id).is_some() {
324            log::info!("Account {account_id} registered");
325            return Ok(());
326        }
327
328        let start = Instant::now();
329        let timeout = Duration::from_secs_f64(timeout_secs);
330        let interval = Duration::from_millis(10);
331
332        loop {
333            tokio::time::sleep(interval).await;
334
335            if self.core.cache().account(&account_id).is_some() {
336                log::info!("Account {account_id} registered");
337                return Ok(());
338            }
339
340            if start.elapsed() >= timeout {
341                anyhow::bail!(
342                    "Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
343                );
344            }
345        }
346    }
347
348    fn get_product_type_for_instrument(&self, instrument_id: InstrumentId) -> BybitProductType {
349        BybitProductType::from_suffix(instrument_id.symbol.as_str()).unwrap_or_else(|| {
350            log::warn!("No product-type suffix on {instrument_id}, defaulting to Linear");
351            BybitProductType::Linear
352        })
353    }
354
355    const fn provides_bulk_position_coverage_for_product_type(
356        product_type: BybitProductType,
357    ) -> bool {
358        !matches!(product_type, BybitProductType::Spot)
359    }
360
361    async fn generate_bulk_position_status_reports(
362        &self,
363        product_types: Vec<BybitProductType>,
364    ) -> anyhow::Result<Vec<PositionStatusReport>> {
365        let mut reports = Vec::new();
366
367        for product_type in product_types {
368            if !Self::provides_bulk_position_coverage_for_product_type(product_type) {
369                continue;
370            }
371
372            let mut fetched = self
373                .http_client
374                .request_position_status_reports(self.core.account_id, product_type, None)
375                .await?;
376            reports.append(&mut fetched);
377        }
378
379        Ok(reports)
380    }
381
382    fn resolve_position_idx(
383        &self,
384        instrument_id: InstrumentId,
385        order_side: BybitOrderSide,
386        is_reduce_only: bool,
387        manual_override: Option<BybitPositionIdx>,
388    ) -> Option<BybitPositionIdx> {
389        let product_type = self.get_product_type_for_instrument(instrument_id);
390        if !matches!(
391            product_type,
392            BybitProductType::Linear | BybitProductType::Inverse
393        ) {
394            return None;
395        }
396        let mode = self
397            .config
398            .position_mode
399            .as_ref()
400            .and_then(|map| map.get(instrument_id.symbol.as_str()).copied());
401        resolve_bybit_position_idx(mode, order_side, is_reduce_only, manual_override)
402    }
403
404    fn resolve_smp_type(
405        &self,
406        manual_override: Option<BybitOrderSmpType>,
407    ) -> Option<BybitOrderSmpType> {
408        manual_override.or(self.config.smp_type)
409    }
410
411    async fn apply_account_configuration(&self) -> anyhow::Result<()> {
412        self.apply_leverages_setting().await;
413        self.apply_position_modes_setting().await;
414        self.apply_margin_mode_setting().await
415    }
416
417    async fn apply_leverages_setting(&self) {
418        let Some(leverages) = &self.config.futures_leverages else {
419            return;
420        };
421
422        for (symbol_str, leverage) in leverages {
423            self.apply_leverage_entry(symbol_str, *leverage).await;
424        }
425    }
426
427    async fn apply_leverage_entry(&self, symbol_str: &str, leverage: u32) {
428        let Some(symbol) = Self::parse_derivative_symbol(symbol_str) else {
429            return;
430        };
431        let lev = leverage.to_string();
432        let result = self
433            .http_client
434            .set_leverage(symbol.product_type(), symbol.raw_symbol(), &lev, &lev)
435            .await;
436
437        match result {
438            Ok(_) => log::info!("Set leverage for {symbol_str} to {leverage}"),
439            Err(e) if Self::is_unchanged_error(&e, "110043") => {
440                log::debug!("Leverage already set for {symbol_str} to {leverage}");
441            }
442            Err(e) => log::error!("Failed to set leverage for {symbol_str}: {e}"),
443        }
444    }
445
446    async fn apply_position_modes_setting(&self) {
447        let Some(modes) = &self.config.position_mode else {
448            return;
449        };
450
451        for (symbol_str, mode) in modes {
452            self.apply_position_mode_entry(symbol_str, *mode).await;
453        }
454    }
455
456    async fn apply_position_mode_entry(&self, symbol_str: &str, mode: BybitPositionMode) {
457        let Some(symbol) = Self::parse_derivative_symbol(symbol_str) else {
458            return;
459        };
460        let result = self
461            .http_client
462            .switch_mode(
463                symbol.product_type(),
464                mode,
465                Some(symbol.raw_symbol().to_string()),
466                None,
467            )
468            .await;
469
470        match result {
471            Ok(_) => log::info!("Set symbol `{symbol_str}` position mode to `{mode:?}`"),
472            Err(e) if Self::is_unchanged_error(&e, "110025") => {
473                log::debug!("Symbol `{symbol_str}` position mode already set to `{mode:?}`");
474            }
475            Err(e) => log::error!("Failed to set position mode for {symbol_str}: {e}"),
476        }
477    }
478
479    async fn apply_margin_mode_setting(&self) -> anyhow::Result<()> {
480        let Some(margin_mode) = self.config.margin_mode else {
481            return Ok(());
482        };
483
484        let result = self.http_client.set_margin_mode(margin_mode).await;
485
486        match result {
487            Ok(_) => {
488                log::info!("Set account margin mode to {margin_mode:?}");
489                Ok(())
490            }
491            Err(e) if Self::is_unchanged_error(&e, "") => {
492                log::debug!("Margin mode already set to {margin_mode:?}");
493                Ok(())
494            }
495            Err(e) if Self::is_low_margin_error(&e) => {
496                log::warn!("Cannot set margin mode: {e}");
497                Ok(())
498            }
499            Err(e) => Err(anyhow::Error::from(e).context("failed to set margin mode")),
500        }
501    }
502
503    fn parse_derivative_symbol(symbol_str: &str) -> Option<BybitSymbol> {
504        let symbol = match BybitSymbol::new(symbol_str) {
505            Ok(s) => s,
506            Err(e) => {
507                log::warn!("Failed to parse symbol {symbol_str}: {e}");
508                return None;
509            }
510        };
511        matches!(
512            symbol.product_type(),
513            BybitProductType::Linear | BybitProductType::Inverse
514        )
515        .then_some(symbol)
516    }
517
518    fn is_unchanged_error<E: std::fmt::Display>(err: &E, code: &str) -> bool {
519        let msg = err.to_string().to_lowercase();
520        if msg.contains("not been modified") {
521            return true;
522        }
523        !code.is_empty() && msg.contains(code)
524    }
525
526    fn is_low_margin_error<E: std::fmt::Display>(err: &E) -> bool {
527        err.to_string()
528            .contains("needs to be equal to or greater than")
529    }
530
531    fn map_order_type(order_type: OrderType) -> anyhow::Result<(BybitOrderType, bool)> {
532        match order_type {
533            OrderType::Market => Ok((BybitOrderType::Market, false)),
534            OrderType::Limit => Ok((BybitOrderType::Limit, false)),
535            OrderType::StopMarket | OrderType::MarketIfTouched => {
536                Ok((BybitOrderType::Market, true))
537            }
538            OrderType::StopLimit | OrderType::LimitIfTouched => Ok((BybitOrderType::Limit, true)),
539            _ => anyhow::bail!("unsupported order type for Bybit: {order_type}"),
540        }
541    }
542
543    fn map_time_in_force(tif: TimeInForce, is_post_only: bool) -> BybitTimeInForce {
544        if is_post_only {
545            return BybitTimeInForce::PostOnly;
546        }
547
548        match tif {
549            TimeInForce::Gtc => BybitTimeInForce::Gtc,
550            TimeInForce::Ioc => BybitTimeInForce::Ioc,
551            TimeInForce::Fok => BybitTimeInForce::Fok,
552            _ => BybitTimeInForce::Gtc,
553        }
554    }
555
556    fn validate_bbo_params(
557        order: &OrderAny,
558        product_type: BybitProductType,
559        tp_sl: &BybitTpSlParams,
560    ) -> anyhow::Result<()> {
561        if !tp_sl.has_bbo() {
562            return Ok(());
563        }
564
565        anyhow::ensure!(
566            matches!(
567                product_type,
568                BybitProductType::Linear | BybitProductType::Inverse
569            ),
570            "`bbo_side_type` and `bbo_level` are only supported for Bybit linear and inverse products"
571        );
572
573        let order_type = order.order_type();
574        anyhow::ensure!(
575            matches!(
576                order_type,
577                OrderType::Limit | OrderType::StopLimit | OrderType::LimitIfTouched
578            ),
579            "`bbo_side_type` and `bbo_level` are not supported for order type {order_type:?}"
580        );
581
582        Ok(())
583    }
584
585    fn build_ws_place_params(
586        order: &OrderAny,
587        product_type: BybitProductType,
588        raw_symbol: &str,
589        tp_sl: &BybitTpSlParams,
590        position_idx: Option<BybitPositionIdx>,
591        smp_type: Option<BybitOrderSmpType>,
592    ) -> anyhow::Result<BybitWsPlaceOrderParams> {
593        let bybit_side = BybitOrderSide::from(order.order_side());
594        let (bybit_order_type, is_conditional) = Self::map_order_type(order.order_type())?;
595        let has_tp_sl = tp_sl.has_tp_sl();
596        let trigger_dir = trigger_direction(order.order_type(), order.order_side(), is_conditional);
597
598        Ok(BybitWsPlaceOrderParams {
599            category: product_type,
600            symbol: Ustr::from(raw_symbol),
601            side: bybit_side,
602            order_type: bybit_order_type,
603            qty: order.quantity().to_string(),
604            is_leverage: spot_leverage(product_type, tp_sl.is_leverage),
605            market_unit: spot_market_unit(
606                product_type,
607                bybit_order_type,
608                order.is_quote_quantity(),
609            ),
610            price: if tp_sl.has_bbo() {
611                None
612            } else {
613                order.price().map(|p: Price| p.to_string())
614            },
615            time_in_force: if bybit_order_type == BybitOrderType::Market {
616                None
617            } else {
618                Some(Self::map_time_in_force(
619                    order.time_in_force(),
620                    order.is_post_only(),
621                ))
622            },
623            order_link_id: Some(order.client_order_id().to_string()),
624            reduce_only: if order.is_reduce_only() {
625                Some(true)
626            } else {
627                None
628            },
629            close_on_trigger: tp_sl.close_on_trigger,
630            trigger_price: order.trigger_price().map(|p: Price| p.to_string()),
631            trigger_by: if is_conditional {
632                Some(resolve_trigger_type(order.trigger_type()))
633            } else {
634                None
635            },
636            trigger_direction: trigger_dir.map(|d| d as i32),
637            tpsl_mode: tp_sl.tpsl_mode.or(has_tp_sl.then_some(BybitTpSlMode::Full)),
638            take_profit: tp_sl.take_profit.map(|p| p.to_string()),
639            stop_loss: tp_sl.stop_loss.map(|p| p.to_string()),
640            tp_trigger_by: tp_sl.tp_trigger_by.or(tp_sl
641                .take_profit
642                .map(|_| resolve_trigger_type(order.trigger_type()))),
643            sl_trigger_by: tp_sl.sl_trigger_by.or(tp_sl
644                .stop_loss
645                .map(|_| resolve_trigger_type(order.trigger_type()))),
646            sl_trigger_price: tp_sl.sl_trigger_price.clone(),
647            tp_trigger_price: tp_sl.tp_trigger_price.clone(),
648            sl_order_type: tp_sl.sl_order_type,
649            tp_order_type: tp_sl.tp_order_type,
650            sl_limit_price: tp_sl.sl_limit_price.clone(),
651            tp_limit_price: tp_sl.tp_limit_price.clone(),
652            order_iv: tp_sl.order_iv.clone(),
653            smp_type,
654            mmp: tp_sl.mmp,
655            position_idx,
656            bbo_side_type: tp_sl.bbo_side_type,
657            bbo_level: tp_sl.bbo_level.clone(),
658        })
659    }
660}
661
662fn submit_rejection_reason(error: &anyhow::Error) -> Option<&str> {
663    for cause in error.chain() {
664        if let Some(submit_error) = cause.downcast_ref::<BybitSubmitOrderError>() {
665            return match submit_error {
666                BybitSubmitOrderError::Rejected { reason } => Some(reason.as_str()),
667                BybitSubmitOrderError::MissingOrderId
668                | BybitSubmitOrderError::PostSubmitLookup { .. } => None,
669            };
670        }
671
672        if let Some(BybitHttpError::BybitError {
673            error_code,
674            message,
675        }) = cause.downcast_ref()
676            && !is_bybit_ambiguous_order_error_code(i64::from(*error_code))
677        {
678            return Some(message.as_str());
679        }
680    }
681    None
682}
683
684#[async_trait(?Send)]
685impl ExecutionClient for BybitExecutionClient {
686    fn is_connected(&self) -> bool {
687        self.core.is_connected()
688    }
689
690    fn client_id(&self) -> ClientId {
691        self.core.client_id
692    }
693
694    fn account_id(&self) -> AccountId {
695        self.core.account_id
696    }
697
698    fn venue(&self) -> Venue {
699        *BYBIT_VENUE
700    }
701
702    fn oms_type(&self) -> OmsType {
703        self.core.oms_type
704    }
705
706    fn get_account(&self) -> Option<AccountAny> {
707        self.core.cache().account_owned(&self.core.account_id)
708    }
709
710    fn provides_bulk_position_coverage(&self, instrument_id: InstrumentId) -> bool {
711        // Resolve the suffix directly rather than through `get_product_type_for_instrument`, whose
712        // Linear fallback would claim coverage for an identifier this adapter cannot classify. A
713        // recognized derivative suffix is still uncovered when its product type is unconfigured,
714        // since the bulk loop only queries `product_types()`.
715        let Some(product_type) = BybitProductType::from_suffix(instrument_id.symbol.as_str())
716        else {
717            return false;
718        };
719
720        self.product_types().contains(&product_type)
721            && Self::provides_bulk_position_coverage_for_product_type(product_type)
722    }
723
724    async fn connect(&mut self) -> anyhow::Result<()> {
725        if self.core.is_connected() && self.pending_tasks.is_open() && self.session_tasks.is_open()
726        {
727            return Ok(());
728        }
729
730        if !self.pending_tasks.is_open() || !self.session_tasks.is_open() {
731            self.teardown_partial_connect().await?;
732        }
733
734        if !self.pending_tasks.is_open() {
735            self.await_pending_tasks().await?;
736            self.pending_tasks
737                .start_generation()
738                .map_err(|e| anyhow::anyhow!("Failed to start Bybit task generation: {e}"))?;
739        }
740
741        if !self.session_tasks.is_open() {
742            self.await_session_tasks().await?;
743            self.session_tasks
744                .start_generation()
745                .map_err(|e| anyhow::anyhow!("Failed to start Bybit session generation: {e}"))?;
746        }
747        let http_client = self.http_client.clone();
748        let ws_private = self.ws_private.clone();
749        let ws_trade = self.ws_trade.clone();
750        let setup_guard =
751            TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
752                http_client.cancel_all_requests();
753                ws_private.begin_shutdown();
754                ws_trade.begin_shutdown();
755            });
756
757        let connect_result: anyhow::Result<()> = async {
758            // Reset after a prior disconnect so REST calls are not short-circuited
759            self.http_client.reset_cancellation_token();
760
761            let product_types = self.product_types();
762
763            if !self.core.instruments_initialized() {
764                let mut all_instruments = Vec::new();
765
766                for product_type in &product_types {
767                    let instruments = self
768                        .http_client
769                        .request_instruments(*product_type, None, None)
770                        .await
771                        .with_context(|| {
772                            format!("failed to request Bybit instruments for {product_type:?}")
773                        })?;
774
775                    if instruments.is_empty() {
776                        log::warn!("No instruments returned for {product_type:?}");
777                        continue;
778                    }
779
780                    log::debug!("Loaded {} {product_type:?} instruments", instruments.len());
781
782                    self.http_client.cache_instruments(&instruments);
783                    all_instruments.extend(instruments);
784                }
785
786                if !all_instruments.is_empty() {
787                    let mut instruments_map = AHashMap::new();
788                    for instrument in &all_instruments {
789                        instruments_map.insert(instrument.id().symbol.inner(), instrument.clone());
790                    }
791                    self.instruments_cache = Arc::new(instruments_map);
792                }
793                self.core.set_instruments_initialized();
794            }
795
796            self.ws_private.set_account_id(self.core.account_id);
797            self.ws_trade.set_account_id(self.core.account_id);
798
799            self.ws_private.connect().await?;
800            self.ws_private.wait_until_active(10.0).await?;
801            log::debug!("Connected to private WebSocket");
802
803            let stream = self.ws_private.stream();
804            let emitter = self.emitter.clone();
805            let account_id = self.core.account_id;
806            let instruments = Arc::clone(&self.instruments_cache);
807            let state = Arc::clone(&self.dispatch_state);
808            let clock = self.clock;
809
810            self.session_tasks.spawn(async move {
811                pin_mut!(stream);
812                while let Some(message) = stream.next().await {
813                    dispatch_ws_message(
814                        &message,
815                        &emitter,
816                        &state,
817                        account_id,
818                        &instruments,
819                        clock,
820                    );
821                }
822            })?;
823
824            // Demo environment does not support Trade WebSocket API
825            if self.config.environment == BybitEnvironment::Demo {
826                log::warn!("Demo mode: Trade WebSocket not available, orders use HTTP REST API");
827            } else {
828                self.ws_trade.connect().await?;
829                self.ws_trade.wait_until_active(10.0).await?;
830                log::debug!("Connected to trade WebSocket");
831
832                let stream = self.ws_trade.stream();
833                let emitter = self.emitter.clone();
834                let account_id = self.core.account_id;
835                let instruments = Arc::clone(&self.instruments_cache);
836                let state = Arc::clone(&self.dispatch_state);
837                let clock = self.clock;
838
839                self.session_tasks.spawn(async move {
840                    pin_mut!(stream);
841                    while let Some(message) = stream.next().await {
842                        dispatch_ws_message(
843                            &message,
844                            &emitter,
845                            &state,
846                            account_id,
847                            &instruments,
848                            clock,
849                        );
850                    }
851                })?;
852            }
853
854            if self.config.auto_repay_spot_borrows {
855                let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
856                self.dispatch_state.set_repay_sender(tx);
857
858                let http_client = self.http_client.clone();
859                let clock = self.clock;
860                self.session_tasks.spawn(async move {
861                    crate::repay::run_spot_repay_consumer(rx, http_client, clock).await;
862                })?;
863            }
864
865            self.ws_private.subscribe_orders().await?;
866            self.ws_private.subscribe_executions().await?;
867            self.ws_private.subscribe_positions().await?;
868            self.ws_private.subscribe_wallet().await?;
869
870            self.apply_account_configuration().await?;
871
872            let account_state = self
873                .http_client
874                .request_account_state(BybitAccountType::Unified, self.core.account_id)
875                .await
876                .context("failed to request Bybit account state")?;
877
878            if !account_state.balances.is_empty() {
879                log::debug!(
880                    "Received account state with {} balance(s)",
881                    account_state.balances.len()
882                );
883            }
884            self.emitter.send_account_state(account_state);
885
886            self.await_account_registered(30.0).await?;
887
888            Ok(())
889        }
890        .await;
891
892        if let Err(e) = connect_result {
893            if let Err(teardown_error) = self.teardown_partial_connect().await {
894                return Err(e.context(format!("Bybit startup teardown failed: {teardown_error}")));
895            }
896            return Err(e);
897        }
898
899        setup_guard.disarm();
900        self.core.set_connected();
901        log::info!("Connected: client_id={}", self.core.client_id);
902        Ok(())
903    }
904
905    async fn disconnect(&mut self) -> anyhow::Result<()> {
906        self.teardown_partial_connect().await?;
907        log::info!("Disconnected: client_id={}", self.core.client_id);
908        Ok(())
909    }
910
911    fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
912        self.update_account_state();
913        Ok(())
914    }
915
916    fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
917        let instrument_id = cmd.instrument_id;
918        let product_type = self.get_product_type_for_instrument(instrument_id);
919        let client_order_id = cmd.client_order_id;
920        let venue_order_id = cmd.venue_order_id;
921        let account_id = self.core.account_id;
922        let http_client = self.http_client.clone();
923        let emitter = self.emitter.clone();
924
925        self.spawn_task("query_order", async move {
926            match http_client
927                .query_order(
928                    account_id,
929                    product_type,
930                    instrument_id,
931                    Some(client_order_id),
932                    venue_order_id,
933                )
934                .await
935            {
936                Ok(Some(report)) => {
937                    emitter.send_order_status_report(report);
938                }
939                Ok(None) => {
940                    log::warn!("Order not found: client_order_id={client_order_id}, venue_order_id={venue_order_id:?}");
941                }
942                Err(e) => {
943                    log::error!("Failed to query order: {e}");
944                }
945            }
946            Ok(())
947        });
948
949        Ok(())
950    }
951
952    fn generate_account_state(
953        &self,
954        balances: Vec<AccountBalance>,
955        margins: Vec<MarginBalance>,
956        reported: bool,
957        ts_event: UnixNanos,
958        info: Option<Params>,
959    ) -> anyhow::Result<()> {
960        self.emitter
961            .emit_account_state(balances, margins, reported, ts_event, info);
962        Ok(())
963    }
964
965    fn start(&mut self) -> anyhow::Result<()> {
966        if self.core.is_started() {
967            return Ok(());
968        }
969
970        let sender = get_exec_event_sender();
971        self.emitter.set_sender(sender);
972        self.core.set_started();
973
974        let http_client = self.http_client.clone();
975        let product_types = self.config.product_types.clone();
976
977        self.session_tasks.spawn(async move {
978            let mut all_instruments = Vec::new();
979
980            for product_type in product_types {
981                match http_client
982                    .request_instruments(product_type, None, None)
983                    .await
984                {
985                    Ok(instruments) => {
986                        if instruments.is_empty() {
987                            log::warn!("No instruments returned for {product_type:?}");
988                            continue;
989                        }
990                        http_client.cache_instruments(&instruments);
991                        all_instruments.extend(instruments);
992                    }
993                    Err(e) => {
994                        log::error!("Failed to request instruments for {product_type:?}: {e}");
995                    }
996                }
997            }
998
999            if all_instruments.is_empty() {
1000                log::warn!(
1001                    "Instrument bootstrap yielded no instruments; WebSocket submissions may fail"
1002                );
1003            } else {
1004                log::debug!("Instruments initialized: count={}", all_instruments.len());
1005            }
1006        })?;
1007
1008        log::info!(
1009            "Started: client_id={}, account_id={}, account_type={:?}, product_types={:?}, environment={:?}, proxy_url={:?}",
1010            self.core.client_id,
1011            self.core.account_id,
1012            self.core.account_type,
1013            self.config.product_types,
1014            self.config.environment,
1015            self.config.proxy_url,
1016        );
1017        Ok(())
1018    }
1019
1020    fn stop(&mut self) -> anyhow::Result<()> {
1021        if self.core.is_stopped() {
1022            return Ok(());
1023        }
1024
1025        self.core.set_stopped();
1026        self.core.set_disconnected();
1027
1028        self.abort_session_tasks();
1029        self.abort_pending_tasks();
1030        self.http_client.cancel_all_requests();
1031        self.ws_private.begin_shutdown();
1032        self.ws_trade.begin_shutdown();
1033        log::info!("Stopped: client_id={}", self.core.client_id);
1034        Ok(())
1035    }
1036
1037    async fn generate_order_status_report(
1038        &self,
1039        cmd: &GenerateOrderStatusReport,
1040    ) -> anyhow::Result<Option<OrderStatusReport>> {
1041        let Some(instrument_id) = cmd.instrument_id else {
1042            log::warn!("generate_order_status_report requires instrument_id: {cmd}");
1043            return Ok(None);
1044        };
1045
1046        let product_type = self.get_product_type_for_instrument(instrument_id);
1047
1048        let mut reports = self
1049            .http_client
1050            .request_order_status_reports(
1051                self.core.account_id,
1052                product_type,
1053                Some(instrument_id),
1054                false,
1055                None,
1056                None,
1057                None,
1058            )
1059            .await?;
1060
1061        if let Some(client_order_id) = cmd.client_order_id {
1062            reports.retain(|report| report.client_order_id == Some(client_order_id));
1063        }
1064
1065        if let Some(venue_order_id) = cmd.venue_order_id {
1066            reports.retain(|report| report.venue_order_id.as_str() == venue_order_id.as_str());
1067        }
1068
1069        let report = reports.into_iter().next();
1070        if let Some(report) = &report {
1071            self.cache_reconciliation_order_identity(report);
1072        }
1073
1074        Ok(report)
1075    }
1076
1077    async fn generate_order_status_reports(
1078        &self,
1079        cmd: &GenerateOrderStatusReports,
1080    ) -> anyhow::Result<Vec<OrderStatusReport>> {
1081        let mut reports = Vec::new();
1082
1083        if let Some(instrument_id) = cmd.instrument_id {
1084            let product_type = self.get_product_type_for_instrument(instrument_id);
1085            let mut fetched = self
1086                .http_client
1087                .request_order_status_reports(
1088                    self.core.account_id,
1089                    product_type,
1090                    Some(instrument_id),
1091                    cmd.open_only,
1092                    None,
1093                    None,
1094                    None,
1095                )
1096                .await?;
1097            reports.append(&mut fetched);
1098        } else {
1099            for product_type in self.product_types() {
1100                let mut fetched = self
1101                    .http_client
1102                    .request_order_status_reports(
1103                        self.core.account_id,
1104                        product_type,
1105                        None,
1106                        cmd.open_only,
1107                        None,
1108                        None,
1109                        None,
1110                    )
1111                    .await?;
1112                reports.append(&mut fetched);
1113            }
1114        }
1115
1116        if let Some(start) = cmd.start {
1117            reports.retain(|r| r.ts_last >= start);
1118        }
1119
1120        if let Some(end) = cmd.end {
1121            reports.retain(|r| r.ts_last <= end);
1122        }
1123
1124        for report in &reports {
1125            self.cache_reconciliation_order_identity(report);
1126        }
1127
1128        Ok(reports)
1129    }
1130
1131    async fn generate_fill_reports(
1132        &self,
1133        cmd: GenerateFillReports,
1134    ) -> anyhow::Result<Vec<FillReport>> {
1135        let start_ms = nanos_to_millis(cmd.start);
1136        let end_ms = nanos_to_millis(cmd.end);
1137        let mut reports = Vec::new();
1138
1139        if let Some(instrument_id) = cmd.instrument_id {
1140            let product_type = self.get_product_type_for_instrument(instrument_id);
1141            let mut fetched = self
1142                .http_client
1143                .request_fill_reports(
1144                    self.core.account_id,
1145                    product_type,
1146                    Some(instrument_id),
1147                    start_ms,
1148                    end_ms,
1149                    None,
1150                )
1151                .await?;
1152            reports.append(&mut fetched);
1153        } else {
1154            for product_type in self.product_types() {
1155                let mut fetched = self
1156                    .http_client
1157                    .request_fill_reports(
1158                        self.core.account_id,
1159                        product_type,
1160                        None,
1161                        start_ms,
1162                        end_ms,
1163                        None,
1164                    )
1165                    .await?;
1166                reports.append(&mut fetched);
1167            }
1168        }
1169
1170        if let Some(venue_order_id) = cmd.venue_order_id {
1171            reports.retain(|report| report.venue_order_id.as_str() == venue_order_id.as_str());
1172        }
1173
1174        Ok(reports)
1175    }
1176
1177    async fn generate_position_status_reports(
1178        &self,
1179        cmd: &GeneratePositionStatusReports,
1180    ) -> anyhow::Result<Vec<PositionStatusReport>> {
1181        if let Some(instrument_id) = cmd.instrument_id {
1182            let product_type = self.get_product_type_for_instrument(instrument_id);
1183            self.http_client
1184                .request_position_status_reports(
1185                    self.core.account_id,
1186                    product_type,
1187                    Some(instrument_id),
1188                )
1189                .await
1190        } else {
1191            self.generate_bulk_position_status_reports(self.product_types())
1192                .await
1193        }
1194    }
1195
1196    async fn generate_mass_status(
1197        &self,
1198        lookback_mins: Option<u64>,
1199    ) -> anyhow::Result<Option<ExecutionMassStatus>> {
1200        log::info!("Generating ExecutionMassStatus (lookback_mins={lookback_mins:?})");
1201
1202        let ts_now = self.clock.get_time_ns();
1203
1204        let start = lookback_mins
1205            .map(DurationNanos::try_from_mins)
1206            .transpose()?
1207            .map(|lookback| ts_now.saturating_sub(lookback));
1208
1209        let order_cmd = GenerateOrderStatusReportsBuilder::default()
1210            .ts_init(ts_now)
1211            .open_only(false)
1212            .start(start)
1213            .build()
1214            .map_err(|e| anyhow::anyhow!("{e}"))?;
1215
1216        let fill_cmd = GenerateFillReportsBuilder::default()
1217            .ts_init(ts_now)
1218            .start(start)
1219            .build()
1220            .map_err(|e| anyhow::anyhow!("{e}"))?;
1221
1222        let position_reports_fut = async {
1223            let product_types = self.product_types();
1224            let skip_spot = product_types.iter().any(|product_type| {
1225                !Self::provides_bulk_position_coverage_for_product_type(*product_type)
1226            });
1227
1228            if skip_spot {
1229                log::warn!(
1230                    "SPOT mass-status position coverage is unavailable because wallet balances cannot be attributed to pairs"
1231                );
1232            }
1233
1234            self.generate_bulk_position_status_reports(product_types)
1235                .await
1236        };
1237
1238        let (order_reports, fill_reports, position_reports) = tokio::try_join!(
1239            self.generate_order_status_reports(&order_cmd),
1240            self.generate_fill_reports(fill_cmd),
1241            position_reports_fut,
1242        )?;
1243
1244        log::info!("Received {} OrderStatusReports", order_reports.len());
1245        log::info!("Received {} FillReports", fill_reports.len());
1246        log::info!("Received {} PositionReports", position_reports.len());
1247
1248        let mut mass_status = ExecutionMassStatus::new(
1249            self.core.client_id,
1250            self.core.account_id,
1251            *BYBIT_VENUE,
1252            ts_now,
1253            None,
1254        );
1255
1256        mass_status.add_order_reports(order_reports);
1257        mass_status.add_fill_reports(fill_reports);
1258        mass_status.add_position_reports(position_reports);
1259
1260        Ok(Some(mass_status))
1261    }
1262
1263    fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
1264        let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
1265        if order.is_closed() {
1266            log::warn!("Cannot submit closed order {}", order.client_order_id());
1267            return Ok(());
1268        }
1269
1270        let instrument_id = order.instrument_id();
1271        let product_type = self.get_product_type_for_instrument(instrument_id);
1272
1273        // Validate order params before emitting submitted event
1274        if Self::map_order_type(order.order_type()).is_err() {
1275            let denied = OrderDeniedReason::UnsupportedOrderType {
1276                order_type: order.order_type(),
1277            };
1278            self.emitter.emit_order_denied(&order, &denied.to_string());
1279            return Ok(());
1280        }
1281
1282        let tp_sl = match parse_bybit_tp_sl_params(cmd.params.as_ref()) {
1283            Ok(p) => p,
1284            Err(e) => {
1285                let denied = OrderDeniedReason::ValidationFailed {
1286                    detail: e.to_string(),
1287                };
1288                self.emitter.emit_order_denied(&order, &denied.to_string());
1289                return Ok(());
1290            }
1291        };
1292
1293        if let Err(e) = Self::validate_bbo_params(&order, product_type, &tp_sl) {
1294            let denied = OrderDeniedReason::ValidationFailed {
1295                detail: e.to_string(),
1296            };
1297            self.emitter.emit_order_denied(&order, &denied.to_string());
1298            return Ok(());
1299        }
1300
1301        // The demo HTTP create-order entry cannot carry TP/SL trigger prices (only the mainnet
1302        // WS path can), so deny rather than submit an order missing the user's trigger prices.
1303        if self.config.environment == BybitEnvironment::Demo
1304            && (tp_sl.tp_trigger_price.is_some() || tp_sl.sl_trigger_price.is_some())
1305        {
1306            let denied = OrderDeniedReason::UnsupportedTpSl {
1307                detail: "TP/SL trigger prices are not supported in demo mode".to_string(),
1308            };
1309            self.emitter.emit_order_denied(&order, &denied.to_string());
1310            return Ok(());
1311        }
1312
1313        log::debug!("OrderSubmitted client_order_id={}", order.client_order_id());
1314        self.emitter.emit_order_submitted(&order);
1315
1316        let client_order_id = order.client_order_id();
1317        let strategy_id = order.strategy_id();
1318        let emitter = self.emitter.clone();
1319        let clock = self.clock;
1320
1321        let bybit_side = BybitOrderSide::from(order.order_side());
1322        let position_idx = self.resolve_position_idx(
1323            instrument_id,
1324            bybit_side,
1325            order.is_reduce_only(),
1326            tp_sl.position_idx,
1327        );
1328        let venue_position_id =
1329            position_idx.and_then(|idx| make_hedge_venue_position_id(instrument_id, idx as i32));
1330
1331        self.dispatch_state.order_identities.insert(
1332            client_order_id,
1333            OrderIdentity {
1334                instrument_id,
1335                strategy_id,
1336                order_side: order.order_side(),
1337                order_type: order.order_type(),
1338                venue_position_id,
1339            },
1340        );
1341
1342        // Seed for BBO reconciliation: venue may replace the submitted price
1343        self.dispatch_state.order_snapshots.insert(
1344            client_order_id,
1345            OrderStateSnapshot {
1346                quantity: order.quantity(),
1347                price: order.price(),
1348                trigger_price: order.trigger_price(),
1349            },
1350        );
1351
1352        let smp_type = self.resolve_smp_type(tp_sl.smp_type);
1353
1354        if self.config.environment == BybitEnvironment::Demo {
1355            let http_client = self.http_client.clone();
1356            let account_id = self.core.account_id;
1357            let order_side = order.order_side();
1358            let order_type = order.order_type();
1359            let quantity = order.quantity();
1360            let time_in_force = order.time_in_force();
1361            let price = order.price();
1362            let trigger_price = order.trigger_price();
1363            let post_only = order.is_post_only();
1364            let reduce_only = order.is_reduce_only();
1365            let is_quote_quantity = order.is_quote_quantity();
1366            let is_leverage = tp_sl.is_leverage;
1367            let bbo_side_type = tp_sl.bbo_side_type;
1368            let bbo_level = tp_sl.bbo_level.clone();
1369            let native_tp_sl = tp_sl.to_native_tp_sl();
1370            let dispatch_state = Arc::clone(&self.dispatch_state);
1371
1372            self.spawn_task("submit_order_http", async move {
1373                let native_tp_sl_ref = (!native_tp_sl.is_empty()).then_some(&native_tp_sl);
1374                let result = http_client
1375                    .submit_order(
1376                        account_id,
1377                        product_type,
1378                        instrument_id,
1379                        client_order_id,
1380                        order_side,
1381                        order_type,
1382                        quantity,
1383                        Some(time_in_force),
1384                        price,
1385                        trigger_price,
1386                        Some(post_only),
1387                        reduce_only,
1388                        is_quote_quantity,
1389                        is_leverage,
1390                        position_idx,
1391                        bbo_side_type,
1392                        bbo_level,
1393                        smp_type,
1394                        native_tp_sl_ref,
1395                    )
1396                    .await;
1397
1398                if let Err(e) = result {
1399                    if let Some(reason) = submit_rejection_reason(&e) {
1400                        dispatch_state.order_identities.remove(&client_order_id);
1401                        dispatch_state.order_snapshots.remove(&client_order_id);
1402                        let ts_event = clock.get_time_ns();
1403                        emitter.emit_order_rejected_event(
1404                            strategy_id,
1405                            instrument_id,
1406                            client_order_id,
1407                            reason,
1408                            ts_event,
1409                            bybit_rejection_due_post_only(reason),
1410                        );
1411                        anyhow::bail!("submit order rejected: {reason}");
1412                    }
1413
1414                    log::warn!(
1415                        "Submit failure without confirmed venue rejection for {client_order_id}: \
1416                         {e}; awaiting reconciliation",
1417                    );
1418                    return Ok(());
1419                }
1420
1421                Ok(())
1422            });
1423
1424            return Ok(());
1425        }
1426
1427        let raw_symbol = extract_raw_symbol(instrument_id.symbol.as_str());
1428        let params = Self::build_ws_place_params(
1429            &order,
1430            product_type,
1431            raw_symbol,
1432            &tp_sl,
1433            position_idx,
1434            smp_type,
1435        )?;
1436
1437        let ws_trade = self.ws_trade.clone();
1438        let dispatch_state = Arc::clone(&self.dispatch_state);
1439
1440        self.spawn_task("submit_order", async move {
1441            let req_id = UUID4::new().to_string();
1442            dispatch_state.pending_requests.insert(
1443                req_id.clone(),
1444                (vec![client_order_id], vec![None], PendingOperation::Place),
1445            );
1446
1447            if let Err(e) = ws_trade.place_order_with_id(params, req_id.clone()).await {
1448                dispatch_state.pending_requests.remove(&req_id);
1449                dispatch_state.order_identities.remove(&client_order_id);
1450                dispatch_state.order_snapshots.remove(&client_order_id);
1451                let reason = e.to_string();
1452                emitter.emit_order_rejected_event(
1453                    strategy_id,
1454                    instrument_id,
1455                    client_order_id,
1456                    &reason,
1457                    clock.get_time_ns(),
1458                    false,
1459                );
1460            }
1461
1462            Ok(())
1463        });
1464
1465        Ok(())
1466    }
1467
1468    fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
1469        if cmd.order_list.client_order_ids.is_empty() {
1470            return Ok(());
1471        }
1472
1473        let tp_sl = match parse_bybit_tp_sl_params(cmd.params.as_ref()) {
1474            Ok(p) => p,
1475            Err(e) => {
1476                let cache = self.core.cache();
1477                let denied = OrderDeniedReason::ValidationFailed {
1478                    detail: e.to_string(),
1479                }
1480                .to_string();
1481
1482                for cid in &cmd.order_list.client_order_ids {
1483                    if let Some(order) = cache.order(cid) {
1484                        self.emitter.emit_order_denied(&order, &denied);
1485                    }
1486                }
1487                return Ok(());
1488            }
1489        };
1490
1491        let instrument_id = cmd.instrument_id;
1492        let product_type = self.get_product_type_for_instrument(instrument_id);
1493
1494        // The demo HTTP create-order entry cannot carry TP/SL trigger prices (only the mainnet
1495        // WS path can), so deny rather than submit orders missing the user's trigger prices.
1496        if self.config.environment == BybitEnvironment::Demo
1497            && (tp_sl.tp_trigger_price.is_some() || tp_sl.sl_trigger_price.is_some())
1498        {
1499            let cache = self.core.cache();
1500            let denied = OrderDeniedReason::UnsupportedTpSl {
1501                detail: "TP/SL trigger prices are not supported in demo mode".to_string(),
1502            }
1503            .to_string();
1504
1505            for cid in &cmd.order_list.client_order_ids {
1506                if let Some(order) = cache.order(cid) {
1507                    self.emitter.emit_order_denied(&order, &denied);
1508                }
1509            }
1510            return Ok(());
1511        }
1512
1513        let strategy_id = cmd.strategy_id;
1514
1515        let mut valid_orders = Vec::with_capacity(cmd.order_list.client_order_ids.len());
1516        {
1517            let cache = self.core.cache();
1518            let order_list_id = cmd.order_list.id;
1519            let list_denied = OrderDeniedReason::OrderListDenied { order_list_id };
1520            // (offending leg, reason for the offending leg, reason for the remaining legs). A
1521            // single offending leg carries its specific reason and the rest render
1522            // `ORDER_LIST_DENIED`; a list-level failure renders the same reason for every leg.
1523            let mut denial: Option<(ClientOrderId, OrderDeniedReason, OrderDeniedReason)> = None;
1524
1525            for cid in &cmd.order_list.client_order_ids {
1526                let Some(order) = cache.order(cid) else {
1527                    let reason = OrderDeniedReason::OrderListIncomplete { order_list_id };
1528                    denial = Some((*cid, reason.clone(), reason));
1529                    break;
1530                };
1531
1532                if order.is_closed() {
1533                    denial = Some((
1534                        *cid,
1535                        OrderDeniedReason::ValidationFailed {
1536                            detail: format!("cannot submit closed order {cid}"),
1537                        },
1538                        list_denied,
1539                    ));
1540                    break;
1541                }
1542
1543                if Self::map_order_type(order.order_type()).is_err() {
1544                    denial = Some((
1545                        *cid,
1546                        OrderDeniedReason::UnsupportedOrderType {
1547                            order_type: order.order_type(),
1548                        },
1549                        list_denied,
1550                    ));
1551                    break;
1552                }
1553
1554                if let Err(e) = Self::validate_bbo_params(&order, product_type, &tp_sl) {
1555                    denial = Some((
1556                        *cid,
1557                        OrderDeniedReason::ValidationFailed {
1558                            detail: e.to_string(),
1559                        },
1560                        list_denied,
1561                    ));
1562                    break;
1563                }
1564
1565                valid_orders.push(order.clone());
1566            }
1567
1568            // Deny the entire list if any leg fails validation
1569            if let Some((offender, offender_reason, rest_reason)) = denial {
1570                let offender_reason = offender_reason.to_string();
1571                let rest_reason = rest_reason.to_string();
1572
1573                for cid in &cmd.order_list.client_order_ids {
1574                    if let Some(order) = cache.order(cid) {
1575                        let reason = if *cid == offender {
1576                            offender_reason.as_str()
1577                        } else {
1578                            rest_reason.as_str()
1579                        };
1580                        self.emitter.emit_order_denied(&order, reason);
1581                    }
1582                }
1583                return Ok(());
1584            }
1585        }
1586
1587        if valid_orders.is_empty() {
1588            return Ok(());
1589        }
1590
1591        for order in &valid_orders {
1592            self.emitter.emit_order_submitted(order);
1593            let bybit_side = BybitOrderSide::from(order.order_side());
1594            let position_idx = self.resolve_position_idx(
1595                instrument_id,
1596                bybit_side,
1597                order.is_reduce_only(),
1598                tp_sl.position_idx,
1599            );
1600            let venue_position_id = position_idx
1601                .and_then(|idx| make_hedge_venue_position_id(instrument_id, idx as i32));
1602            self.dispatch_state.order_identities.insert(
1603                order.client_order_id(),
1604                OrderIdentity {
1605                    instrument_id,
1606                    strategy_id,
1607                    order_side: order.order_side(),
1608                    order_type: order.order_type(),
1609                    venue_position_id,
1610                },
1611            );
1612            self.dispatch_state.order_snapshots.insert(
1613                order.client_order_id(),
1614                OrderStateSnapshot {
1615                    quantity: order.quantity(),
1616                    price: order.price(),
1617                    trigger_price: order.trigger_price(),
1618                },
1619            );
1620        }
1621
1622        let emitter = self.emitter.clone();
1623        let clock = self.clock;
1624
1625        let smp_type = self.resolve_smp_type(tp_sl.smp_type);
1626
1627        // Demo mode: submit individually via HTTP
1628        if self.config.environment == BybitEnvironment::Demo {
1629            let http_client = self.http_client.clone();
1630            let account_id = self.core.account_id;
1631            let is_leverage = tp_sl.is_leverage;
1632            let bbo_side_type = tp_sl.bbo_side_type;
1633            let bbo_level = tp_sl.bbo_level.clone();
1634            let native_tp_sl = tp_sl.to_native_tp_sl();
1635            let dispatch_state = Arc::clone(&self.dispatch_state);
1636
1637            let order_data: Vec<_> = valid_orders
1638                .iter()
1639                .map(|o| {
1640                    let bybit_side = BybitOrderSide::from(o.order_side());
1641                    let position_idx = self.resolve_position_idx(
1642                        instrument_id,
1643                        bybit_side,
1644                        o.is_reduce_only(),
1645                        tp_sl.position_idx,
1646                    );
1647                    (
1648                        o.client_order_id(),
1649                        o.order_side(),
1650                        o.order_type(),
1651                        o.quantity(),
1652                        o.time_in_force(),
1653                        o.price(),
1654                        o.trigger_price(),
1655                        o.is_post_only(),
1656                        o.is_reduce_only(),
1657                        o.is_quote_quantity(),
1658                        position_idx,
1659                    )
1660                })
1661                .collect();
1662
1663            self.spawn_task("submit_order_list_http", async move {
1664                let native_tp_sl_ref = (!native_tp_sl.is_empty()).then_some(&native_tp_sl);
1665
1666                for (
1667                    cid,
1668                    side,
1669                    otype,
1670                    qty,
1671                    tif,
1672                    price,
1673                    trigger,
1674                    post_only,
1675                    reduce,
1676                    quote_qty,
1677                    position_idx,
1678                ) in order_data
1679                {
1680                    if let Err(e) = http_client
1681                        .submit_order(
1682                            account_id,
1683                            product_type,
1684                            instrument_id,
1685                            cid,
1686                            side,
1687                            otype,
1688                            qty,
1689                            Some(tif),
1690                            price,
1691                            trigger,
1692                            Some(post_only),
1693                            reduce,
1694                            quote_qty,
1695                            is_leverage,
1696                            position_idx,
1697                            bbo_side_type,
1698                            bbo_level.clone(),
1699                            smp_type,
1700                            native_tp_sl_ref,
1701                        )
1702                        .await
1703                    {
1704                        if let Some(reason) = submit_rejection_reason(&e) {
1705                            dispatch_state.order_identities.remove(&cid);
1706                            dispatch_state.order_snapshots.remove(&cid);
1707                            let ts_event = clock.get_time_ns();
1708                            emitter.emit_order_rejected_event(
1709                                strategy_id,
1710                                instrument_id,
1711                                cid,
1712                                reason,
1713                                ts_event,
1714                                bybit_rejection_due_post_only(reason),
1715                            );
1716                            continue;
1717                        }
1718
1719                        log::warn!(
1720                            "Submit failure without confirmed venue rejection for {cid}: {e}; \
1721                             awaiting reconciliation",
1722                        );
1723                    }
1724                }
1725                Ok(())
1726            });
1727
1728            return Ok(());
1729        }
1730
1731        // Live mode: batch submit via WebSocket
1732        let raw_symbol = extract_raw_symbol(instrument_id.symbol.as_str());
1733
1734        let mut order_params = Vec::with_capacity(valid_orders.len());
1735        let mut client_order_ids = Vec::with_capacity(valid_orders.len());
1736
1737        for order in &valid_orders {
1738            let bybit_side = BybitOrderSide::from(order.order_side());
1739            let position_idx = self.resolve_position_idx(
1740                instrument_id,
1741                bybit_side,
1742                order.is_reduce_only(),
1743                tp_sl.position_idx,
1744            );
1745            let params = Self::build_ws_place_params(
1746                order,
1747                product_type,
1748                raw_symbol,
1749                &tp_sl,
1750                position_idx,
1751                smp_type,
1752            )
1753            .expect("validated above");
1754            order_params.push(params);
1755            client_order_ids.push(order.client_order_id());
1756        }
1757
1758        let ws_trade = self.ws_trade.clone();
1759        let dispatch_state = Arc::clone(&self.dispatch_state);
1760
1761        self.spawn_task("submit_order_list", async move {
1762            let req_ids = BybitWebSocketClient::batch_request_ids(product_type, order_params.len());
1763            for (req_id, chunk_cids) in req_ids.iter().zip(
1764                client_order_ids
1765                    .chunks(batch_send_limit(product_type))
1766                    .map(|chunk| chunk.to_vec()),
1767            ) {
1768                let chunk_voids = vec![None; chunk_cids.len()];
1769                dispatch_state.pending_requests.insert(
1770                    req_id.clone(),
1771                    (chunk_cids, chunk_voids, PendingOperation::Place),
1772                );
1773            }
1774
1775            if let Err(e) = ws_trade
1776                .batch_place_orders_with_ids(order_params, req_ids.clone())
1777                .await
1778            {
1779                for req_id in req_ids {
1780                    dispatch_state.pending_requests.remove(&req_id);
1781                }
1782                let reason = e.to_string();
1783                let ts_event = clock.get_time_ns();
1784
1785                for client_order_id in client_order_ids {
1786                    dispatch_state.order_identities.remove(&client_order_id);
1787                    dispatch_state.order_snapshots.remove(&client_order_id);
1788                    emitter.emit_order_rejected_event(
1789                        strategy_id,
1790                        instrument_id,
1791                        client_order_id,
1792                        &reason,
1793                        ts_event,
1794                        false,
1795                    );
1796                }
1797            }
1798            Ok(())
1799        });
1800
1801        Ok(())
1802    }
1803
1804    fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1805        let instrument_id = cmd.instrument_id;
1806        let product_type = self.get_product_type_for_instrument(instrument_id);
1807        let client_order_id = cmd.client_order_id;
1808        let strategy_id = cmd.strategy_id;
1809        let venue_order_id = cmd.venue_order_id;
1810        let emitter = self.emitter.clone();
1811        let clock = self.clock;
1812
1813        let has_order_iv = cmd
1814            .params
1815            .as_ref()
1816            .and_then(|p| p.get("order_iv"))
1817            .is_some();
1818
1819        if self.config.environment == BybitEnvironment::Demo && has_order_iv {
1820            log::warn!(
1821                "Modify command failed local validation for {client_order_id}: {}",
1822                "Option params (order_iv) are not supported in demo mode",
1823            );
1824            return Ok(());
1825        }
1826
1827        if self.config.environment == BybitEnvironment::Demo {
1828            let http_client = self.http_client.clone();
1829            let account_id = self.core.account_id;
1830            let quantity = cmd.quantity;
1831            let price = cmd.price;
1832
1833            self.spawn_task("modify_order_http", async move {
1834                let result = http_client
1835                    .modify_order(
1836                        account_id,
1837                        product_type,
1838                        instrument_id,
1839                        Some(client_order_id),
1840                        venue_order_id,
1841                        quantity,
1842                        price,
1843                    )
1844                    .await;
1845
1846                if let Err(e) = result {
1847                    match classify_modify_http_failure(&e) {
1848                        CommandFailure::VenueRejected(reason) => {
1849                            let ts_event = clock.get_time_ns();
1850                            emitter.emit_order_modify_rejected_event(
1851                                strategy_id,
1852                                instrument_id,
1853                                client_order_id,
1854                                venue_order_id,
1855                                &format!("modify-order-error: {reason}"),
1856                                ts_event,
1857                            );
1858                            anyhow::bail!("modify order rejected: {reason}");
1859                        }
1860                        CommandFailure::NotSent(reason) => {
1861                            log::warn!(
1862                                "HTTP modify command failed local validation for {client_order_id}: {reason}"
1863                            );
1864                        }
1865                        CommandFailure::Ambiguous(reason) => {
1866                            log::warn!(
1867                                "Ambiguous HTTP modify failure for {client_order_id}, awaiting reconciliation: {reason}"
1868                            );
1869                        }
1870                    }
1871                }
1872
1873                Ok(())
1874            });
1875
1876            return Ok(());
1877        }
1878
1879        let raw_symbol = extract_raw_symbol(instrument_id.symbol.as_str());
1880
1881        let order_iv = if let Some(value) = cmd.params.as_ref().and_then(|p| p.get("order_iv")) {
1882            match get_price_str(cmd.params.as_ref().unwrap(), "order_iv") {
1883                Some(s) => Some(s),
1884                None => {
1885                    log::warn!(
1886                        "Modify command failed local validation for {client_order_id}: invalid type for 'order_iv': {value}, expected string or number",
1887                    );
1888                    return Ok(());
1889                }
1890            }
1891        } else {
1892            None
1893        };
1894
1895        let params = BybitWsAmendOrderParams {
1896            category: product_type,
1897            symbol: Ustr::from(raw_symbol),
1898            order_id: cmd.venue_order_id.map(|v| v.to_string()),
1899            order_link_id: Some(cmd.client_order_id.to_string()),
1900            qty: cmd.quantity.map(|q| q.to_string()),
1901            price: cmd.price.map(|p| p.to_string()),
1902            trigger_price: None,
1903            take_profit: None,
1904            stop_loss: None,
1905            tp_trigger_by: None,
1906            sl_trigger_by: None,
1907            order_iv,
1908        };
1909
1910        let ws_trade = self.ws_trade.clone();
1911        let dispatch_state = Arc::clone(&self.dispatch_state);
1912
1913        self.spawn_task("modify_order", async move {
1914            let req_id = UUID4::new().to_string();
1915            dispatch_state.pending_requests.insert(
1916                req_id.clone(),
1917                (
1918                    vec![client_order_id],
1919                    vec![venue_order_id],
1920                    PendingOperation::Amend,
1921                ),
1922            );
1923
1924            if let Err(e) = ws_trade.amend_order_with_id(params, req_id.clone()).await {
1925                dispatch_state.pending_requests.remove(&req_id);
1926                emitter.emit_order_modify_rejected_event(
1927                    strategy_id,
1928                    instrument_id,
1929                    client_order_id,
1930                    venue_order_id,
1931                    &e.to_string(),
1932                    clock.get_time_ns(),
1933                );
1934            }
1935
1936            Ok(())
1937        });
1938
1939        Ok(())
1940    }
1941
1942    fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
1943        let instrument_id = cmd.instrument_id;
1944        let product_type = self.get_product_type_for_instrument(instrument_id);
1945        let client_order_id = cmd.client_order_id;
1946        let strategy_id = cmd.strategy_id;
1947        let venue_order_id = cmd.venue_order_id;
1948        let emitter = self.emitter.clone();
1949        let clock = self.clock;
1950
1951        if self.config.environment == BybitEnvironment::Demo {
1952            let http_client = self.http_client.clone();
1953            let account_id = self.core.account_id;
1954
1955            self.spawn_task("cancel_order_http", async move {
1956                let result = http_client
1957                    .cancel_order(
1958                        account_id,
1959                        product_type,
1960                        instrument_id,
1961                        Some(client_order_id),
1962                        venue_order_id,
1963                    )
1964                    .await;
1965
1966                if let Err(e) = result {
1967                    match classify_cancel_http_failure(&e) {
1968                        CommandFailure::VenueRejected(reason) => {
1969                            let ts_event = clock.get_time_ns();
1970                            emitter.emit_order_cancel_rejected_event(
1971                                strategy_id,
1972                                instrument_id,
1973                                client_order_id,
1974                                venue_order_id,
1975                                &format!("cancel-order-error: {reason}"),
1976                                ts_event,
1977                            );
1978                            anyhow::bail!("cancel order rejected: {reason}");
1979                        }
1980                        CommandFailure::NotSent(reason) => {
1981                            log::warn!(
1982                                "HTTP cancel command failed local validation for {client_order_id}: {reason}"
1983                            );
1984                        }
1985                        CommandFailure::Ambiguous(reason) => {
1986                            log::warn!(
1987                                "Ambiguous HTTP cancel failure for {client_order_id}, awaiting reconciliation: {reason}"
1988                            );
1989                        }
1990                    }
1991                }
1992
1993                Ok(())
1994            });
1995
1996            return Ok(());
1997        }
1998
1999        let raw_symbol = extract_raw_symbol(instrument_id.symbol.as_str());
2000
2001        let params = BybitWsCancelOrderParams {
2002            category: product_type,
2003            symbol: Ustr::from(raw_symbol),
2004            order_id: cmd.venue_order_id.map(|v| v.to_string()),
2005            order_link_id: Some(cmd.client_order_id.to_string()),
2006        };
2007
2008        let ws_trade = self.ws_trade.clone();
2009        let dispatch_state = Arc::clone(&self.dispatch_state);
2010
2011        self.spawn_task("cancel_order", async move {
2012            let req_id = UUID4::new().to_string();
2013            dispatch_state.pending_requests.insert(
2014                req_id.clone(),
2015                (
2016                    vec![client_order_id],
2017                    vec![venue_order_id],
2018                    PendingOperation::Cancel,
2019                ),
2020            );
2021
2022            if let Err(e) = ws_trade.cancel_order_with_id(params, req_id.clone()).await {
2023                dispatch_state.pending_requests.remove(&req_id);
2024                emitter.emit_order_cancel_rejected_event(
2025                    strategy_id,
2026                    instrument_id,
2027                    client_order_id,
2028                    venue_order_id,
2029                    &e.to_string(),
2030                    clock.get_time_ns(),
2031                );
2032            }
2033
2034            Ok(())
2035        });
2036
2037        Ok(())
2038    }
2039
2040    fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
2041        if cmd.order_side.is_some() {
2042            // Bybit cancel-all has no side parameter, so select matching open
2043            // orders and cancel their explicit IDs through the batch path.
2044            let cancels: Vec<CancelOrder> = {
2045                let cache = self.core.cache();
2046                cache
2047                    .orders_open(None, Some(&cmd.instrument_id), None, None, cmd.order_side)
2048                    .iter()
2049                    .map(|order| CancelOrder {
2050                        trader_id: order.trader_id(),
2051                        client_id: cmd.client_id,
2052                        strategy_id: order.strategy_id(),
2053                        instrument_id: order.instrument_id(),
2054                        client_order_id: order.client_order_id(),
2055                        venue_order_id: order.venue_order_id(),
2056                        command_id: cmd.command_id,
2057                        ts_init: cmd.ts_init,
2058                        params: cmd.params.clone(),
2059                        correlation_id: cmd.correlation_id,
2060                        causation_id: cmd.causation_id,
2061                    })
2062                    .collect()
2063            };
2064
2065            if cancels.is_empty() {
2066                log::debug!("No open orders to cancel for {}", cmd.instrument_id);
2067                return Ok(());
2068            }
2069
2070            return self.batch_cancel_orders(BatchCancelOrders {
2071                trader_id: cmd.trader_id,
2072                client_id: cmd.client_id,
2073                strategy_id: cmd.strategy_id,
2074                instrument_id: cmd.instrument_id,
2075                cancels,
2076                command_id: cmd.command_id,
2077                ts_init: cmd.ts_init,
2078                params: cmd.params,
2079                correlation_id: cmd.correlation_id,
2080                causation_id: cmd.causation_id,
2081            });
2082        }
2083
2084        let instrument_id = cmd.instrument_id;
2085        let product_type = self.get_product_type_for_instrument(instrument_id);
2086        let account_id = self.core.account_id;
2087        let http_client = self.http_client.clone();
2088
2089        self.spawn_task("cancel_all_orders", async move {
2090            match http_client
2091                .cancel_all_orders(account_id, product_type, instrument_id)
2092                .await
2093            {
2094                Ok(reports) => {
2095                    for report in reports {
2096                        log::debug!("Cancelled order: {report}");
2097                    }
2098                }
2099                Err(e) => {
2100                    log::error!("Failed to cancel all orders for {instrument_id}: {e}");
2101                }
2102            }
2103            Ok(())
2104        });
2105
2106        Ok(())
2107    }
2108
2109    fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
2110        if cmd.cancels.is_empty() {
2111            return Ok(());
2112        }
2113
2114        let instrument_id = cmd.instrument_id;
2115        let product_type = self.get_product_type_for_instrument(instrument_id);
2116        let emitter = self.emitter.clone();
2117        let clock = self.clock;
2118
2119        // Demo mode: cancel individually via HTTP (batch not supported)
2120        if self.config.environment == BybitEnvironment::Demo {
2121            let http_client = self.http_client.clone();
2122            let account_id = self.core.account_id;
2123            let cancels: Vec<_> = cmd
2124                .cancels
2125                .iter()
2126                .map(|c| (c.strategy_id, c.client_order_id, c.venue_order_id))
2127                .collect();
2128
2129            self.spawn_task("batch_cancel_orders_http", async move {
2130                for (strategy_id, client_order_id, venue_order_id) in cancels {
2131                    if let Err(e) = http_client
2132                        .cancel_order(
2133                            account_id,
2134                            product_type,
2135                            instrument_id,
2136                            Some(client_order_id),
2137                            venue_order_id,
2138                        )
2139                        .await
2140                    {
2141                        match classify_cancel_http_failure(&e) {
2142                            CommandFailure::VenueRejected(reason) => {
2143                                let ts_event = clock.get_time_ns();
2144                                emitter.emit_order_cancel_rejected_event(
2145                                    strategy_id,
2146                                    instrument_id,
2147                                    client_order_id,
2148                                    venue_order_id,
2149                                    &format!("cancel-order-error: {reason}"),
2150                                    ts_event,
2151                                );
2152                            }
2153                            CommandFailure::NotSent(reason) => {
2154                                log::warn!(
2155                                    "HTTP batch cancel command failed local validation for {client_order_id}: {reason}"
2156                                );
2157                            }
2158                            CommandFailure::Ambiguous(reason) => {
2159                                log::warn!(
2160                                    "Ambiguous HTTP batch cancel failure for {client_order_id}, awaiting reconciliation: {reason}"
2161                                );
2162                            }
2163                        }
2164                    }
2165                }
2166                Ok(())
2167            });
2168
2169            return Ok(());
2170        }
2171
2172        let raw_symbol = Ustr::from(extract_raw_symbol(instrument_id.symbol.as_str()));
2173
2174        let mut cancel_params = Vec::with_capacity(cmd.cancels.len());
2175        let strategy_ids: Vec<_> = cmd.cancels.iter().map(|c| c.strategy_id).collect();
2176        let client_order_ids: Vec<_> = cmd.cancels.iter().map(|c| c.client_order_id).collect();
2177        let venue_order_ids: Vec<_> = cmd.cancels.iter().map(|c| c.venue_order_id).collect();
2178
2179        for cancel in &cmd.cancels {
2180            cancel_params.push(BybitWsCancelOrderParams {
2181                category: product_type,
2182                symbol: raw_symbol,
2183                order_id: cancel.venue_order_id.map(|v| v.to_string()),
2184                order_link_id: Some(cancel.client_order_id.to_string()),
2185            });
2186        }
2187
2188        let ws_trade = self.ws_trade.clone();
2189        let dispatch_state = Arc::clone(&self.dispatch_state);
2190
2191        self.spawn_task("batch_cancel_orders", async move {
2192            let req_ids =
2193                BybitWebSocketClient::batch_request_ids(product_type, cancel_params.len());
2194            let chunk_limit = batch_send_limit(product_type);
2195
2196            for (req_id, (chunk_cids, chunk_voids)) in req_ids.iter().zip(
2197                client_order_ids
2198                    .chunks(chunk_limit)
2199                    .map(|chunk| chunk.to_vec())
2200                    .zip(
2201                        venue_order_ids
2202                            .chunks(chunk_limit)
2203                            .map(|chunk| chunk.to_vec()),
2204                    ),
2205            ) {
2206                dispatch_state.pending_requests.insert(
2207                    req_id.clone(),
2208                    (chunk_cids, chunk_voids, PendingOperation::Cancel),
2209                );
2210            }
2211
2212            if let Err(e) = ws_trade
2213                .batch_cancel_orders_with_ids(cancel_params, req_ids.clone())
2214                .await
2215            {
2216                for req_id in req_ids {
2217                    dispatch_state.pending_requests.remove(&req_id);
2218                }
2219                let reason = e.to_string();
2220                let ts_event = clock.get_time_ns();
2221
2222                for ((strategy_id, client_order_id), venue_order_id) in strategy_ids
2223                    .into_iter()
2224                    .zip(client_order_ids)
2225                    .zip(venue_order_ids)
2226                {
2227                    emitter.emit_order_cancel_rejected_event(
2228                        strategy_id,
2229                        instrument_id,
2230                        client_order_id,
2231                        venue_order_id,
2232                        &reason,
2233                        ts_event,
2234                    );
2235                }
2236            }
2237            Ok(())
2238        });
2239
2240        Ok(())
2241    }
2242}
2243
2244fn classify_cancel_http_failure(error: &anyhow::Error) -> CommandFailure {
2245    if error
2246        .chain()
2247        .any(|cause| cause.downcast_ref::<BybitCancelOrderError>().is_some())
2248    {
2249        return CommandFailure::ambiguous(error.to_string());
2250    }
2251
2252    classify_http_failure(error)
2253}
2254
2255fn classify_modify_http_failure(error: &anyhow::Error) -> CommandFailure {
2256    if error
2257        .chain()
2258        .any(|cause| cause.downcast_ref::<BybitModifyOrderError>().is_some())
2259    {
2260        return CommandFailure::ambiguous(error.to_string());
2261    }
2262
2263    classify_http_failure(error)
2264}
2265
2266fn classify_http_failure(error: &anyhow::Error) -> CommandFailure {
2267    let reason = error.to_string();
2268
2269    for cause in error.chain() {
2270        let Some(http_error) = cause.downcast_ref::<BybitHttpError>() else {
2271            continue;
2272        };
2273
2274        return match http_error {
2275            BybitHttpError::BybitError { error_code, .. }
2276                if is_bybit_ambiguous_order_error_code(i64::from(*error_code)) =>
2277            {
2278                CommandFailure::Ambiguous(reason)
2279            }
2280            BybitHttpError::BybitError { .. } => CommandFailure::VenueRejected(reason),
2281            BybitHttpError::MissingCredentials
2282            | BybitHttpError::ValidationError(_)
2283            | BybitHttpError::BuildError(_) => CommandFailure::NotSent(reason),
2284            BybitHttpError::JsonError(_)
2285            | BybitHttpError::Canceled(_)
2286            | BybitHttpError::NetworkError(_)
2287            | BybitHttpError::UnexpectedStatus { .. } => CommandFailure::Ambiguous(reason),
2288        };
2289    }
2290
2291    CommandFailure::NotSent(reason)
2292}
2293
2294impl BybitExecutionClient {
2295    fn cache_reconciliation_order_identity(&self, report: &OrderStatusReport) {
2296        let Some(client_order_id) = report.client_order_id else {
2297            return;
2298        };
2299
2300        if report.order_status.is_closed() {
2301            self.dispatch_state
2302                .order_identities
2303                .remove(&client_order_id);
2304            return;
2305        }
2306
2307        let cache = self.core.cache();
2308        let Some(order) = cache.order(&client_order_id) else {
2309            return;
2310        };
2311
2312        let identity = OrderIdentity {
2313            instrument_id: report.instrument_id,
2314            strategy_id: order.strategy_id(),
2315            order_side: order.order_side(),
2316            order_type: order.order_type(),
2317            venue_position_id: report.venue_position_id,
2318        };
2319        self.dispatch_state
2320            .order_identities
2321            .insert(client_order_id, identity);
2322    }
2323}
2324
2325#[cfg(test)]
2326mod tests {
2327    use std::{cell::RefCell, rc::Rc, time::Duration};
2328
2329    use nautilus_common::{
2330        cache::Cache,
2331        clients::ExecutionClient,
2332        messages::{
2333            ExecutionEvent,
2334            execution::{
2335                BatchCancelOrders, CancelAllOrders, CancelOrder, ModifyOrder, SubmitOrder,
2336                SubmitOrderList,
2337            },
2338        },
2339    };
2340    use nautilus_core::{Params, UUID4};
2341    use nautilus_live::ExecutionClientCore;
2342    use nautilus_model::{
2343        enums::{AccountType, OrderSide, OrderStatus},
2344        events::OrderEventAny,
2345        identifiers::{ClientOrderId, OrderListId, PositionId, StrategyId, TraderId, VenueOrderId},
2346        orders::{OrderList, builder::OrderTestBuilder, stubs::TestOrderEventStubs},
2347        types::Quantity,
2348    };
2349    use rstest::rstest;
2350
2351    use super::*;
2352    use crate::common::{
2353        consts::{BYBIT_CLIENT_ID, BYBIT_VENUE},
2354        enums::BybitMarketUnit,
2355    };
2356
2357    fn test_execution_client() -> (BybitExecutionClient, Rc<RefCell<Cache>>) {
2358        let cache = Rc::new(RefCell::new(Cache::default()));
2359        let core = ExecutionClientCore::new(
2360            TraderId::from("TESTER-001"),
2361            *BYBIT_CLIENT_ID,
2362            *BYBIT_VENUE,
2363            OmsType::Netting,
2364            AccountId::from("BYBIT-001"),
2365            AccountType::Margin,
2366            None,
2367            cache.clone(),
2368        );
2369        let config = BybitExecutionClientConfig {
2370            api_key: Some("test_key".into()),
2371            api_secret: Some("test_secret".into()),
2372            ..Default::default()
2373        };
2374
2375        (BybitExecutionClient::new(core, config).unwrap(), cache)
2376    }
2377
2378    async fn wait_for_spawned_tasks(client: &BybitExecutionClient) {
2379        for _ in 0..20 {
2380            if client.pending_tasks.all_finished() {
2381                return;
2382            }
2383
2384            tokio::time::sleep(Duration::from_millis(25)).await;
2385        }
2386
2387        panic!("timed out waiting for spawned Bybit execution tasks");
2388    }
2389
2390    fn assert_next_submitted(
2391        rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2392        client_order_id: ClientOrderId,
2393    ) {
2394        let event = rx.try_recv().expect("expected OrderSubmitted event");
2395        assert!(
2396            matches!(event, ExecutionEvent::Order(OrderEventAny::Submitted(ref submitted)) if submitted.client_order_id == client_order_id),
2397            "expected OrderSubmitted for {client_order_id}, was {event:?}",
2398        );
2399    }
2400
2401    fn assert_next_rejected(
2402        rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2403        client_order_id: ClientOrderId,
2404    ) {
2405        let event = rx.try_recv().expect("expected OrderRejected event");
2406        assert!(
2407            matches!(event, ExecutionEvent::Order(OrderEventAny::Rejected(ref rejected))
2408                if rejected.client_order_id == client_order_id),
2409            "expected OrderRejected for {client_order_id}, was {event:?}",
2410        );
2411    }
2412
2413    fn assert_next_cancel_rejected(
2414        rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2415        client_order_id: ClientOrderId,
2416    ) {
2417        let event = rx.try_recv().expect("expected OrderCancelRejected event");
2418        assert!(
2419            matches!(event, ExecutionEvent::Order(OrderEventAny::CancelRejected(ref rejected))
2420                if rejected.client_order_id == client_order_id),
2421            "expected OrderCancelRejected for {client_order_id}, was {event:?}",
2422        );
2423    }
2424
2425    fn assert_next_modify_rejected(
2426        rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2427        client_order_id: ClientOrderId,
2428    ) {
2429        let event = rx.try_recv().expect("expected OrderModifyRejected event");
2430        assert!(
2431            matches!(event, ExecutionEvent::Order(OrderEventAny::ModifyRejected(ref rejected))
2432                if rejected.client_order_id == client_order_id),
2433            "expected OrderModifyRejected for {client_order_id}, was {event:?}",
2434        );
2435    }
2436
2437    fn assert_no_order_cancel_rejected(
2438        rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2439    ) {
2440        while let Ok(event) = rx.try_recv() {
2441            assert!(
2442                !matches!(
2443                    event,
2444                    ExecutionEvent::Order(OrderEventAny::CancelRejected(_))
2445                ),
2446                "unexpected OrderCancelRejected event: {event:?}",
2447            );
2448        }
2449    }
2450
2451    fn assert_no_order_modify_rejected(
2452        rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2453    ) {
2454        while let Ok(event) = rx.try_recv() {
2455            assert!(
2456                !matches!(
2457                    event,
2458                    ExecutionEvent::Order(OrderEventAny::ModifyRejected(_))
2459                ),
2460                "unexpected OrderModifyRejected event: {event:?}",
2461            );
2462        }
2463    }
2464
2465    fn cancel_command(client_order_id: ClientOrderId) -> CancelOrder {
2466        CancelOrder::new(
2467            TraderId::from("TESTER-001"),
2468            Some(*BYBIT_CLIENT_ID),
2469            StrategyId::from("S-001"),
2470            InstrumentId::from("BTCUSDT-LINEAR.BYBIT"),
2471            client_order_id,
2472            Some(VenueOrderId::from("venue-cancel-1")),
2473            UUID4::new(),
2474            UnixNanos::default(),
2475            None,
2476            None,
2477        )
2478    }
2479
2480    fn modify_command(client_order_id: ClientOrderId, params: Option<Params>) -> ModifyOrder {
2481        ModifyOrder::new(
2482            TraderId::from("TESTER-001"),
2483            Some(*BYBIT_CLIENT_ID),
2484            StrategyId::from("S-001"),
2485            InstrumentId::from("BTCUSDT-LINEAR.BYBIT"),
2486            client_order_id,
2487            Some(VenueOrderId::from("venue-modify-1")),
2488            Some(Quantity::from("1")),
2489            Some(Price::from("10001.00")),
2490            None,
2491            UUID4::new(),
2492            UnixNanos::default(),
2493            params,
2494            None,
2495        )
2496    }
2497
2498    #[rstest]
2499    fn test_cancel_http_failure_classification_matches_policy() {
2500        let venue_reject = anyhow::Error::from(BybitHttpError::BybitError {
2501            error_code: 110001,
2502            message: "Order does not exist".to_string(),
2503        });
2504        assert_eq!(
2505            classify_cancel_http_failure(&venue_reject),
2506            CommandFailure::VenueRejected(venue_reject.to_string()),
2507        );
2508
2509        let rate_limit = anyhow::Error::from(BybitHttpError::BybitError {
2510            error_code: 10006,
2511            message: "Too many visits".to_string(),
2512        });
2513        assert_eq!(
2514            classify_cancel_http_failure(&rate_limit),
2515            CommandFailure::VenueRejected(rate_limit.to_string()),
2516        );
2517
2518        let post_lookup = anyhow::Error::from(BybitCancelOrderError::PostCancelLookup {
2519            source: anyhow::anyhow!("history lookup failed"),
2520        });
2521        assert_eq!(
2522            classify_cancel_http_failure(&post_lookup),
2523            CommandFailure::Ambiguous(post_lookup.to_string()),
2524        );
2525
2526        let transport = anyhow::Error::from(BybitHttpError::NetworkError(
2527            "connection closed".to_string(),
2528        ));
2529        assert_eq!(
2530            classify_cancel_http_failure(&transport),
2531            CommandFailure::Ambiguous(transport.to_string()),
2532        );
2533    }
2534
2535    #[rstest]
2536    fn test_modify_http_failure_classification_matches_policy() {
2537        let venue_reject = anyhow::Error::from(BybitHttpError::BybitError {
2538            error_code: 110003,
2539            message: "Order price exceeds allowable range".to_string(),
2540        });
2541        assert_eq!(
2542            classify_modify_http_failure(&venue_reject),
2543            CommandFailure::VenueRejected(venue_reject.to_string()),
2544        );
2545
2546        let server_error = anyhow::Error::from(BybitHttpError::BybitError {
2547            error_code: 10016,
2548            message: "Server error".to_string(),
2549        });
2550        assert_eq!(
2551            classify_modify_http_failure(&server_error),
2552            CommandFailure::Ambiguous(server_error.to_string()),
2553        );
2554
2555        let post_lookup = anyhow::Error::from(BybitModifyOrderError::PostModifyLookup {
2556            source: anyhow::anyhow!("realtime lookup failed"),
2557        });
2558        assert_eq!(
2559            classify_modify_http_failure(&post_lookup),
2560            CommandFailure::Ambiguous(post_lookup.to_string()),
2561        );
2562
2563        let status = anyhow::Error::from(BybitHttpError::UnexpectedStatus {
2564            status: 503,
2565            body: "service unavailable".to_string(),
2566        });
2567        assert_eq!(
2568            classify_modify_http_failure(&status),
2569            CommandFailure::Ambiguous(status.to_string()),
2570        );
2571    }
2572
2573    #[rstest]
2574    #[tokio::test]
2575    async fn test_ws_cancel_queue_failure_emits_cancel_rejected() {
2576        let (mut client, _cache) = test_execution_client();
2577        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2578        client.emitter.set_sender(tx);
2579        client.ws_trade.close().await.unwrap();
2580        let client_order_id = ClientOrderId::from("O-CANCEL-WS-FAIL");
2581
2582        client
2583            .cancel_order(cancel_command(client_order_id))
2584            .unwrap();
2585        wait_for_spawned_tasks(&client).await;
2586
2587        assert_next_cancel_rejected(&mut rx, client_order_id);
2588    }
2589
2590    #[rstest]
2591    #[tokio::test]
2592    async fn test_ws_modify_queue_failure_emits_modify_rejected() {
2593        let (mut client, _cache) = test_execution_client();
2594        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2595        client.emitter.set_sender(tx);
2596        client.ws_trade.close().await.unwrap();
2597        let client_order_id = ClientOrderId::from("O-MODIFY-WS-FAIL");
2598
2599        client
2600            .modify_order(modify_command(client_order_id, None))
2601            .unwrap();
2602        wait_for_spawned_tasks(&client).await;
2603
2604        assert_next_modify_rejected(&mut rx, client_order_id);
2605    }
2606
2607    #[rstest]
2608    #[tokio::test]
2609    async fn test_http_cancel_local_validation_failure_does_not_emit_cancel_rejected() {
2610        let (mut client, _cache) = test_execution_client();
2611        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2612        client.emitter.set_sender(tx);
2613        client.config.environment = BybitEnvironment::Demo;
2614
2615        client
2616            .cancel_order(cancel_command(ClientOrderId::from(
2617                "O-CANCEL-LOCAL-VALIDATION",
2618            )))
2619            .unwrap();
2620        wait_for_spawned_tasks(&client).await;
2621
2622        assert_no_order_cancel_rejected(&mut rx);
2623    }
2624
2625    #[rstest]
2626    #[tokio::test]
2627    async fn test_modify_local_validation_failure_does_not_emit_modify_rejected() {
2628        let (mut client, _cache) = test_execution_client();
2629        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2630        client.emitter.set_sender(tx);
2631
2632        let mut params = Params::new();
2633        params.insert("order_iv".to_string(), serde_json::json!({ "bad": true }));
2634
2635        client
2636            .modify_order(modify_command(
2637                ClientOrderId::from("O-MODIFY-LOCAL-VALIDATION"),
2638                Some(params),
2639            ))
2640            .unwrap();
2641
2642        assert_no_order_modify_rejected(&mut rx);
2643    }
2644
2645    #[rstest]
2646    #[tokio::test]
2647    async fn test_ws_submit_queue_failure_emits_order_rejected() {
2648        let (mut client, cache) = test_execution_client();
2649        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2650        client.emitter.set_sender(tx);
2651
2652        let client_order_id = ClientOrderId::from("O-WS-SEND-FAIL");
2653        let instrument_id = InstrumentId::from("BTCUSDT-LINEAR.BYBIT");
2654        let mut builder = OrderTestBuilder::new(OrderType::Limit);
2655        let order = builder
2656            .instrument_id(instrument_id)
2657            .client_order_id(client_order_id)
2658            .side(OrderSide::Buy)
2659            .quantity(Quantity::from("1"))
2660            .price(Price::from("10000.00"))
2661            .build();
2662        let init = order.init_event().clone();
2663        let trader_id = order.trader_id();
2664        let strategy_id = order.strategy_id();
2665
2666        cache
2667            .borrow_mut()
2668            .add_order(order, None, Some(*BYBIT_CLIENT_ID), false)
2669            .unwrap();
2670
2671        let command = SubmitOrder::new(
2672            trader_id,
2673            Some(*BYBIT_CLIENT_ID),
2674            strategy_id,
2675            instrument_id,
2676            client_order_id,
2677            init,
2678            None,
2679            None,
2680            None,
2681            UUID4::new(),
2682            UnixNanos::default(),
2683            None, // correlation_id
2684        );
2685
2686        client.submit_order(command).unwrap();
2687
2688        assert_next_submitted(&mut rx, client_order_id);
2689        wait_for_spawned_tasks(&client).await;
2690
2691        assert!(
2692            !client
2693                .dispatch_state
2694                .order_identities
2695                .contains_key(&client_order_id)
2696        );
2697        assert!(
2698            !client
2699                .dispatch_state
2700                .order_snapshots
2701                .contains_key(&client_order_id)
2702        );
2703        assert_next_rejected(&mut rx, client_order_id);
2704    }
2705
2706    #[rstest]
2707    #[tokio::test]
2708    async fn test_ws_submit_order_list_queue_failure_emits_order_rejected() {
2709        let (mut client, cache) = test_execution_client();
2710        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2711        client.emitter.set_sender(tx);
2712
2713        let instrument_id = InstrumentId::from("BTCUSDT-LINEAR.BYBIT");
2714        let strategy_id = StrategyId::from("S-001");
2715        let client_order_id_1 = ClientOrderId::from("O-WS-LIST-SEND-FAIL-1");
2716        let client_order_id_2 = ClientOrderId::from("O-WS-LIST-SEND-FAIL-2");
2717
2718        let mut builder_1 = OrderTestBuilder::new(OrderType::Limit);
2719        let order_1 = builder_1
2720            .strategy_id(strategy_id)
2721            .instrument_id(instrument_id)
2722            .client_order_id(client_order_id_1)
2723            .side(OrderSide::Buy)
2724            .quantity(Quantity::from("1"))
2725            .price(Price::from("10000.00"))
2726            .build();
2727        let init_1 = order_1.init_event().clone();
2728        let trader_id = order_1.trader_id();
2729
2730        let mut builder_2 = OrderTestBuilder::new(OrderType::Limit);
2731        let order_2 = builder_2
2732            .strategy_id(strategy_id)
2733            .instrument_id(instrument_id)
2734            .client_order_id(client_order_id_2)
2735            .side(OrderSide::Sell)
2736            .quantity(Quantity::from("1"))
2737            .price(Price::from("10001.00"))
2738            .build();
2739        let init_2 = order_2.init_event().clone();
2740
2741        cache
2742            .borrow_mut()
2743            .add_order(order_1, None, Some(*BYBIT_CLIENT_ID), false)
2744            .unwrap();
2745        cache
2746            .borrow_mut()
2747            .add_order(order_2, None, Some(*BYBIT_CLIENT_ID), false)
2748            .unwrap();
2749
2750        let order_list = OrderList::new(
2751            OrderListId::from("OL-WS-SEND-FAIL"),
2752            instrument_id,
2753            strategy_id,
2754            vec![client_order_id_1, client_order_id_2],
2755            UnixNanos::default(),
2756        );
2757        let command = SubmitOrderList::new(
2758            trader_id,
2759            Some(*BYBIT_CLIENT_ID),
2760            strategy_id,
2761            order_list,
2762            vec![init_1, init_2],
2763            None,
2764            None,
2765            None,
2766            UUID4::new(),
2767            UnixNanos::default(),
2768            None, // correlation_id
2769        );
2770
2771        client.submit_order_list(command).unwrap();
2772
2773        assert_next_submitted(&mut rx, client_order_id_1);
2774        assert_next_submitted(&mut rx, client_order_id_2);
2775        wait_for_spawned_tasks(&client).await;
2776
2777        for client_order_id in [client_order_id_1, client_order_id_2] {
2778            assert!(
2779                !client
2780                    .dispatch_state
2781                    .order_identities
2782                    .contains_key(&client_order_id)
2783            );
2784            assert!(
2785                !client
2786                    .dispatch_state
2787                    .order_snapshots
2788                    .contains_key(&client_order_id)
2789            );
2790            assert_next_rejected(&mut rx, client_order_id);
2791        }
2792    }
2793
2794    fn sample_order_status_report(
2795        client_order_id: ClientOrderId,
2796        instrument_id: InstrumentId,
2797        order_status: OrderStatus,
2798        venue_position_id: Option<PositionId>,
2799    ) -> OrderStatusReport {
2800        let mut report = OrderStatusReport::new(
2801            AccountId::from("BYBIT-001"),
2802            instrument_id,
2803            Some(client_order_id),
2804            VenueOrderId::from("BYBIT-ORDER-001"),
2805            OrderSide::Buy.into(),
2806            OrderType::Limit,
2807            TimeInForce::Gtc,
2808            order_status,
2809            Quantity::from("1"),
2810            Quantity::from("0"),
2811            UnixNanos::default(),
2812            UnixNanos::default(),
2813            UnixNanos::default(),
2814            None,
2815        );
2816        report.venue_position_id = venue_position_id;
2817        report
2818    }
2819
2820    #[rstest]
2821    fn test_cache_reconciliation_order_identity_caches_and_clears_hedge_report() {
2822        let (client, cache) = test_execution_client();
2823        let client_order_id = ClientOrderId::from("O-HEDGE-RECON");
2824        let instrument_id = InstrumentId::from("BTCUSDT-LINEAR.BYBIT");
2825        let venue_position_id = PositionId::from("BTCUSDT-LINEAR.BYBIT-LONG");
2826        let mut builder = OrderTestBuilder::new(OrderType::Limit);
2827        let order = builder
2828            .instrument_id(instrument_id)
2829            .client_order_id(client_order_id)
2830            .side(OrderSide::Buy)
2831            .quantity(Quantity::from("1"))
2832            .price(Price::from("10000.00"))
2833            .build();
2834        cache
2835            .borrow_mut()
2836            .add_order(order.clone(), None, None, false)
2837            .unwrap();
2838
2839        let report = sample_order_status_report(
2840            client_order_id,
2841            instrument_id,
2842            OrderStatus::Accepted,
2843            Some(venue_position_id),
2844        );
2845        client.cache_reconciliation_order_identity(&report);
2846
2847        {
2848            let identity = client
2849                .dispatch_state
2850                .order_identities
2851                .get(&client_order_id)
2852                .unwrap();
2853            assert_eq!(identity.instrument_id, instrument_id);
2854            assert_eq!(identity.strategy_id, order.strategy_id());
2855            assert_eq!(identity.order_side, order.order_side());
2856            assert_eq!(identity.order_type, order.order_type());
2857            assert_eq!(identity.venue_position_id, Some(venue_position_id));
2858        }
2859
2860        let terminal_report = sample_order_status_report(
2861            client_order_id,
2862            instrument_id,
2863            OrderStatus::Filled,
2864            Some(venue_position_id),
2865        );
2866        client.cache_reconciliation_order_identity(&terminal_report);
2867
2868        assert!(
2869            client
2870                .dispatch_state
2871                .order_identities
2872                .get(&client_order_id)
2873                .is_none()
2874        );
2875    }
2876
2877    #[rstest]
2878    fn test_cache_reconciliation_order_identity_keeps_one_way_local_report() {
2879        let (client, cache) = test_execution_client();
2880        let client_order_id = ClientOrderId::from("O-ONEWAY-RECON");
2881        let instrument_id = InstrumentId::from("BTCUSDT-LINEAR.BYBIT");
2882        let mut builder = OrderTestBuilder::new(OrderType::Limit);
2883        let order = builder
2884            .instrument_id(instrument_id)
2885            .client_order_id(client_order_id)
2886            .side(OrderSide::Buy)
2887            .quantity(Quantity::from("1"))
2888            .price(Price::from("10000.00"))
2889            .build();
2890        cache
2891            .borrow_mut()
2892            .add_order(order, None, None, false)
2893            .unwrap();
2894
2895        let report =
2896            sample_order_status_report(client_order_id, instrument_id, OrderStatus::Accepted, None);
2897        client.cache_reconciliation_order_identity(&report);
2898
2899        let identity = client
2900            .dispatch_state
2901            .order_identities
2902            .get(&client_order_id)
2903            .unwrap();
2904        assert_eq!(identity.instrument_id, instrument_id);
2905        assert_eq!(identity.venue_position_id, None);
2906    }
2907
2908    #[rstest]
2909    #[case::spot_market_base(
2910        BybitProductType::Spot,
2911        BybitOrderType::Market,
2912        false,
2913        Some(BybitMarketUnit::BaseCoin)
2914    )]
2915    #[case::spot_market_quote(
2916        BybitProductType::Spot,
2917        BybitOrderType::Market,
2918        true,
2919        Some(BybitMarketUnit::QuoteCoin)
2920    )]
2921    #[case::spot_limit(BybitProductType::Spot, BybitOrderType::Limit, true, None)]
2922    #[case::linear_market(BybitProductType::Linear, BybitOrderType::Market, true, None)]
2923    fn test_ws_params_market_unit(
2924        #[case] product_type: BybitProductType,
2925        #[case] order_type: BybitOrderType,
2926        #[case] is_quote_quantity: bool,
2927        #[case] expected: Option<BybitMarketUnit>,
2928    ) {
2929        let params = BybitWsPlaceOrderParams {
2930            category: product_type,
2931            symbol: ustr::Ustr::from("BTCUSDT"),
2932            side: BybitOrderSide::Buy,
2933            order_type,
2934            qty: "1.0".to_string(),
2935            is_leverage: None,
2936            market_unit: spot_market_unit(product_type, order_type, is_quote_quantity),
2937            price: None,
2938            time_in_force: None,
2939            order_link_id: None,
2940            reduce_only: None,
2941            close_on_trigger: None,
2942            trigger_price: None,
2943            trigger_by: None,
2944            trigger_direction: None,
2945            tpsl_mode: None,
2946            take_profit: None,
2947            stop_loss: None,
2948            tp_trigger_by: None,
2949            sl_trigger_by: None,
2950            sl_trigger_price: None,
2951            tp_trigger_price: None,
2952            sl_order_type: None,
2953            tp_order_type: None,
2954            sl_limit_price: None,
2955            tp_limit_price: None,
2956            order_iv: None,
2957            smp_type: None,
2958            mmp: None,
2959            position_idx: None,
2960            bbo_side_type: None,
2961            bbo_level: None,
2962        };
2963
2964        assert_eq!(params.market_unit, expected);
2965    }
2966
2967    #[rstest]
2968    #[case::market(OrderType::Market, BybitOrderType::Market, false)]
2969    #[case::limit(OrderType::Limit, BybitOrderType::Limit, false)]
2970    #[case::stop_market(OrderType::StopMarket, BybitOrderType::Market, true)]
2971    #[case::stop_limit(OrderType::StopLimit, BybitOrderType::Limit, true)]
2972    #[case::market_if_touched(OrderType::MarketIfTouched, BybitOrderType::Market, true)]
2973    #[case::limit_if_touched(OrderType::LimitIfTouched, BybitOrderType::Limit, true)]
2974    fn test_map_order_type(
2975        #[case] input: OrderType,
2976        #[case] expected_type: BybitOrderType,
2977        #[case] expected_conditional: bool,
2978    ) {
2979        let (bybit_type, is_conditional) = BybitExecutionClient::map_order_type(input).unwrap();
2980        assert_eq!(bybit_type, expected_type);
2981        assert_eq!(is_conditional, expected_conditional);
2982    }
2983
2984    #[rstest]
2985    fn test_map_order_type_rejects_trailing_stop() {
2986        BybitExecutionClient::map_order_type(OrderType::TrailingStopMarket).unwrap_err();
2987    }
2988
2989    #[rstest]
2990    #[case::linear("BTCUSDT-LINEAR", true)]
2991    #[case::inverse("BTCUSD-INVERSE", true)]
2992    #[case::spot("BTCUSDT-SPOT", false)]
2993    #[case::option("BTC-30JUN25-100000-C-OPTION", false)]
2994    fn test_parse_derivative_symbol_filters_product_type(
2995        #[case] symbol_str: &str,
2996        #[case] keeps: bool,
2997    ) {
2998        let result = BybitExecutionClient::parse_derivative_symbol(symbol_str);
2999        assert_eq!(result.is_some(), keeps);
3000    }
3001
3002    #[rstest]
3003    fn test_parse_derivative_symbol_rejects_malformed() {
3004        assert!(BybitExecutionClient::parse_derivative_symbol("not-a-real-symbol").is_none());
3005    }
3006
3007    #[rstest]
3008    #[case::matches_msg("Position mode has not been modified", "110025", true)]
3009    #[case::matches_code("retCode 110025: noop", "110025", true)]
3010    #[case::matches_msg_only("Already not been modified", "", true)]
3011    #[case::wrong_code("retCode 99999: other", "110025", false)]
3012    #[case::empty_no_modified_msg("retCode 99999", "", false)]
3013    fn test_is_unchanged_error(#[case] msg: &str, #[case] code: &str, #[case] expected: bool) {
3014        let err = anyhow::anyhow!("{msg}");
3015        assert_eq!(
3016            BybitExecutionClient::is_unchanged_error(&err, code),
3017            expected
3018        );
3019    }
3020
3021    #[rstest]
3022    #[case::matches("Margin needs to be equal to or greater than 0.5", true)]
3023    #[case::no_match("Some other error", false)]
3024    fn test_is_low_margin_error(#[case] msg: &str, #[case] expected: bool) {
3025        let err = anyhow::anyhow!("{msg}");
3026        assert_eq!(BybitExecutionClient::is_low_margin_error(&err), expected);
3027    }
3028
3029    #[rstest]
3030    fn test_submit_rejection_reason_matches_confirmed_rejection() {
3031        let err = anyhow::Error::from(BybitSubmitOrderError::Rejected {
3032            reason: "EC_PostOnlyWillTakeLiquidity".to_string(),
3033        });
3034
3035        assert_eq!(
3036            submit_rejection_reason(&err),
3037            Some("EC_PostOnlyWillTakeLiquidity"),
3038        );
3039    }
3040
3041    #[rstest]
3042    fn test_submit_rejection_reason_ignores_post_submit_lookup_failure() {
3043        let err = anyhow::Error::from(BybitSubmitOrderError::PostSubmitLookup {
3044            source: anyhow::Error::from(BybitHttpError::BybitError {
3045                error_code: 110017,
3046                message: "current position is zero, cannot fix reduce-only order qty".to_string(),
3047            }),
3048        })
3049        .context("Submit order failed");
3050
3051        assert_eq!(submit_rejection_reason(&err), None);
3052    }
3053
3054    #[rstest]
3055    fn test_submit_rejection_reason_ignores_missing_order_id() {
3056        let err = anyhow::Error::from(BybitSubmitOrderError::MissingOrderId);
3057
3058        assert_eq!(submit_rejection_reason(&err), None);
3059    }
3060
3061    #[rstest]
3062    fn test_submit_rejection_reason_matches_venue_http_error() {
3063        let err = anyhow::Error::from(BybitHttpError::BybitError {
3064            error_code: 110017,
3065            message: "current position is zero, cannot fix reduce-only order qty".to_string(),
3066        });
3067
3068        assert_eq!(
3069            submit_rejection_reason(&err),
3070            Some("current position is zero, cannot fix reduce-only order qty"),
3071        );
3072    }
3073
3074    #[rstest]
3075    fn test_submit_rejection_reason_ignores_ambiguous_http_error() {
3076        let err = anyhow::Error::from(BybitHttpError::BybitError {
3077            error_code: 10016,
3078            message: "rate limit exceeded".to_string(),
3079        });
3080
3081        assert_eq!(submit_rejection_reason(&err), None);
3082    }
3083
3084    #[rstest]
3085    #[case::buy(Some(OrderSide::Buy), vec![0, 2])]
3086    #[case::sell(Some(OrderSide::Sell), vec![1, 3])]
3087    #[tokio::test]
3088    async fn test_cancel_all_orders_filters_by_side_and_preserves_owners(
3089        #[case] order_side: Option<OrderSide>,
3090        #[case] expected_indices: Vec<usize>,
3091    ) {
3092        let (mut client, cache) = test_execution_client();
3093        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3094        client.emitter.set_sender(tx);
3095        client.ws_trade.close().await.unwrap();
3096
3097        let instrument_id = InstrumentId::from("BTCUSDT-LINEAR.BYBIT");
3098        let other_instrument_id = InstrumentId::from("ETHUSDT-LINEAR.BYBIT");
3099        let account_id = client.core.account_id;
3100        let mut orders = Vec::new();
3101
3102        for (index, (instrument, owner, side, open)) in [
3103            (instrument_id, "S-001", OrderSide::Buy, true),
3104            (instrument_id, "S-001", OrderSide::Sell, true),
3105            (instrument_id, "S-002", OrderSide::Buy, true),
3106            (instrument_id, "S-002", OrderSide::Sell, true),
3107            (other_instrument_id, "S-001", OrderSide::Buy, true),
3108            (other_instrument_id, "S-002", OrderSide::Sell, true),
3109            (instrument_id, "S-001", OrderSide::Buy, false),
3110            (instrument_id, "S-002", OrderSide::Sell, false),
3111        ]
3112        .into_iter()
3113        .enumerate()
3114        {
3115            let client_order_id = ClientOrderId::from(format!("O-CANCEL-ALL-{index}"));
3116            let mut builder = OrderTestBuilder::new(OrderType::Limit);
3117            let order = builder
3118                .instrument_id(instrument)
3119                .client_order_id(client_order_id)
3120                .strategy_id(StrategyId::from(owner))
3121                .side(side)
3122                .quantity(Quantity::from("1"))
3123                .price(Price::from("10000.00"))
3124                .build();
3125            let venue_order_id = VenueOrderId::from(format!("BYBIT-{index}"));
3126            let accepted = TestOrderEventStubs::accepted(&order, account_id, venue_order_id);
3127            cache
3128                .borrow_mut()
3129                .add_order(order.clone(), None, Some(*BYBIT_CLIENT_ID), false)
3130                .unwrap();
3131            let order = cache.borrow_mut().update_order(&accepted).unwrap();
3132
3133            if !open {
3134                let canceled =
3135                    TestOrderEventStubs::canceled(&order, account_id, Some(venue_order_id));
3136                cache.borrow_mut().update_order(&canceled).unwrap();
3137            }
3138
3139            orders.push((order, venue_order_id));
3140        }
3141
3142        client
3143            .cancel_all_orders(CancelAllOrders::new(
3144                TraderId::from("TESTER-001"),
3145                Some(*BYBIT_CLIENT_ID),
3146                StrategyId::from("S-001"),
3147                instrument_id,
3148                order_side,
3149                UUID4::new(),
3150                UnixNanos::default(),
3151                None,
3152                None,
3153            ))
3154            .unwrap();
3155        wait_for_spawned_tasks(&client).await;
3156
3157        // The closed transport rejects the batch, exposing selection and
3158        // event ownership through the existing failure path.
3159        let mut actual = Vec::new();
3160
3161        for _ in &expected_indices {
3162            let event = rx.try_recv().expect("expected OrderCancelRejected event");
3163            match event {
3164                ExecutionEvent::Order(OrderEventAny::CancelRejected(rejected)) => {
3165                    actual.push((rejected.client_order_id, rejected.strategy_id));
3166                }
3167                event => panic!("expected OrderCancelRejected, was {event:?}"),
3168            }
3169        }
3170
3171        let mut expected: Vec<_> = expected_indices
3172            .iter()
3173            .map(|&index| {
3174                let (order, _) = &orders[index];
3175                (order.client_order_id(), order.strategy_id())
3176            })
3177            .collect();
3178
3179        actual.sort();
3180        expected.sort();
3181
3182        assert_eq!(actual, expected);
3183        assert!(rx.try_recv().is_err());
3184    }
3185
3186    #[rstest]
3187    #[tokio::test]
3188    async fn test_cancel_all_orders_with_side_and_empty_cache_sends_nothing() {
3189        let (mut client, _cache) = test_execution_client();
3190        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3191        client.emitter.set_sender(tx);
3192        client.ws_trade.close().await.unwrap();
3193
3194        client
3195            .cancel_all_orders(CancelAllOrders::new(
3196                TraderId::from("TESTER-001"),
3197                Some(*BYBIT_CLIENT_ID),
3198                StrategyId::from("S-001"),
3199                InstrumentId::from("BTCUSDT-LINEAR.BYBIT"),
3200                Some(OrderSide::Buy),
3201                UUID4::new(),
3202                UnixNanos::default(),
3203                None,
3204                None,
3205            ))
3206            .unwrap();
3207        wait_for_spawned_tasks(&client).await;
3208
3209        assert!(rx.try_recv().is_err());
3210        assert!(client.dispatch_state.pending_requests.is_empty());
3211    }
3212
3213    #[rstest]
3214    #[tokio::test]
3215    async fn test_batch_cancel_orders_failure_preserves_owning_strategies() {
3216        let (mut client, _cache) = test_execution_client();
3217        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3218        client.emitter.set_sender(tx);
3219        client.ws_trade.close().await.unwrap();
3220
3221        let instrument_id = InstrumentId::from("BTCUSDT-LINEAR.BYBIT");
3222
3223        let cancels = ["S-001", "S-002"]
3224            .into_iter()
3225            .enumerate()
3226            .map(|(index, owner)| {
3227                CancelOrder::new(
3228                    TraderId::from("TESTER-001"),
3229                    Some(*BYBIT_CLIENT_ID),
3230                    StrategyId::from(owner),
3231                    instrument_id,
3232                    ClientOrderId::from(format!("O-BATCH-OWNER-{index}")),
3233                    Some(VenueOrderId::from(format!("BYBIT-BATCH-{index}"))),
3234                    UUID4::new(),
3235                    UnixNanos::default(),
3236                    None,
3237                    None,
3238                )
3239            })
3240            .collect::<Vec<_>>();
3241
3242        client
3243            .batch_cancel_orders(BatchCancelOrders::new(
3244                TraderId::from("TESTER-001"),
3245                Some(*BYBIT_CLIENT_ID),
3246                StrategyId::from("S-001"),
3247                instrument_id,
3248                cancels.clone(),
3249                UUID4::new(),
3250                UnixNanos::default(),
3251                None,
3252                None,
3253            ))
3254            .unwrap();
3255        wait_for_spawned_tasks(&client).await;
3256
3257        let mut actual = Vec::new();
3258
3259        for _ in &cancels {
3260            let event = rx.try_recv().expect("expected OrderCancelRejected event");
3261            match event {
3262                ExecutionEvent::Order(OrderEventAny::CancelRejected(rejected)) => {
3263                    actual.push((rejected.client_order_id, rejected.strategy_id));
3264                }
3265                event => panic!("expected OrderCancelRejected, was {event:?}"),
3266            }
3267        }
3268
3269        let mut expected: Vec<_> = cancels
3270            .iter()
3271            .map(|cancel| (cancel.client_order_id, cancel.strategy_id))
3272            .collect();
3273
3274        actual.sort();
3275        expected.sort();
3276
3277        assert_eq!(actual, expected);
3278        assert!(rx.try_recv().is_err());
3279    }
3280}