Skip to main content

nautilus_bitmex/
execution.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live execution client implementation for the BitMEX adapter.
17
18use 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    /// Creates a new [`BitmexExecutionClient`].
117    ///
118    /// # Errors
119    ///
120    /// Returns an error if the broadcaster pool sizes are invalid, API credentials are unavailable,
121    /// or either the HTTP or WebSocket client fails to construct.
122    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    /// Populates `order_identities` for an order if not already present.
273    ///
274    /// Needed for cancel/modify commands on orders loaded via reconciliation
275    /// (which bypass `submit_order` and therefore have no identity entry).
276    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        // Abort and await loop shutdown so disconnect does not block on sleep/HTTP timeout.
347        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        // Build symbol-keyed instrument map, preferring core cache then HTTP client cache
466        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                    // The WS dispatch handles all lifecycle events for tracked orders.
611                    // Forwarding the HTTP response as a report would cause the ExecEngine
612                    // to generate inferred fills that conflict with real fills from the
613                    // Execution table WS stream.
614                }
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        // Reset cancellation token so HTTP requests succeed after reconnect
718        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        // Start submitter/canceller after WS connection succeeds
729        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        // Disarm DMS before cancelling requests (needs working HTTP)
754        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    // A duplicate clOrdID can mean the original success response was lost
1222    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}