1use std::{
19 future::Future,
20 sync::{
21 Arc,
22 atomic::{AtomicBool, Ordering},
23 },
24 time::{Duration, Instant},
25};
26
27use ahash::AHashMap;
28use anyhow::Context;
29use async_trait::async_trait;
30use futures_util::{StreamExt, pin_mut};
31use nautilus_common::{
32 clients::ExecutionClient,
33 enums::LogLevel,
34 live::{get_runtime, runner::get_exec_event_sender, task::TaskHandles},
35 messages::execution::{
36 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
37 GenerateFillReportsBuilder, GenerateOrderStatusReport, GenerateOrderStatusReports,
38 GenerateOrderStatusReportsBuilder, GeneratePositionStatusReports,
39 GeneratePositionStatusReportsBuilder, ModifyOrder, QueryAccount, QueryOrder, SubmitOrder,
40 SubmitOrderList,
41 },
42};
43use nautilus_core::{
44 UnixNanos,
45 time::{AtomicTime, get_atomic_clock_realtime},
46};
47use nautilus_live::{ExecutionClientCore, ExecutionEventEmitter};
48use nautilus_model::{
49 accounts::AccountAny,
50 enums::{AccountType, OmsType, OrderSide, OrderType, TrailingOffsetType},
51 identifiers::{
52 AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, Venue, VenueOrderId,
53 },
54 instruments::{Instrument, InstrumentAny},
55 orders::{Order, OrderAny},
56 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
57 types::{AccountBalance, MarginBalance},
58};
59use rust_decimal::prelude::ToPrimitive;
60use tokio::task::JoinHandle;
61use ustr::Ustr;
62
63use crate::{
64 broadcast::{
65 canceller::{CancelBroadcaster, CancelBroadcasterConfig},
66 submitter::{DEFINITIVE_SUBMIT_REJECTION, SubmitBroadcaster, SubmitBroadcasterConfig},
67 },
68 common::{
69 enums::{BitmexContingencyType, BitmexOrderType, BitmexPegPriceType, BitmexTimeInForce},
70 parse::{parse_peg_offset_value, parse_peg_price_type},
71 },
72 config::BitmexExecClientConfig,
73 http::{client::BitmexHttpClient, error::BitmexHttpError},
74 websocket::{
75 client::BitmexWebSocketClient,
76 dispatch::{self, OrderIdentity, WsDispatchState},
77 },
78};
79
80#[derive(Debug)]
81pub struct BitmexExecutionClient {
82 core: ExecutionClientCore,
83 clock: &'static AtomicTime,
84 config: BitmexExecClientConfig,
85 emitter: ExecutionEventEmitter,
86 http_client: BitmexHttpClient,
87 ws_client: BitmexWebSocketClient,
88 ws_dispatch_state: Arc<WsDispatchState>,
89 _submitter: SubmitBroadcaster,
90 _canceller: CancelBroadcaster,
91 ws_stream_handle: Option<JoinHandle<()>>,
92 pending_tasks: TaskHandles,
93 dms_task_handle: Option<JoinHandle<()>>,
94 dms_running: Arc<AtomicBool>,
95}
96
97impl BitmexExecutionClient {
98 fn log_report_receipt(count: usize, report_type: &str, log_level: LogLevel) {
99 let plural = if count == 1 { "" } else { "s" };
100 let message = format!("Received {count} {report_type}{plural}");
101
102 match log_level {
103 LogLevel::Off => {}
104 LogLevel::Trace => log::trace!("{message}"),
105 LogLevel::Debug => log::debug!("{message}"),
106 LogLevel::Info => log::info!("{message}"),
107 LogLevel::Warning => log::warn!("{message}"),
108 LogLevel::Error => log::error!("{message}"),
109 }
110 }
111
112 pub fn new(
118 mut core: ExecutionClientCore,
119 config: BitmexExecClientConfig,
120 ) -> anyhow::Result<Self> {
121 if !config.has_api_credentials() {
122 anyhow::bail!("BitMEX execution client requires API key and secret");
123 }
124
125 if let Some(account_id) = config.account_id {
126 core.set_account_id(account_id);
127 }
128
129 let trader_id = core.trader_id;
130 let account_id = core.account_id;
131 let clock = get_atomic_clock_realtime();
132 let emitter =
133 ExecutionEventEmitter::new(clock, trader_id, account_id, AccountType::Margin, None);
134 let http_client = BitmexHttpClient::new(
135 Some(config.http_base_url()),
136 config.api_key.clone(),
137 config.api_secret.clone(),
138 config.environment,
139 config.http_timeout_secs,
140 config.max_retries,
141 config.retry_delay_initial_ms,
142 config.retry_delay_max_ms,
143 config.recv_window_ms,
144 config.max_requests_per_second,
145 config.max_requests_per_minute,
146 config.proxy_url.clone(),
147 )
148 .context("failed to construct BitMEX HTTP client")?;
149 let ws_client = BitmexWebSocketClient::new_with_env(
150 Some(config.ws_url()),
151 config.api_key.clone(),
152 config.api_secret.clone(),
153 Some(account_id),
154 config.heartbeat_interval_secs,
155 config.auth_timeout_secs,
156 config.environment,
157 config.transport_backend,
158 config.proxy_url.clone(),
159 )
160 .context("failed to construct BitMEX execution websocket client")?;
161
162 let pool_size = config.submitter_pool_size.unwrap_or(1);
163 let submitter_proxy_urls = match &config.submitter_proxy_urls {
164 Some(urls) => urls.iter().map(|url| Some(url.clone())).collect(),
165 None => vec![config.proxy_url.clone(); pool_size],
166 };
167
168 let submitter_config = SubmitBroadcasterConfig {
169 pool_size,
170 api_key: config.api_key.clone(),
171 api_secret: config.api_secret.clone(),
172 base_url: config.base_url_http.clone(),
173 environment: config.environment,
174 timeout_secs: config.http_timeout_secs,
175 max_retries: config.max_retries,
176 retry_delay_ms: config.retry_delay_initial_ms,
177 retry_delay_max_ms: config.retry_delay_max_ms,
178 recv_window_ms: config.recv_window_ms,
179 max_requests_per_second: config.max_requests_per_second,
180 max_requests_per_minute: config.max_requests_per_minute,
181 proxy_urls: submitter_proxy_urls,
182 ..Default::default()
183 };
184
185 let _submitter = SubmitBroadcaster::new(submitter_config)
186 .context("failed to create SubmitBroadcaster")?;
187
188 let canceller_pool_size = config.canceller_pool_size.unwrap_or(1);
189 let canceller_proxy_urls = match &config.canceller_proxy_urls {
190 Some(urls) => urls.iter().map(|url| Some(url.clone())).collect(),
191 None => vec![config.proxy_url.clone(); canceller_pool_size],
192 };
193
194 let canceller_config = CancelBroadcasterConfig {
195 pool_size: canceller_pool_size,
196 api_key: config.api_key.clone(),
197 api_secret: config.api_secret.clone(),
198 base_url: config.base_url_http.clone(),
199 environment: config.environment,
200 timeout_secs: config.http_timeout_secs,
201 max_retries: config.max_retries,
202 retry_delay_ms: config.retry_delay_initial_ms,
203 retry_delay_max_ms: config.retry_delay_max_ms,
204 recv_window_ms: config.recv_window_ms,
205 max_requests_per_second: config.max_requests_per_second,
206 max_requests_per_minute: config.max_requests_per_minute,
207 proxy_urls: canceller_proxy_urls,
208 ..Default::default()
209 };
210
211 let _canceller = CancelBroadcaster::new(canceller_config)
212 .context("failed to create CancelBroadcaster")?;
213
214 Ok(Self {
215 core,
216 clock,
217 config,
218 emitter,
219 http_client,
220 ws_client,
221 ws_dispatch_state: Arc::new(WsDispatchState::default()),
222 _submitter,
223 _canceller,
224 ws_stream_handle: None,
225 pending_tasks: TaskHandles::default(),
226 dms_task_handle: None,
227 dms_running: Arc::new(AtomicBool::new(false)),
228 })
229 }
230
231 fn spawn_task<F>(&self, label: &'static str, fut: F)
232 where
233 F: Future<Output = anyhow::Result<()>> + Send + 'static,
234 {
235 let handle = get_runtime().spawn(async move {
236 if let Err(e) = fut.await {
237 log::error!("{label}: {e:?}");
238 }
239 });
240
241 self.pending_tasks.push(handle);
242 }
243
244 fn abort_pending_tasks(&self) {
245 self.pending_tasks.abort_all();
246 }
247
248 fn ensure_order_identity(
253 &self,
254 client_order_id: ClientOrderId,
255 strategy_id: StrategyId,
256 instrument_id: InstrumentId,
257 ) {
258 if self
259 .ws_dispatch_state
260 .order_identities
261 .contains_key(&client_order_id)
262 {
263 return;
264 }
265
266 let cache = self.core.cache();
267 let (order_side, order_type) = cache
268 .order(&client_order_id)
269 .map_or((OrderSide::NoOrderSide, OrderType::Market), |o| {
270 (o.order_side(), o.order_type())
271 });
272 drop(cache);
273
274 self.ws_dispatch_state.order_identities.insert(
275 client_order_id,
276 OrderIdentity {
277 instrument_id,
278 strategy_id,
279 order_side,
280 order_type,
281 },
282 );
283 self.ws_dispatch_state.insert_accepted(client_order_id);
284 }
285
286 fn start_deadmans_switch(&mut self) {
287 let Some(timeout_secs) = self.config.deadmans_switch_timeout_secs else {
288 return;
289 };
290
291 let timeout_ms = timeout_secs * 1000;
292 let interval_secs = (timeout_secs / 4).max(1);
293
294 log::info!(
295 "Starting dead man's switch: timeout={timeout_secs}s, refresh_interval={interval_secs}s",
296 );
297
298 self.dms_running.store(true, Ordering::SeqCst);
299 let running = self.dms_running.clone();
300 let http_client = self.http_client.clone();
301
302 let handle = get_runtime().spawn(async move {
303 while running.load(Ordering::SeqCst) {
304 if let Err(e) = http_client.cancel_all_after(timeout_ms).await {
305 log::warn!("Dead man's switch heartbeat failed: {e}");
306 }
307 tokio::time::sleep(Duration::from_secs(interval_secs)).await;
308 }
309 });
310
311 self.dms_task_handle = Some(handle);
312 }
313
314 async fn stop_deadmans_switch(&mut self) {
315 if self.config.deadmans_switch_timeout_secs.is_none() {
316 return;
317 }
318
319 self.dms_running.store(false, Ordering::SeqCst);
320
321 if let Some(handle) = self.dms_task_handle.take() {
323 handle.abort();
324 let _ = handle.await;
325 }
326
327 log::info!("Disarming dead man's switch");
328
329 if let Err(e) = self.http_client.cancel_all_after(0).await {
330 log::warn!("Failed to disarm dead man's switch: {e}");
331 }
332 }
333
334 async fn ensure_instruments_initialized_async(&self) -> anyhow::Result<()> {
335 if self.core.instruments_initialized() {
336 return Ok(());
337 }
338
339 let mut instruments: Vec<InstrumentAny> = {
340 let cache = self.core.cache();
341 cache
342 .instruments(&self.core.venue, None)
343 .into_iter()
344 .cloned()
345 .collect()
346 };
347
348 if instruments.is_empty() {
349 let http = self.http_client.clone();
350 instruments = http
351 .request_instruments(self.config.active_only)
352 .await
353 .context("failed to request BitMEX instruments")?;
354 } else {
355 log::debug!(
356 "Reusing {} cached BitMEX instruments for execution client initialization",
357 instruments.len()
358 );
359 }
360
361 instruments.sort_by_key(|instrument| instrument.id());
362
363 self.http_client.cache_instruments(&instruments);
364 self.ws_client.cache_instruments(&instruments);
365 for instrument in &instruments {
366 self._submitter.cache_instrument(instrument);
367 self._canceller.cache_instrument(instrument);
368 }
369
370 self.core.set_instruments_initialized();
371 Ok(())
372 }
373
374 async fn refresh_account_state(&mut self) -> anyhow::Result<()> {
375 let account_state = self
376 .http_client
377 .request_account_state(self.core.account_id)
378 .await
379 .context("failed to request BitMEX account state")?;
380
381 self.apply_account_id(account_state.account_id);
382 self.emitter.send_account_state(account_state);
383 Ok(())
384 }
385
386 fn apply_account_id(&mut self, account_id: AccountId) {
387 if self.core.account_id != account_id {
388 log::debug!(
389 "Discovered BitMEX account ID: account_id={} (was {})",
390 account_id,
391 self.core.account_id
392 );
393 }
394
395 self.core.set_account_id(account_id);
396 self.emitter.set_account_id(account_id);
397 self.ws_client.set_account_id(account_id);
398 }
399
400 async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
401 let account_id = self.core.account_id;
402
403 if self.core.cache().account(&account_id).is_some() {
404 log::info!("Account {account_id} registered");
405 return Ok(());
406 }
407
408 let start = Instant::now();
409 let timeout = Duration::from_secs_f64(timeout_secs);
410 let interval = Duration::from_millis(10);
411
412 loop {
413 tokio::time::sleep(interval).await;
414
415 if self.core.cache().account(&account_id).is_some() {
416 log::info!("Account {account_id} registered");
417 return Ok(());
418 }
419
420 if start.elapsed() >= timeout {
421 anyhow::bail!(
422 "Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
423 );
424 }
425 }
426 }
427
428 fn start_ws_stream(&mut self) {
429 if self.ws_stream_handle.is_some() {
430 return;
431 }
432
433 let stream = self.ws_client.stream();
434 let emitter = self.emitter.clone();
435 let state = Arc::clone(&self.ws_dispatch_state);
436 state.order_rows_clear();
437 let account_id = self.core.account_id;
438 let clock = self.clock;
439
440 let mut instruments_by_symbol: AHashMap<Ustr, InstrumentAny> = self
442 .core
443 .cache()
444 .instruments(&self.core.venue, None)
445 .into_iter()
446 .map(|inst| (inst.symbol().inner(), inst.clone()))
447 .collect();
448
449 if instruments_by_symbol.is_empty() {
450 for (key, inst) in self.http_client.instruments_cache.load().iter() {
451 instruments_by_symbol.insert(*key, inst.clone());
452 }
453 }
454
455 let handle = get_runtime().spawn(async move {
456 pin_mut!(stream);
457 let mut order_type_cache: AHashMap<ClientOrderId, OrderType> = AHashMap::new();
458 let mut order_symbol_cache: AHashMap<ClientOrderId, Ustr> = AHashMap::new();
459 let mut insts_by_symbol = instruments_by_symbol;
460
461 while let Some(message) = stream.next().await {
462 dispatch::dispatch_ws_message(
463 clock.get_time_ns(),
464 message,
465 &emitter,
466 &state,
467 &mut insts_by_symbol,
468 &mut order_type_cache,
469 &mut order_symbol_cache,
470 account_id,
471 );
472 }
473 });
474
475 self.ws_stream_handle = Some(handle);
476 }
477
478 fn submit_cached_order(
479 &self,
480 order: &OrderAny,
481 submit_tries: Option<usize>,
482 peg_price_type: Option<BitmexPegPriceType>,
483 peg_offset_value: Option<f64>,
484 task_label: &'static str,
485 ) {
486 if order.is_closed() {
487 log::warn!("Cannot submit closed order {}", order.client_order_id());
488 return;
489 }
490
491 if let Err(e) = validate_order_for_bitmex_submit(order, peg_price_type, peg_offset_value) {
492 self.emitter.emit_order_denied(order, &e.to_string());
493 return;
494 }
495
496 self.emitter.emit_order_submitted(order);
497
498 let strategy_id = order.strategy_id();
499 let instrument_id = order.instrument_id();
500 let client_order_id = order.client_order_id();
501 let order_side = order.order_side();
502 let order_type = order.order_type();
503
504 self.ws_dispatch_state.order_identities.insert(
505 client_order_id,
506 OrderIdentity {
507 instrument_id,
508 strategy_id,
509 order_side,
510 order_type,
511 },
512 );
513
514 let use_broadcaster = submit_tries.is_some_and(|n| n > 1);
515 let http_client = self.http_client.clone();
516 let submitter = self._submitter.clone_for_async();
517 let ws_dispatch_state = self.ws_dispatch_state.clone();
518 let emitter = self.emitter.clone();
519 let clock = self.clock;
520 let quantity = order.quantity();
521 let time_in_force = order.time_in_force();
522 let price = order.price();
523 let trigger_price = order.trigger_price();
524 let trigger_type = order.trigger_type();
525 let trailing_offset = order.trailing_offset().and_then(|d| d.to_f64());
526 let trailing_offset_type = order.trailing_offset_type();
527 let display_qty = order.display_qty();
528 let post_only = order.is_post_only();
529 let reduce_only = order.is_reduce_only();
530 let order_list_id = order.order_list_id();
531 let contingency_type = order.contingency_type();
532
533 self.spawn_task(task_label, async move {
534 let result = if use_broadcaster {
535 submitter
536 .broadcast_submit(
537 instrument_id,
538 client_order_id,
539 order_side,
540 order_type,
541 quantity,
542 time_in_force,
543 price,
544 trigger_price,
545 trigger_type,
546 trailing_offset,
547 trailing_offset_type,
548 display_qty,
549 post_only,
550 reduce_only,
551 order_list_id,
552 contingency_type,
553 submit_tries,
554 peg_price_type,
555 peg_offset_value,
556 )
557 .await
558 } else {
559 http_client
560 .submit_order(
561 instrument_id,
562 client_order_id,
563 order_side,
564 order_type,
565 quantity,
566 time_in_force,
567 price,
568 trigger_price,
569 trigger_type,
570 trailing_offset,
571 trailing_offset_type,
572 display_qty,
573 post_only,
574 reduce_only,
575 order_list_id,
576 contingency_type,
577 peg_price_type,
578 peg_offset_value,
579 )
580 .await
581 };
582
583 match result {
584 Ok(_report) => {
585 }
590 Err(e) => handle_submit_failure(&SubmitFailure {
591 err: &e,
592 ws_dispatch_state: &ws_dispatch_state,
593 emitter: &emitter,
594 clock,
595 strategy_id,
596 instrument_id,
597 client_order_id,
598 post_only,
599 }),
600 }
601 Ok(())
602 });
603 }
604}
605
606#[async_trait(?Send)]
607impl ExecutionClient for BitmexExecutionClient {
608 fn is_connected(&self) -> bool {
609 self.core.is_connected()
610 }
611
612 fn client_id(&self) -> ClientId {
613 self.core.client_id
614 }
615
616 fn account_id(&self) -> AccountId {
617 self.core.account_id
618 }
619
620 fn venue(&self) -> Venue {
621 self.core.venue
622 }
623
624 fn oms_type(&self) -> OmsType {
625 self.core.oms_type
626 }
627
628 fn get_account(&self) -> Option<AccountAny> {
629 self.core.cache().account_owned(&self.core.account_id)
630 }
631
632 fn generate_account_state(
633 &self,
634 balances: Vec<AccountBalance>,
635 margins: Vec<MarginBalance>,
636 reported: bool,
637 ts_event: UnixNanos,
638 ) -> anyhow::Result<()> {
639 self.emitter
640 .emit_account_state(balances, margins, reported, ts_event);
641 Ok(())
642 }
643
644 fn start(&mut self) -> anyhow::Result<()> {
645 if self.core.is_started() {
646 return Ok(());
647 }
648
649 self.emitter.set_sender(get_exec_event_sender());
650 self.core.set_started();
651 log::info!(
652 "BitMEX execution client started: client_id={}, account_id={}, environment={}, submitter_pool_size={:?}, canceller_pool_size={:?}, proxy_url={:?}, submitter_proxy_urls={:?}, canceller_proxy_urls={:?}",
653 self.core.client_id,
654 self.core.account_id,
655 self.config.environment,
656 self.config.submitter_pool_size,
657 self.config.canceller_pool_size,
658 self.config.proxy_url,
659 self.config.submitter_proxy_urls,
660 self.config.canceller_proxy_urls,
661 );
662 Ok(())
663 }
664
665 fn stop(&mut self) -> anyhow::Result<()> {
666 if self.core.is_stopped() {
667 return Ok(());
668 }
669
670 self.core.set_stopped();
671 self.core.set_disconnected();
672
673 if let Some(handle) = self.ws_stream_handle.take() {
674 handle.abort();
675 }
676
677 if let Some(handle) = self.dms_task_handle.take() {
678 handle.abort();
679 }
680 self.dms_running.store(false, Ordering::SeqCst);
681 self.abort_pending_tasks();
682 log::info!("BitMEX execution client {} stopped", self.core.client_id);
683 Ok(())
684 }
685
686 async fn connect(&mut self) -> anyhow::Result<()> {
687 if self.core.is_connected() {
688 return Ok(());
689 }
690
691 self.http_client.reset_cancellation_token();
693
694 self.ensure_instruments_initialized_async().await?;
695
696 self.refresh_account_state().await?;
697 self.await_account_registered(30.0).await?;
698
699 self.ws_client.connect().await?;
700 self.ws_client.wait_until_active(10.0).await?;
701
702 self._submitter.start().await?;
704 self._canceller.start().await?;
705
706 self.ws_client.subscribe_orders().await?;
707 self.ws_client.subscribe_executions().await?;
708 self.ws_client.subscribe_positions().await?;
709 self.ws_client.subscribe_wallet().await?;
710 if let Err(e) = self.ws_client.subscribe_margin().await {
711 log::debug!("Margin subscription unavailable: {e:?}");
712 }
713
714 self.start_ws_stream();
715
716 self.core.set_connected();
717 self.start_deadmans_switch();
718 log::info!("Connected: client_id={}", self.core.client_id);
719 Ok(())
720 }
721
722 async fn disconnect(&mut self) -> anyhow::Result<()> {
723 if self.core.is_disconnected() {
724 return Ok(());
725 }
726
727 self.stop_deadmans_switch().await;
729
730 self.http_client.cancel_all_requests();
731 self._submitter.stop().await;
732 self._canceller.stop().await;
733
734 if let Err(e) = self.ws_client.close().await {
735 log::warn!("Error while closing BitMEX execution websocket: {e:?}");
736 }
737
738 if let Some(handle) = self.ws_stream_handle.take() {
739 handle.abort();
740 }
741
742 self.abort_pending_tasks();
743 self.core.set_disconnected();
744 log::info!("Disconnected: client_id={}", self.core.client_id);
745 Ok(())
746 }
747
748 async fn generate_order_status_report(
749 &self,
750 cmd: &GenerateOrderStatusReport,
751 ) -> anyhow::Result<Option<OrderStatusReport>> {
752 let instrument_id = cmd
753 .instrument_id
754 .context("BitMEX generate_order_status_report requires an instrument identifier")?;
755
756 self.http_client
757 .query_order(
758 instrument_id,
759 cmd.client_order_id,
760 cmd.venue_order_id.map(|id| VenueOrderId::from(id.as_str())),
761 )
762 .await
763 .context("failed to query BitMEX order status")
764 }
765
766 async fn generate_order_status_reports(
767 &self,
768 cmd: &GenerateOrderStatusReports,
769 ) -> anyhow::Result<Vec<OrderStatusReport>> {
770 let start_dt = cmd.start.map(|nanos| nanos.to_datetime_utc());
771 let end_dt = cmd.end.map(|nanos| nanos.to_datetime_utc());
772
773 let mut reports = self
774 .http_client
775 .request_order_status_reports(cmd.instrument_id, cmd.open_only, start_dt, end_dt, None)
776 .await
777 .context("failed to request BitMEX order status reports")?;
778
779 if let Some(start) = cmd.start {
780 reports.retain(|report| report.ts_last >= start);
781 }
782
783 if let Some(end) = cmd.end {
784 reports.retain(|report| report.ts_last <= end);
785 }
786
787 Self::log_report_receipt(reports.len(), "OrderStatusReport", cmd.log_receipt_level);
788
789 Ok(reports)
790 }
791
792 async fn generate_fill_reports(
793 &self,
794 cmd: GenerateFillReports,
795 ) -> anyhow::Result<Vec<FillReport>> {
796 let start_dt = cmd.start.map(|nanos| nanos.to_datetime_utc());
797 let end_dt = cmd.end.map(|nanos| nanos.to_datetime_utc());
798
799 let mut reports = self
800 .http_client
801 .request_fill_reports(cmd.instrument_id, start_dt, end_dt, None)
802 .await
803 .context("failed to request BitMEX fill reports")?;
804
805 if let Some(order_id) = cmd.venue_order_id {
806 reports.retain(|report| report.venue_order_id.as_str() == order_id.as_str());
807 }
808
809 if let Some(start) = cmd.start {
810 reports.retain(|report| report.ts_event >= start);
811 }
812
813 if let Some(end) = cmd.end {
814 reports.retain(|report| report.ts_event <= end);
815 }
816
817 Self::log_report_receipt(reports.len(), "FillReport", cmd.log_receipt_level);
818
819 Ok(reports)
820 }
821
822 async fn generate_position_status_reports(
823 &self,
824 cmd: &GeneratePositionStatusReports,
825 ) -> anyhow::Result<Vec<PositionStatusReport>> {
826 let mut reports = self
827 .http_client
828 .request_position_status_reports()
829 .await
830 .context("failed to request BitMEX position reports")?;
831
832 if let Some(instrument_id) = cmd.instrument_id {
833 reports.retain(|report| report.instrument_id == instrument_id);
834 }
835
836 if let Some(start) = cmd.start {
837 reports.retain(|report| report.ts_last >= start);
838 }
839
840 if let Some(end) = cmd.end {
841 reports.retain(|report| report.ts_last <= end);
842 }
843
844 Self::log_report_receipt(reports.len(), "PositionStatusReport", cmd.log_receipt_level);
845
846 Ok(reports)
847 }
848
849 async fn generate_mass_status(
850 &self,
851 lookback_mins: Option<u64>,
852 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
853 log::info!("Generating ExecutionMassStatus (lookback_mins={lookback_mins:?})");
854
855 let ts_now = self.clock.get_time_ns();
856 let start = lookback_mins.map(|mins| {
857 let lookback_ns = mins.saturating_mul(60).saturating_mul(1_000_000_000);
858 UnixNanos::from(ts_now.as_u64().saturating_sub(lookback_ns))
859 });
860
861 let order_cmd = GenerateOrderStatusReportsBuilder::default()
862 .ts_init(ts_now)
863 .open_only(false)
864 .start(start)
865 .build()
866 .map_err(|e| anyhow::anyhow!("{e}"))?;
867
868 let fill_cmd = GenerateFillReportsBuilder::default()
869 .ts_init(ts_now)
870 .start(start)
871 .build()
872 .map_err(|e| anyhow::anyhow!("{e}"))?;
873
874 let position_cmd = GeneratePositionStatusReportsBuilder::default()
875 .ts_init(ts_now)
876 .start(start)
877 .build()
878 .map_err(|e| anyhow::anyhow!("{e}"))?;
879
880 let (order_reports, fill_reports, position_reports) = tokio::try_join!(
881 self.generate_order_status_reports(&order_cmd),
882 self.generate_fill_reports(fill_cmd),
883 self.generate_position_status_reports(&position_cmd),
884 )?;
885
886 let mut mass_status = ExecutionMassStatus::new(
887 self.core.client_id,
888 self.core.account_id,
889 self.core.venue,
890 ts_now,
891 None,
892 );
893 mass_status.add_order_reports(order_reports);
894 mass_status.add_fill_reports(fill_reports);
895 mass_status.add_position_reports(position_reports);
896
897 Ok(Some(mass_status))
898 }
899
900 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
901 let http_client = self.http_client.clone();
902 let emitter = self.emitter.clone();
903 let account_id = self.core.account_id;
904
905 self.spawn_task("query_account", async move {
906 match http_client.request_account_state(account_id).await {
907 Ok(account_state) => emitter.send_account_state(account_state),
908 Err(e) => log::error!("BitMEX query account failed: {e:?}"),
909 }
910 Ok(())
911 });
912
913 Ok(())
914 }
915
916 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
917 let http_client = self.http_client.clone();
918 let instrument_id = cmd.instrument_id;
919 let client_order_id = Some(cmd.client_order_id);
920 let venue_order_id = cmd.venue_order_id;
921 let emitter = self.emitter.clone();
922
923 self.spawn_task("query_order", async move {
924 match http_client
925 .request_order_status_report(instrument_id, client_order_id, venue_order_id)
926 .await
927 {
928 Ok(report) => emitter.send_order_status_report(report),
929 Err(e) => log::error!("BitMEX query order failed: {e:?}"),
930 }
931 Ok(())
932 });
933
934 Ok(())
935 }
936
937 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
938 let submit_tries = cmd
939 .params
940 .as_ref()
941 .and_then(|p| p.get_usize("submit_tries"))
942 .filter(|&n| n > 0);
943
944 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
945
946 let peg_price_type = match parse_peg_price_type(cmd.params.as_ref()) {
947 Ok(value) => value,
948 Err(e) => {
949 self.emitter.emit_order_denied(&order, &e.to_string());
950 return Ok(());
951 }
952 };
953 let peg_offset_value = match parse_peg_offset_value(cmd.params.as_ref()) {
954 Ok(value) => value,
955 Err(e) => {
956 self.emitter.emit_order_denied(&order, &e.to_string());
957 return Ok(());
958 }
959 };
960
961 self.submit_cached_order(
962 &order,
963 submit_tries,
964 peg_price_type,
965 peg_offset_value,
966 "submit_order",
967 );
968 Ok(())
969 }
970
971 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
972 if cmd.order_list.client_order_ids.is_empty() {
973 log::debug!("submit_order_list called with empty order list");
974 return Ok(());
975 }
976
977 let submit_tries = cmd
978 .params
979 .as_ref()
980 .and_then(|p| p.get_usize("submit_tries"))
981 .filter(|&n| n > 0);
982
983 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
984
985 let peg_price_type = match parse_peg_price_type(cmd.params.as_ref()) {
986 Ok(value) => value,
987 Err(e) => {
988 for order in &orders {
989 self.emitter.emit_order_denied(order, &e.to_string());
990 }
991 return Ok(());
992 }
993 };
994 let peg_offset_value = match parse_peg_offset_value(cmd.params.as_ref()) {
995 Ok(value) => value,
996 Err(e) => {
997 for order in &orders {
998 self.emitter.emit_order_denied(order, &e.to_string());
999 }
1000 return Ok(());
1001 }
1002 };
1003
1004 log::debug!(
1005 "Submitting BitMEX order list: order_list_id={}, count={}",
1006 cmd.order_list.id,
1007 orders.len(),
1008 );
1009
1010 for order in orders {
1011 self.submit_cached_order(
1012 &order,
1013 submit_tries,
1014 peg_price_type,
1015 peg_offset_value,
1016 "submit_order_list_item",
1017 );
1018 }
1019
1020 Ok(())
1021 }
1022
1023 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1024 self.ensure_order_identity(cmd.client_order_id, cmd.strategy_id, cmd.instrument_id);
1025 let http_client = self.http_client.clone();
1026 let emitter = self.emitter.clone();
1027 let clock = self.clock;
1028 let instrument_id = cmd.instrument_id;
1029 let client_order_id = cmd.client_order_id;
1030 let client_order_id_opt = Some(client_order_id);
1031 let venue_order_id = cmd.venue_order_id;
1032 let quantity = cmd.quantity;
1033 let price = cmd.price;
1034 let trigger_price = cmd.trigger_price;
1035 let strategy_id = cmd.strategy_id;
1036
1037 self.spawn_task("modify_order", async move {
1038 match http_client
1039 .modify_order(
1040 instrument_id,
1041 client_order_id_opt,
1042 venue_order_id,
1043 quantity,
1044 price,
1045 trigger_price,
1046 )
1047 .await
1048 {
1049 Ok(_) => {
1050 log::debug!(
1051 "BitMEX modify accepted by REST, awaiting websocket confirmation: client_order_id={client_order_id}"
1052 );
1053 }
1054 Err(e) => handle_modify_failure(&ModifyFailure {
1055 err: &e,
1056 emitter: &emitter,
1057 clock,
1058 strategy_id,
1059 instrument_id,
1060 client_order_id,
1061 venue_order_id,
1062 }),
1063 }
1064 Ok(())
1065 });
1066
1067 Ok(())
1068 }
1069
1070 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
1071 self.ensure_order_identity(cmd.client_order_id, cmd.strategy_id, cmd.instrument_id);
1072 let canceller = self._canceller.clone_for_async();
1073 let emitter = self.emitter.clone();
1074 let dispatch_state = Arc::clone(&self.ws_dispatch_state);
1075 let instrument_id = cmd.instrument_id;
1076 let client_order_id = Some(cmd.client_order_id);
1077 let venue_order_id = cmd.venue_order_id;
1078
1079 self.spawn_task("cancel_order", async move {
1080 match canceller
1081 .broadcast_cancel(instrument_id, client_order_id, venue_order_id)
1082 .await
1083 {
1084 Ok(Some(report)) => {
1085 if let Some(cid) = &report.client_order_id {
1086 dispatch_state.tombstone_order(cid);
1087 }
1088 emitter.send_order_status_report(report);
1089 }
1090 Ok(None) => {
1091 log::debug!("Order already cancelled: {client_order_id:?}");
1092 }
1093 Err(e) => log::error!("BitMEX cancel order failed: {e:?}"),
1094 }
1095 Ok(())
1096 });
1097
1098 Ok(())
1099 }
1100
1101 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
1102 let canceller = self._canceller.clone_for_async();
1103 let emitter = self.emitter.clone();
1104 let dispatch_state = Arc::clone(&self.ws_dispatch_state);
1105 let instrument_id = cmd.instrument_id;
1106 let order_side = if cmd.order_side == OrderSide::NoOrderSide {
1107 log::debug!(
1108 "BitMEX cancel_all_orders received NoOrderSide for {instrument_id}, using unfiltered cancel-all",
1109 );
1110 None
1111 } else {
1112 Some(cmd.order_side)
1113 };
1114
1115 self.spawn_task("cancel_all_orders", async move {
1116 match canceller
1117 .broadcast_cancel_all(instrument_id, order_side)
1118 .await
1119 {
1120 Ok(reports) => {
1121 for report in &reports {
1122 if let Some(cid) = &report.client_order_id {
1123 dispatch_state.tombstone_order(cid);
1124 }
1125 }
1126
1127 for report in reports {
1128 emitter.send_order_status_report(report);
1129 }
1130 }
1131 Err(e) => log::error!("BitMEX cancel all failed: {e:?}"),
1132 }
1133 Ok(())
1134 });
1135
1136 Ok(())
1137 }
1138
1139 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1140 let canceller = self._canceller.clone_for_async();
1141 let emitter = self.emitter.clone();
1142 let dispatch_state = Arc::clone(&self.ws_dispatch_state);
1143 let instrument_id = cmd.instrument_id;
1144
1145 let client_ids: Vec<ClientOrderId> = cmd
1146 .cancels
1147 .iter()
1148 .map(|cancel| cancel.client_order_id)
1149 .collect();
1150
1151 let venue_ids: Vec<VenueOrderId> = cmd
1152 .cancels
1153 .iter()
1154 .filter_map(|cancel| cancel.venue_order_id)
1155 .collect();
1156
1157 let client_ids_opt = if client_ids.is_empty() {
1158 None
1159 } else {
1160 Some(client_ids)
1161 };
1162
1163 let venue_ids_opt = if venue_ids.is_empty() {
1164 None
1165 } else {
1166 Some(venue_ids)
1167 };
1168
1169 self.spawn_task("batch_cancel_orders", async move {
1170 match canceller
1171 .broadcast_batch_cancel(instrument_id, client_ids_opt, venue_ids_opt)
1172 .await
1173 {
1174 Ok(reports) => {
1175 for report in &reports {
1176 if let Some(cid) = &report.client_order_id {
1177 dispatch_state.tombstone_order(cid);
1178 }
1179 }
1180
1181 for report in reports {
1182 emitter.send_order_status_report(report);
1183 }
1184 }
1185 Err(e) => log::error!("BitMEX batch cancel failed: {e:?}"),
1186 }
1187 Ok(())
1188 });
1189
1190 Ok(())
1191 }
1192}
1193
1194struct SubmitFailure<'a> {
1195 err: &'a anyhow::Error,
1196 ws_dispatch_state: &'a Arc<WsDispatchState>,
1197 emitter: &'a ExecutionEventEmitter,
1198 clock: &'static AtomicTime,
1199 strategy_id: StrategyId,
1200 instrument_id: InstrumentId,
1201 client_order_id: ClientOrderId,
1202 post_only: bool,
1203}
1204
1205fn handle_submit_failure(failure: &SubmitFailure<'_>) {
1206 let error_msg = failure.err.to_string();
1207
1208 if is_bitmex_duplicate_clordid_submit_failure(failure.err) {
1210 log::warn!(
1211 "Order {} may exist (duplicate clOrdID), \
1212 awaiting WebSocket confirmation",
1213 failure.client_order_id,
1214 );
1215 return;
1216 }
1217
1218 if is_definitive_bitmex_submit_rejection(failure.err) {
1219 failure
1220 .ws_dispatch_state
1221 .order_identities
1222 .remove(&failure.client_order_id);
1223 let ts_event = failure.clock.get_time_ns();
1224 let rejection_reason = error_msg
1225 .strip_prefix(DEFINITIVE_SUBMIT_REJECTION)
1226 .map_or(error_msg.as_str(), |msg| {
1227 msg.trim_start_matches(':').trim_start()
1228 });
1229 failure.emitter.emit_order_rejected_event(
1230 failure.strategy_id,
1231 failure.instrument_id,
1232 failure.client_order_id,
1233 &format!("submit-order-error: {rejection_reason}"),
1234 ts_event,
1235 failure.post_only,
1236 );
1237 } else {
1238 log::warn!(
1239 "Ambiguous BitMEX submit failure for {}, awaiting reconciliation: {:?}",
1240 failure.client_order_id,
1241 failure.err,
1242 );
1243 }
1244}
1245
1246struct ModifyFailure<'a> {
1247 err: &'a anyhow::Error,
1248 emitter: &'a ExecutionEventEmitter,
1249 clock: &'static AtomicTime,
1250 strategy_id: StrategyId,
1251 instrument_id: InstrumentId,
1252 client_order_id: ClientOrderId,
1253 venue_order_id: Option<VenueOrderId>,
1254}
1255
1256fn handle_modify_failure(failure: &ModifyFailure<'_>) {
1257 if is_definitive_bitmex_modify_rejection(failure.err) {
1258 let ts_event = failure.clock.get_time_ns();
1259 failure.emitter.emit_order_modify_rejected_event(
1260 failure.strategy_id,
1261 failure.instrument_id,
1262 failure.client_order_id,
1263 failure.venue_order_id,
1264 &format!("modify-order-error: {}", failure.err),
1265 ts_event,
1266 );
1267 } else {
1268 log::warn!(
1269 "Ambiguous BitMEX modify failure for {}, awaiting reconciliation: {:?}",
1270 failure.client_order_id,
1271 failure.err,
1272 );
1273 }
1274}
1275
1276fn validate_order_for_bitmex_submit(
1277 order: &OrderAny,
1278 peg_price_type: Option<BitmexPegPriceType>,
1279 peg_offset_value: Option<f64>,
1280) -> anyhow::Result<()> {
1281 if order.order_side() == OrderSide::NoOrderSide {
1282 anyhow::bail!("Order side must be Buy or Sell");
1283 }
1284
1285 BitmexOrderType::try_from_order_type(order.order_type())?;
1286 BitmexTimeInForce::try_from_time_in_force(order.time_in_force())?;
1287
1288 let is_trailing_stop = matches!(
1289 order.order_type(),
1290 OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
1291 );
1292
1293 if is_trailing_stop
1294 && let Some(offset_type) = order.trailing_offset_type()
1295 && offset_type != TrailingOffsetType::Price
1296 {
1297 anyhow::bail!("BitMEX only supports PRICE trailing offset type, was {offset_type:?}");
1298 }
1299
1300 if peg_price_type.is_none() && peg_offset_value.is_some() {
1301 anyhow::bail!("`peg_offset_value` requires `peg_price_type`");
1302 }
1303
1304 if peg_price_type.is_some() && order.order_type() != OrderType::Limit {
1305 let order_type = order.order_type();
1306 anyhow::bail!("Pegged orders only supported for LIMIT order type, was {order_type:?}");
1307 }
1308
1309 if let Some(contingency_type) = order.contingency_type() {
1310 BitmexContingencyType::try_from(contingency_type)?;
1311 }
1312
1313 Ok(())
1314}
1315
1316fn is_definitive_bitmex_submit_rejection(err: &anyhow::Error) -> bool {
1317 if is_bitmex_duplicate_clordid_submit_failure(err) {
1318 return false;
1319 }
1320
1321 if has_bitmex_api_refusal(err) {
1322 return true;
1323 }
1324
1325 let message = err.to_string();
1326 message.starts_with("Order rejected:") || message.starts_with(DEFINITIVE_SUBMIT_REJECTION)
1327}
1328
1329fn is_bitmex_duplicate_clordid_submit_failure(err: &anyhow::Error) -> bool {
1330 if err.to_string().contains("IDEMPOTENT_DUPLICATE") {
1331 return true;
1332 }
1333
1334 err.chain().any(|cause| {
1335 cause
1336 .downcast_ref::<BitmexHttpError>()
1337 .is_some_and(|e| {
1338 matches!(e, BitmexHttpError::BitmexError { message, .. } if message.contains("Duplicate clOrdID"))
1339 })
1340 })
1341}
1342
1343fn is_definitive_bitmex_modify_rejection(err: &anyhow::Error) -> bool {
1344 if has_bitmex_api_refusal(err) {
1345 return true;
1346 }
1347
1348 err.to_string().starts_with("Order modification rejected:")
1349}
1350
1351fn has_bitmex_api_refusal(err: &anyhow::Error) -> bool {
1352 err.chain().any(|cause| {
1353 cause
1354 .downcast_ref::<BitmexHttpError>()
1355 .is_some_and(|e| matches!(e, BitmexHttpError::BitmexError { .. }))
1356 })
1357}
1358
1359#[cfg(test)]
1360mod tests {
1361 use std::{cell::RefCell, rc::Rc};
1362
1363 use nautilus_common::{
1364 cache::Cache,
1365 clients::ExecutionClient,
1366 messages::{ExecutionEvent, ExecutionReport},
1367 };
1368 use nautilus_core::{Params, UUID4};
1369 use nautilus_model::{
1370 enums::TimeInForce,
1371 events::OrderEventAny,
1372 identifiers::{Symbol, TraderId},
1373 instruments::crypto_perpetual::CryptoPerpetual,
1374 orders::builder::OrderTestBuilder,
1375 types::{Currency, Price, Quantity},
1376 };
1377 use nautilus_network::http::StatusCode;
1378 use rstest::rstest;
1379
1380 use super::*;
1381 use crate::{
1382 common::{
1383 consts::{BITMEX_CLIENT_ID, BITMEX_VENUE},
1384 testing::load_test_json,
1385 },
1386 websocket::{
1387 enums::BitmexAction,
1388 messages::{
1389 BitmexExecutionMsg, BitmexOrderMsg, BitmexTableMessage, BitmexWalletMsg,
1390 BitmexWsMessage, OrderData,
1391 },
1392 },
1393 };
1394
1395 fn bitmex_api_error() -> anyhow::Error {
1396 anyhow::Error::new(BitmexHttpError::BitmexError {
1397 error_name: "HTTPError".to_string(),
1398 message: "Invalid price".to_string(),
1399 })
1400 }
1401
1402 fn test_execution_client() -> (BitmexExecutionClient, Rc<RefCell<Cache>>) {
1403 let cache = Rc::new(RefCell::new(Cache::default()));
1404 let core = ExecutionClientCore::new(
1405 TraderId::from("TESTER-001"),
1406 *BITMEX_CLIENT_ID,
1407 *BITMEX_VENUE,
1408 OmsType::Netting,
1409 AccountId::from("BITMEX-001"),
1410 AccountType::Margin,
1411 None,
1412 cache.clone(),
1413 );
1414 let config = BitmexExecClientConfig {
1415 api_key: Some("test_key".to_string()),
1416 api_secret: Some("test_secret".to_string()),
1417 base_url_http: Some("http://127.0.0.1:9/api/v1".to_string()),
1418 base_url_ws: Some("ws://127.0.0.1:9/realtime".to_string()),
1419 ..Default::default()
1420 };
1421
1422 (BitmexExecutionClient::new(core, config).unwrap(), cache)
1423 }
1424
1425 fn make_emitter() -> (
1426 ExecutionEventEmitter,
1427 tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
1428 ) {
1429 let mut emitter = ExecutionEventEmitter::new(
1430 get_atomic_clock_realtime(),
1431 TraderId::from("TESTER-001"),
1432 AccountId::from("BITMEX-001"),
1433 AccountType::Margin,
1434 None,
1435 );
1436 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
1437 emitter.set_sender(tx);
1438 (emitter, rx)
1439 }
1440
1441 fn limit_order() -> OrderAny {
1442 limit_order_with_id(ClientOrderId::from("O-LIMIT"))
1443 }
1444
1445 fn limit_order_with_id(client_order_id: ClientOrderId) -> OrderAny {
1446 let mut builder = OrderTestBuilder::new(OrderType::Limit);
1447 builder
1448 .instrument_id(InstrumentId::from("XBTUSD.BITMEX"))
1449 .client_order_id(client_order_id)
1450 .side(OrderSide::Buy)
1451 .quantity(Quantity::from("1"))
1452 .price(Price::from("100.0"))
1453 .build()
1454 }
1455
1456 fn test_perpetual_instrument() -> InstrumentAny {
1457 InstrumentAny::CryptoPerpetual(CryptoPerpetual::new(
1458 InstrumentId::from("XBTUSD.BITMEX"),
1459 Symbol::new("XBTUSD"),
1460 Currency::BTC(),
1461 Currency::USD(),
1462 Currency::BTC(),
1463 true,
1464 1,
1465 0,
1466 Price::new(0.5, 1),
1467 Quantity::new(1.0, 0),
1468 None,
1469 None,
1470 None,
1471 None,
1472 None,
1473 None,
1474 None,
1475 None,
1476 None,
1477 None,
1478 None,
1479 None,
1480 None,
1481 None,
1482 UnixNanos::default(),
1483 UnixNanos::default(),
1484 ))
1485 }
1486
1487 fn market_order() -> OrderAny {
1488 let mut builder = OrderTestBuilder::new(OrderType::Market);
1489 builder
1490 .instrument_id(InstrumentId::from("XBTUSD.BITMEX"))
1491 .quantity(Quantity::from("1"))
1492 .build()
1493 }
1494
1495 fn order_identity(order: &OrderAny) -> OrderIdentity {
1496 OrderIdentity {
1497 instrument_id: order.instrument_id(),
1498 strategy_id: order.strategy_id(),
1499 order_side: order.order_side(),
1500 order_type: order.order_type(),
1501 }
1502 }
1503
1504 fn submit_command(order: &OrderAny, params: Option<Params>) -> SubmitOrder {
1505 SubmitOrder::new(
1506 order.trader_id(),
1507 Some(*BITMEX_CLIENT_ID),
1508 order.strategy_id(),
1509 order.instrument_id(),
1510 order.client_order_id(),
1511 order.init_event().clone(),
1512 None,
1513 None,
1514 params,
1515 UUID4::new(),
1516 UnixNanos::default(),
1517 None,
1518 )
1519 }
1520
1521 fn drain_order_events(
1522 rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
1523 ) -> Vec<OrderEventAny> {
1524 let mut events = Vec::new();
1525
1526 while let Ok(event) = rx.try_recv() {
1527 if let ExecutionEvent::Order(event) = event {
1528 events.push(event);
1529 }
1530 }
1531 events
1532 }
1533
1534 fn dispatch_execution_fixture(
1535 state: &WsDispatchState,
1536 emitter: &ExecutionEventEmitter,
1537 account_id: AccountId,
1538 ) {
1539 let exec_msg: BitmexExecutionMsg =
1540 serde_json::from_str(&load_test_json("ws_execution.json")).unwrap();
1541 let mut instruments_by_symbol = AHashMap::new();
1542 instruments_by_symbol.insert(Ustr::from("XBTUSD"), test_perpetual_instrument());
1543 let mut order_type_cache = AHashMap::new();
1544 let mut order_symbol_cache = AHashMap::new();
1545
1546 dispatch::dispatch_ws_message(
1547 UnixNanos::default(),
1548 BitmexWsMessage::Table(BitmexTableMessage::Execution {
1549 action: BitmexAction::Insert,
1550 data: vec![exec_msg],
1551 }),
1552 emitter,
1553 state,
1554 &mut instruments_by_symbol,
1555 &mut order_type_cache,
1556 &mut order_symbol_cache,
1557 account_id,
1558 );
1559 }
1560
1561 #[rstest]
1562 fn test_bitmex_api_error_is_definitive_submit_rejection() {
1563 let err = bitmex_api_error();
1564
1565 assert!(is_definitive_bitmex_submit_rejection(&err));
1566 }
1567
1568 #[rstest]
1569 fn test_config_account_id_seeds_core_account_id() {
1570 let cache = Rc::new(RefCell::new(Cache::default()));
1571 let core = ExecutionClientCore::new(
1572 TraderId::from("TESTER-001"),
1573 *BITMEX_CLIENT_ID,
1574 *BITMEX_VENUE,
1575 OmsType::Netting,
1576 AccountId::from("BITMEX-001"),
1577 AccountType::Margin,
1578 None,
1579 cache,
1580 );
1581 let config = BitmexExecClientConfig {
1582 api_key: Some("test_key".to_string()),
1583 api_secret: Some("test_secret".to_string()),
1584 account_id: Some(AccountId::from("BITMEX-319111")),
1585 base_url_http: Some("http://127.0.0.1:9/api/v1".to_string()),
1586 base_url_ws: Some("ws://127.0.0.1:9/realtime".to_string()),
1587 ..Default::default()
1588 };
1589
1590 let client = BitmexExecutionClient::new(core, config).unwrap();
1591
1592 assert_eq!(client.account_id(), AccountId::from("BITMEX-319111"));
1593 }
1594
1595 #[rstest]
1596 fn test_apply_account_id_updates_core_emitter_and_websocket_client() {
1597 let (mut client, _) = test_execution_client();
1598 let account_id = AccountId::from("BITMEX-319111");
1599
1600 client.apply_account_id(account_id);
1601
1602 assert_eq!(client.account_id(), account_id);
1603 assert_eq!(client.emitter.account_id(), account_id);
1604 assert_eq!(client.ws_client.account_id(), account_id);
1605 }
1606
1607 #[rstest]
1608 fn test_dispatch_tracked_fill_uses_bitmex_account_id() {
1609 let (emitter, mut rx) = make_emitter();
1610 let state = WsDispatchState::default();
1611 let account_id = AccountId::from("BITMEX-1234567");
1612 let client_order_id = ClientOrderId::from("mm_bitmex_2b/oemUeQ4CAJZgP3fjHsB");
1613 state.order_identities.insert(
1614 client_order_id,
1615 OrderIdentity {
1616 instrument_id: InstrumentId::from("XBTUSD.BITMEX"),
1617 strategy_id: StrategyId::from("S-001"),
1618 order_side: OrderSide::Sell,
1619 order_type: OrderType::Limit,
1620 },
1621 );
1622
1623 dispatch_execution_fixture(&state, &emitter, account_id);
1624
1625 let events = drain_order_events(&mut rx);
1626 assert_eq!(events.len(), 2);
1627 match &events[..] {
1628 [
1629 OrderEventAny::Accepted(accepted),
1630 OrderEventAny::Filled(filled),
1631 ] => {
1632 assert_eq!(accepted.account_id, account_id);
1633 assert_eq!(filled.account_id, account_id);
1634 }
1635 events => panic!("expected accepted and filled events, was {events:?}"),
1636 }
1637 }
1638
1639 #[rstest]
1640 fn test_dispatch_untracked_fill_report_uses_bitmex_account_id() {
1641 let (emitter, mut rx) = make_emitter();
1642 let state = WsDispatchState::default();
1643 let account_id = AccountId::from("BITMEX-1234567");
1644
1645 dispatch_execution_fixture(&state, &emitter, account_id);
1646
1647 match rx.try_recv().unwrap() {
1648 ExecutionEvent::Report(ExecutionReport::Fill(report)) => {
1649 assert_eq!(report.account_id, account_id);
1650 }
1651 event => panic!("expected fill report, was {event:?}"),
1652 }
1653 assert!(rx.try_recv().is_err());
1654 }
1655
1656 #[rstest]
1657 #[case::continuous(false)]
1658 #[case::reconnected(true)]
1659 fn test_dispatch_sparse_terminal_update_respects_cache_lifecycle(#[case] reconnect: bool) {
1660 let (emitter, mut rx) = make_emitter();
1661 let state = WsDispatchState::default();
1662 let account_id = AccountId::from("BITMEX-1234567");
1663 let client_order_id = ClientOrderId::from("mm_bitmex_1a/oemUeQ4CAJZgP3fjHsA");
1664 let order: BitmexOrderMsg = serde_json::from_str(&load_test_json("ws_order.json")).unwrap();
1665 let update: BitmexTableMessage =
1666 serde_json::from_str(&load_test_json("ws_order_update_canceled.json")).unwrap();
1667 let mut instruments_by_symbol = AHashMap::new();
1668 instruments_by_symbol.insert(Ustr::from("XBTUSD"), test_perpetual_instrument());
1669 let mut order_type_cache = AHashMap::new();
1670 let mut order_symbol_cache = AHashMap::new();
1671 state.order_identities.insert(
1672 client_order_id,
1673 OrderIdentity {
1674 instrument_id: InstrumentId::from("XBTUSD.BITMEX"),
1675 strategy_id: StrategyId::from("S-001"),
1676 order_side: OrderSide::Buy,
1677 order_type: OrderType::Limit,
1678 },
1679 );
1680
1681 dispatch::dispatch_ws_message(
1682 UnixNanos::default(),
1683 BitmexWsMessage::Table(BitmexTableMessage::Order {
1684 action: BitmexAction::Partial,
1685 data: vec![OrderData::Full(order)],
1686 }),
1687 &emitter,
1688 &state,
1689 &mut instruments_by_symbol,
1690 &mut order_type_cache,
1691 &mut order_symbol_cache,
1692 account_id,
1693 );
1694
1695 if reconnect {
1696 dispatch::dispatch_ws_message(
1697 UnixNanos::default(),
1698 BitmexWsMessage::Reconnected,
1699 &emitter,
1700 &state,
1701 &mut instruments_by_symbol,
1702 &mut order_type_cache,
1703 &mut order_symbol_cache,
1704 account_id,
1705 );
1706 }
1707 dispatch::dispatch_ws_message(
1708 UnixNanos::default(),
1709 BitmexWsMessage::Table(update),
1710 &emitter,
1711 &state,
1712 &mut instruments_by_symbol,
1713 &mut order_type_cache,
1714 &mut order_symbol_cache,
1715 account_id,
1716 );
1717
1718 let events = drain_order_events(&mut rx);
1719 match (reconnect, &events[..]) {
1720 (false, [OrderEventAny::Accepted(_), OrderEventAny::Canceled(_)])
1721 | (true, [OrderEventAny::Accepted(_)]) => {}
1722 (_, events) => panic!("unexpected order lifecycle events: {events:?}"),
1723 }
1724 }
1725
1726 #[rstest]
1727 fn test_dispatch_untracked_sparse_terminal_update_does_not_emit_report() {
1728 let (emitter, mut rx) = make_emitter();
1729 let state = WsDispatchState::default();
1730 let account_id = AccountId::from("BITMEX-1234567");
1731 let order: BitmexOrderMsg = serde_json::from_str(&load_test_json("ws_order.json")).unwrap();
1732 let update: BitmexTableMessage =
1733 serde_json::from_str(&load_test_json("ws_order_update_canceled.json")).unwrap();
1734 let mut instruments_by_symbol = AHashMap::new();
1735 instruments_by_symbol.insert(Ustr::from("XBTUSD"), test_perpetual_instrument());
1736 let mut order_type_cache = AHashMap::new();
1737 let mut order_symbol_cache = AHashMap::new();
1738
1739 dispatch::dispatch_ws_message(
1740 UnixNanos::default(),
1741 BitmexWsMessage::Table(BitmexTableMessage::Order {
1742 action: BitmexAction::Partial,
1743 data: vec![OrderData::Full(order)],
1744 }),
1745 &emitter,
1746 &state,
1747 &mut instruments_by_symbol,
1748 &mut order_type_cache,
1749 &mut order_symbol_cache,
1750 account_id,
1751 );
1752 assert!(matches!(
1753 rx.try_recv(),
1754 Ok(ExecutionEvent::Report(ExecutionReport::Order(_)))
1755 ));
1756
1757 dispatch::dispatch_ws_message(
1758 UnixNanos::default(),
1759 BitmexWsMessage::Table(update),
1760 &emitter,
1761 &state,
1762 &mut instruments_by_symbol,
1763 &mut order_type_cache,
1764 &mut order_symbol_cache,
1765 account_id,
1766 );
1767
1768 assert!(rx.try_recv().is_err());
1769 }
1770
1771 #[rstest]
1772 fn test_dispatch_wallet_account_state_uses_bitmex_account_id() {
1773 let (emitter, mut rx) = make_emitter();
1774 let state = WsDispatchState::default();
1775 let account_id = AccountId::from("BITMEX-1234567");
1776 let wallet_msg: BitmexWalletMsg =
1777 serde_json::from_str(&load_test_json("ws_wallet.json")).unwrap();
1778 let mut instruments_by_symbol = AHashMap::new();
1779 let mut order_type_cache = AHashMap::new();
1780 let mut order_symbol_cache = AHashMap::new();
1781
1782 dispatch::dispatch_ws_message(
1783 UnixNanos::default(),
1784 BitmexWsMessage::Table(BitmexTableMessage::Wallet {
1785 action: BitmexAction::Insert,
1786 data: vec![wallet_msg],
1787 }),
1788 &emitter,
1789 &state,
1790 &mut instruments_by_symbol,
1791 &mut order_type_cache,
1792 &mut order_symbol_cache,
1793 account_id,
1794 );
1795
1796 match rx.try_recv().unwrap() {
1797 ExecutionEvent::Account(state) => {
1798 assert_eq!(state.account_id, account_id);
1799 }
1800 event => panic!("expected account state, was {event:?}"),
1801 }
1802 assert!(rx.try_recv().is_err());
1803 }
1804
1805 #[rstest]
1806 fn test_bitmex_api_error_is_definitive_modify_rejection() {
1807 let err = bitmex_api_error();
1808
1809 assert!(is_definitive_bitmex_modify_rejection(&err));
1810 }
1811
1812 #[rstest]
1813 fn test_parsed_submit_reject_is_definitive_submit_rejection() {
1814 let err = anyhow::anyhow!("Order rejected: Price is invalid");
1815
1816 assert!(is_definitive_bitmex_submit_rejection(&err));
1817 assert!(!is_definitive_bitmex_modify_rejection(&err));
1818 }
1819
1820 #[rstest]
1821 fn test_broadcast_submit_refusal_is_definitive_submit_rejection() {
1822 let err =
1823 anyhow::anyhow!("{DEFINITIVE_SUBMIT_REJECTION}: All submit requests were refused");
1824
1825 assert!(is_definitive_bitmex_submit_rejection(&err));
1826 assert!(!is_definitive_bitmex_modify_rejection(&err));
1827 }
1828
1829 #[rstest]
1830 fn test_duplicate_clordid_is_ambiguous_submit_failure() {
1831 let err = anyhow::Error::new(BitmexHttpError::BitmexError {
1832 error_name: "HTTPError".to_string(),
1833 message: "Duplicate clOrdID".to_string(),
1834 });
1835
1836 assert!(is_bitmex_duplicate_clordid_submit_failure(&err));
1837 assert!(!is_definitive_bitmex_submit_rejection(&err));
1838 }
1839
1840 #[rstest]
1841 fn test_parsed_modify_reject_is_definitive_modify_rejection() {
1842 let err = anyhow::anyhow!("Order modification rejected: Price is invalid");
1843
1844 assert!(is_definitive_bitmex_modify_rejection(&err));
1845 assert!(!is_definitive_bitmex_submit_rejection(&err));
1846 }
1847
1848 #[rstest]
1849 fn test_network_error_is_ambiguous_command_failure() {
1850 let err = anyhow::Error::new(BitmexHttpError::NetworkError("timeout".to_string()));
1851
1852 assert!(!is_definitive_bitmex_submit_rejection(&err));
1853 assert!(!is_definitive_bitmex_modify_rejection(&err));
1854 }
1855
1856 #[rstest]
1857 fn test_canceled_request_is_ambiguous_command_failure() {
1858 let err = anyhow::Error::new(BitmexHttpError::Canceled("shutdown".to_string()));
1859
1860 assert!(!is_definitive_bitmex_submit_rejection(&err));
1861 assert!(!is_definitive_bitmex_modify_rejection(&err));
1862 }
1863
1864 #[rstest]
1865 fn test_unstructured_http_status_is_ambiguous_command_failure() {
1866 let err = anyhow::Error::new(BitmexHttpError::UnexpectedStatus {
1867 status: StatusCode::BAD_GATEWAY,
1868 body: "bad gateway".to_string(),
1869 });
1870
1871 assert!(!is_definitive_bitmex_submit_rejection(&err));
1872 assert!(!is_definitive_bitmex_modify_rejection(&err));
1873 }
1874
1875 #[rstest]
1876 fn test_validate_order_for_bitmex_submit_requires_peg_type_for_offset() {
1877 let order = limit_order();
1878 let err = validate_order_for_bitmex_submit(&order, None, Some(1.0)).unwrap_err();
1879
1880 assert!(err.to_string().contains("`peg_offset_value` requires"));
1881 }
1882
1883 #[rstest]
1884 fn test_validate_order_for_bitmex_submit_rejects_pegged_market_order() {
1885 let order = market_order();
1886 let err = validate_order_for_bitmex_submit(&order, Some(BitmexPegPriceType::LastPeg), None)
1887 .unwrap_err();
1888
1889 assert!(err.to_string().contains("Pegged orders only supported"));
1890 }
1891
1892 #[rstest]
1893 fn test_submit_order_invalid_peg_params_emits_denied_without_submitted() {
1894 let (mut client, cache) = test_execution_client();
1895 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1896 client.emitter.set_sender(tx);
1897
1898 let order = limit_order_with_id(ClientOrderId::from("O-INVALID-PEG"));
1899 cache
1900 .borrow_mut()
1901 .add_order(order.clone(), None, Some(*BITMEX_CLIENT_ID), false)
1902 .unwrap();
1903
1904 let mut params = Params::new();
1905 params.insert("peg_price_type".to_string(), serde_json::json!("BadPeg"));
1906
1907 client
1908 .submit_order(submit_command(&order, Some(params)))
1909 .unwrap();
1910
1911 let events = drain_order_events(&mut rx);
1912 assert_eq!(events.len(), 1);
1913 match &events[0] {
1914 OrderEventAny::Denied(denied) => {
1915 assert_eq!(denied.client_order_id, order.client_order_id());
1916 assert_eq!(denied.reason.to_string(), "Invalid peg_price_type: BadPeg");
1917 }
1918 event => panic!("expected OrderDenied event, was {event:?}"),
1919 }
1920 assert!(
1921 !client
1922 .ws_dispatch_state
1923 .order_identities
1924 .contains_key(&order.client_order_id())
1925 );
1926 }
1927
1928 #[rstest]
1929 fn test_submit_order_gtd_time_in_force_emits_denied_without_submitted() {
1930 let (mut client, cache) = test_execution_client();
1931 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1932 client.emitter.set_sender(tx);
1933
1934 let mut builder = OrderTestBuilder::new(OrderType::Limit);
1935 let order = builder
1936 .instrument_id(InstrumentId::from("XBTUSD.BITMEX"))
1937 .client_order_id(ClientOrderId::from("O-GTD"))
1938 .side(OrderSide::Buy)
1939 .quantity(Quantity::from("1"))
1940 .price(Price::from("100.0"))
1941 .time_in_force(TimeInForce::Gtd)
1942 .expire_time(UnixNanos::from(1_000_000_000_u64))
1943 .build();
1944 cache
1945 .borrow_mut()
1946 .add_order(order.clone(), None, Some(*BITMEX_CLIENT_ID), false)
1947 .unwrap();
1948
1949 client.submit_order(submit_command(&order, None)).unwrap();
1950
1951 let events = drain_order_events(&mut rx);
1952 assert_eq!(events.len(), 1);
1953 match &events[0] {
1954 OrderEventAny::Denied(denied) => {
1955 assert_eq!(denied.client_order_id, order.client_order_id());
1956 assert!(
1957 denied
1958 .reason
1959 .to_string()
1960 .contains("GTD time in force is not supported")
1961 );
1962 }
1963 event => panic!("expected OrderDenied event, was {event:?}"),
1964 }
1965 assert!(
1966 !client
1967 .ws_dispatch_state
1968 .order_identities
1969 .contains_key(&order.client_order_id())
1970 );
1971 }
1972
1973 #[rstest]
1974 fn test_submit_failure_definitive_refusal_removes_identity_and_emits_rejected() {
1975 let (emitter, mut rx) = make_emitter();
1976 let ws_dispatch_state = Arc::new(WsDispatchState::default());
1977 let order = limit_order_with_id(ClientOrderId::from("O-SUBMIT-REJECTED"));
1978 ws_dispatch_state
1979 .order_identities
1980 .insert(order.client_order_id(), order_identity(&order));
1981
1982 let err = anyhow::anyhow!(
1983 "{DEFINITIVE_SUBMIT_REJECTION}: All submit requests were refused by BitMEX"
1984 );
1985
1986 handle_submit_failure(&SubmitFailure {
1987 err: &err,
1988 ws_dispatch_state: &ws_dispatch_state,
1989 emitter: &emitter,
1990 clock: get_atomic_clock_realtime(),
1991 strategy_id: order.strategy_id(),
1992 instrument_id: order.instrument_id(),
1993 client_order_id: order.client_order_id(),
1994 post_only: false,
1995 });
1996
1997 assert!(
1998 !ws_dispatch_state
1999 .order_identities
2000 .contains_key(&order.client_order_id())
2001 );
2002
2003 let events = drain_order_events(&mut rx);
2004 assert_eq!(events.len(), 1);
2005 match &events[0] {
2006 OrderEventAny::Rejected(rejected) => {
2007 assert_eq!(rejected.client_order_id, order.client_order_id());
2008 assert_eq!(
2009 rejected.reason.to_string(),
2010 "submit-order-error: All submit requests were refused by BitMEX"
2011 );
2012 assert!(!rejected.due_post_only);
2013 }
2014 event => panic!("expected OrderRejected event, was {event:?}"),
2015 }
2016 }
2017
2018 #[rstest]
2019 fn test_submit_failure_duplicate_clordid_keeps_identity_and_emits_no_rejection() {
2020 let (emitter, mut rx) = make_emitter();
2021 let ws_dispatch_state = Arc::new(WsDispatchState::default());
2022 let order = limit_order_with_id(ClientOrderId::from("O-DUPLICATE"));
2023 ws_dispatch_state
2024 .order_identities
2025 .insert(order.client_order_id(), order_identity(&order));
2026 let err = anyhow::Error::new(BitmexHttpError::BitmexError {
2027 error_name: "HTTPError".to_string(),
2028 message: "Duplicate clOrdID".to_string(),
2029 });
2030
2031 handle_submit_failure(&SubmitFailure {
2032 err: &err,
2033 ws_dispatch_state: &ws_dispatch_state,
2034 emitter: &emitter,
2035 clock: get_atomic_clock_realtime(),
2036 strategy_id: order.strategy_id(),
2037 instrument_id: order.instrument_id(),
2038 client_order_id: order.client_order_id(),
2039 post_only: false,
2040 });
2041
2042 assert!(
2043 ws_dispatch_state
2044 .order_identities
2045 .contains_key(&order.client_order_id())
2046 );
2047 assert!(drain_order_events(&mut rx).is_empty());
2048 }
2049
2050 #[rstest]
2051 fn test_submit_failure_network_error_keeps_identity_and_emits_no_rejection() {
2052 let (emitter, mut rx) = make_emitter();
2053 let ws_dispatch_state = Arc::new(WsDispatchState::default());
2054 let order = limit_order_with_id(ClientOrderId::from("O-SUBMIT-NETWORK"));
2055 ws_dispatch_state
2056 .order_identities
2057 .insert(order.client_order_id(), order_identity(&order));
2058 let err = anyhow::Error::new(BitmexHttpError::NetworkError("timeout".to_string()));
2059
2060 handle_submit_failure(&SubmitFailure {
2061 err: &err,
2062 ws_dispatch_state: &ws_dispatch_state,
2063 emitter: &emitter,
2064 clock: get_atomic_clock_realtime(),
2065 strategy_id: order.strategy_id(),
2066 instrument_id: order.instrument_id(),
2067 client_order_id: order.client_order_id(),
2068 post_only: false,
2069 });
2070
2071 assert!(
2072 ws_dispatch_state
2073 .order_identities
2074 .contains_key(&order.client_order_id())
2075 );
2076 assert!(drain_order_events(&mut rx).is_empty());
2077 }
2078
2079 #[rstest]
2080 fn test_modify_failure_definitive_refusal_emits_modify_rejected() {
2081 let (emitter, mut rx) = make_emitter();
2082 let order = limit_order_with_id(ClientOrderId::from("O-MODIFY-REJECTED"));
2083 let venue_order_id = Some(VenueOrderId::from("V-001"));
2084 let err = bitmex_api_error();
2085
2086 handle_modify_failure(&ModifyFailure {
2087 err: &err,
2088 emitter: &emitter,
2089 clock: get_atomic_clock_realtime(),
2090 strategy_id: order.strategy_id(),
2091 instrument_id: order.instrument_id(),
2092 client_order_id: order.client_order_id(),
2093 venue_order_id,
2094 });
2095
2096 let events = drain_order_events(&mut rx);
2097 assert_eq!(events.len(), 1);
2098 match &events[0] {
2099 OrderEventAny::ModifyRejected(rejected) => {
2100 assert_eq!(rejected.client_order_id, order.client_order_id());
2101 assert_eq!(rejected.venue_order_id, venue_order_id);
2102 assert_eq!(
2103 rejected.reason.to_string(),
2104 "modify-order-error: BitMEX error HTTPError: Invalid price"
2105 );
2106 }
2107 event => panic!("expected OrderModifyRejected event, was {event:?}"),
2108 }
2109 }
2110
2111 #[rstest]
2112 fn test_modify_failure_network_error_emits_no_modify_rejected() {
2113 let (emitter, mut rx) = make_emitter();
2114 let order = limit_order_with_id(ClientOrderId::from("O-MODIFY-NETWORK"));
2115 let err = anyhow::Error::new(BitmexHttpError::NetworkError("timeout".to_string()));
2116
2117 handle_modify_failure(&ModifyFailure {
2118 err: &err,
2119 emitter: &emitter,
2120 clock: get_atomic_clock_realtime(),
2121 strategy_id: order.strategy_id(),
2122 instrument_id: order.instrument_id(),
2123 client_order_id: order.client_order_id(),
2124 venue_order_id: None,
2125 });
2126
2127 assert!(drain_order_events(&mut rx).is_empty());
2128 }
2129}