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