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