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