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