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 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#[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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, );
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, );
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 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}