Skip to main content

nautilus_bitmex/broadcast/
canceller.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//! Cancel request broadcaster for redundant order cancellation.
17//!
18//! This module provides the [`CancelBroadcaster`] which fans out cancel requests
19//! to multiple HTTP clients in parallel for redundancy. Key design patterns:
20//!
21//! - **Dependency injection via traits**: Uses `CancelExecutor` trait to abstract
22//!   the HTTP client, enabling testing without `#[cfg(test)]` conditional compilation.
23//! - **Trait objects over generics**: Uses `Arc<dyn CancelExecutor>` to avoid
24//!   generic type parameters on the public API (simpler Python FFI).
25//! - **Short-circuit on first success**: Aborts remaining requests once any client
26//!   succeeds, minimizing latency.
27//! - **Idempotent success handling**: Recognizes "already cancelled" responses as
28//!   successful outcomes.
29
30// TODO: Replace boxed futures in `CancelExecutor` once stable async trait object support
31// lands so we can drop the per-call heap allocation
32
33use std::{
34    fmt::Debug,
35    future::Future,
36    pin::Pin,
37    sync::{
38        Arc,
39        atomic::{AtomicBool, AtomicU64, Ordering},
40    },
41    time::Duration,
42};
43
44use futures_util::future;
45use nautilus_common::live::get_runtime;
46use nautilus_core::string::secret::SecretString;
47use nautilus_model::{
48    enums::OrderSide,
49    identifiers::{ClientOrderId, InstrumentId, VenueOrderId},
50    instruments::InstrumentAny,
51    reports::OrderStatusReport,
52};
53use tokio::{sync::RwLock, task::JoinHandle, time::interval};
54
55use crate::{
56    common::{consts::BITMEX_HTTP_TESTNET_URL, enums::BitmexEnvironment},
57    config::validate_broadcaster_pool_size,
58    http::client::BitmexHttpClient,
59};
60
61const IDEMPOTENT_ALREADY_CANCELED: &str = "AlreadyCanceled";
62const IDEMPOTENT_ORDER_NOT_FOUND: &str = "orderID not found";
63const IDEMPOTENT_UNABLE_DUE_TO_STATE: &str = "Unable to cancel order due to existing state";
64
65/// Trait for order cancellation operations.
66///
67/// This trait abstracts the execution layer to enable dependency injection and testing
68/// without conditional compilation. The broadcaster holds executors as `Arc<dyn CancelExecutor>`
69/// to avoid generic type parameters that would complicate the Python FFI boundary.
70///
71/// # Thread Safety
72///
73/// All methods must be safe to call concurrently from multiple threads. Implementations
74/// should use interior mutability (e.g., `Arc<Mutex<T>>`) if mutable state is required.
75///
76/// # Error Handling
77///
78/// Methods return `anyhow::Result` for flexibility. Implementers should provide
79/// meaningful error messages that can be logged and tracked by the broadcaster.
80///
81/// # Implementation Note
82///
83/// This trait does not require `Clone` because executors are wrapped in `Arc` at the
84/// `TransportClient` level. This allows `BitmexHttpClient` (which doesn't implement
85/// `Clone`) to be used without modification.
86trait CancelExecutor: Send + Sync {
87    /// Adds an instrument for caching.
88    fn add_instrument(&self, instrument: InstrumentAny);
89
90    /// Performs a health check on the executor.
91    fn health_check(&self) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send + '_>>;
92
93    /// Cancels a single order.
94    fn cancel_order(
95        &self,
96        instrument_id: InstrumentId,
97        client_order_id: Option<ClientOrderId>,
98        venue_order_id: Option<VenueOrderId>,
99    ) -> Pin<Box<dyn Future<Output = anyhow::Result<OrderStatusReport>> + Send + '_>>;
100
101    /// Cancels multiple orders.
102    fn cancel_orders(
103        &self,
104        instrument_id: InstrumentId,
105        client_order_ids: Option<Vec<ClientOrderId>>,
106        venue_order_ids: Option<Vec<VenueOrderId>>,
107    ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>>;
108
109    /// Cancels all orders for an instrument.
110    fn cancel_all_orders(
111        &self,
112        instrument_id: InstrumentId,
113        order_side: Option<OrderSide>,
114    ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>>;
115}
116
117impl CancelExecutor for BitmexHttpClient {
118    fn add_instrument(&self, instrument: InstrumentAny) {
119        Self::cache_instrument(self, instrument);
120    }
121
122    fn health_check(&self) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send + '_>> {
123        Box::pin(async move {
124            Self::get_server_time(self)
125                .await
126                .map(|_| ())
127                .map_err(|e| anyhow::anyhow!("{e}"))
128        })
129    }
130
131    fn cancel_order(
132        &self,
133        instrument_id: InstrumentId,
134        client_order_id: Option<ClientOrderId>,
135        venue_order_id: Option<VenueOrderId>,
136    ) -> Pin<Box<dyn Future<Output = anyhow::Result<OrderStatusReport>> + Send + '_>> {
137        Box::pin(async move {
138            Self::cancel_order(self, instrument_id, client_order_id, venue_order_id).await
139        })
140    }
141
142    fn cancel_orders(
143        &self,
144        instrument_id: InstrumentId,
145        client_order_ids: Option<Vec<ClientOrderId>>,
146        venue_order_ids: Option<Vec<VenueOrderId>>,
147    ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>> {
148        Box::pin(async move {
149            Self::cancel_orders(self, instrument_id, client_order_ids, venue_order_ids).await
150        })
151    }
152
153    fn cancel_all_orders(
154        &self,
155        instrument_id: InstrumentId,
156        order_side: Option<OrderSide>,
157    ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>> {
158        Box::pin(async move { Self::cancel_all_orders(self, instrument_id, order_side).await })
159    }
160}
161
162/// Configuration for the cancel broadcaster.
163#[derive(Debug, Clone)]
164pub struct CancelBroadcasterConfig {
165    /// Number of HTTP clients in the pool (range `[1, 16]`).
166    pub pool_size: usize,
167    /// BitMEX API key (None will source from environment).
168    pub api_key: Option<SecretString>,
169    /// BitMEX API secret (None will source from environment).
170    pub api_secret: Option<SecretString>,
171    /// Base URL for BitMEX HTTP API.
172    pub base_url: Option<String>,
173    /// BitMEX environment (mainnet or testnet).
174    pub environment: BitmexEnvironment,
175    /// Timeout in seconds for HTTP requests.
176    pub timeout_secs: u64,
177    /// Maximum number of retry attempts for failed requests.
178    pub max_retries: u32,
179    /// Initial delay in milliseconds between retry attempts.
180    pub retry_delay_ms: u64,
181    /// Maximum delay in milliseconds between retry attempts.
182    pub retry_delay_max_ms: u64,
183    /// Expiration window in milliseconds for signed requests.
184    pub recv_window_ms: u64,
185    /// Maximum REST burst rate (requests per second).
186    pub max_requests_per_second: u32,
187    /// Maximum REST rolling rate (requests per minute).
188    pub max_requests_per_minute: u32,
189    /// Interval in seconds between health check pings.
190    pub health_check_interval_secs: u64,
191    /// Timeout in seconds for health check requests.
192    pub health_check_timeout_secs: u64,
193    /// Substrings to identify expected cancel rejections for debug-level logging.
194    pub expected_reject_patterns: Vec<String>,
195    /// Substrings to identify idempotent success (order already cancelled/not found).
196    pub idempotent_success_patterns: Vec<String>,
197    /// Optional list of proxy URLs for path diversity.
198    ///
199    /// Each transport instance uses the proxy at its index. If the list is shorter
200    /// than pool_size, remaining transports will use no proxy. If longer, extra proxies
201    /// are ignored.
202    pub proxy_urls: Vec<Option<SecretString>>,
203}
204
205impl Default for CancelBroadcasterConfig {
206    fn default() -> Self {
207        Self {
208            pool_size: 2,
209            api_key: None,
210            api_secret: None,
211            base_url: None,
212            environment: BitmexEnvironment::Mainnet,
213            timeout_secs: 60,
214            max_retries: 3,
215            retry_delay_ms: 1_000,
216            retry_delay_max_ms: 5_000,
217            recv_window_ms: 10_000,
218            max_requests_per_second: 10,
219            max_requests_per_minute: 120,
220            health_check_interval_secs: 30,
221            health_check_timeout_secs: 5,
222            expected_reject_patterns: vec![
223                "Order had execInst of ParticipateDoNotInitiate".to_string(),
224            ],
225            idempotent_success_patterns: vec![
226                IDEMPOTENT_ALREADY_CANCELED.to_string(),
227                IDEMPOTENT_ORDER_NOT_FOUND.to_string(),
228                IDEMPOTENT_UNABLE_DUE_TO_STATE.to_string(),
229            ],
230            proxy_urls: vec![],
231        }
232    }
233}
234
235/// Transport client wrapper with health monitoring.
236#[derive(Clone)]
237struct TransportClient {
238    /// Executor wrapped in Arc to enable cloning without requiring Clone on CancelExecutor.
239    ///
240    /// BitmexHttpClient doesn't implement Clone, so we use reference counting to share
241    /// the executor across multiple TransportClient clones.
242    executor: Arc<dyn CancelExecutor>,
243    client_id: String,
244    healthy: Arc<AtomicBool>,
245    cancel_count: Arc<AtomicU64>,
246    error_count: Arc<AtomicU64>,
247}
248
249impl Debug for TransportClient {
250    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
251        f.debug_struct(stringify!(TransportClient))
252            .field("client_id", &self.client_id)
253            .field("healthy", &self.healthy)
254            .field("cancel_count", &self.cancel_count)
255            .field("error_count", &self.error_count)
256            .finish()
257    }
258}
259
260impl TransportClient {
261    fn new<E: CancelExecutor + 'static>(executor: E, client_id: String) -> Self {
262        Self {
263            executor: Arc::new(executor),
264            client_id,
265            healthy: Arc::new(AtomicBool::new(true)),
266            cancel_count: Arc::new(AtomicU64::new(0)),
267            error_count: Arc::new(AtomicU64::new(0)),
268        }
269    }
270
271    fn is_healthy(&self) -> bool {
272        self.healthy.load(Ordering::Relaxed)
273    }
274
275    fn mark_healthy(&self) {
276        self.healthy.store(true, Ordering::Relaxed);
277    }
278
279    fn mark_unhealthy(&self) {
280        self.healthy.store(false, Ordering::Relaxed);
281    }
282
283    fn get_cancel_count(&self) -> u64 {
284        self.cancel_count.load(Ordering::Relaxed)
285    }
286
287    fn get_error_count(&self) -> u64 {
288        self.error_count.load(Ordering::Relaxed)
289    }
290
291    async fn health_check(&self, timeout_secs: u64) -> bool {
292        match tokio::time::timeout(
293            Duration::from_secs(timeout_secs),
294            self.executor.health_check(),
295        )
296        .await
297        {
298            Ok(Ok(())) => {
299                self.mark_healthy();
300                true
301            }
302            Ok(Err(e)) => {
303                log::warn!("Health check failed for client {}: {e:?}", self.client_id);
304                self.mark_unhealthy();
305                false
306            }
307            Err(_) => {
308                log::warn!("Health check timeout for client {}", self.client_id);
309                self.mark_unhealthy();
310                false
311            }
312        }
313    }
314
315    async fn cancel_order(
316        &self,
317        instrument_id: InstrumentId,
318        client_order_id: Option<ClientOrderId>,
319        venue_order_id: Option<VenueOrderId>,
320    ) -> anyhow::Result<OrderStatusReport> {
321        self.cancel_count.fetch_add(1, Ordering::Relaxed);
322
323        match self
324            .executor
325            .cancel_order(instrument_id, client_order_id, venue_order_id)
326            .await
327        {
328            Ok(report) => {
329                self.mark_healthy();
330                Ok(report)
331            }
332            Err(e) => {
333                self.error_count.fetch_add(1, Ordering::Relaxed);
334                Err(e)
335            }
336        }
337    }
338}
339
340/// Broadcasts cancel requests to multiple HTTP clients for redundancy.
341///
342/// This broadcaster fans out cancel requests to multiple pre-warmed HTTP clients
343/// in parallel, short-circuits when the first successful acknowledgement is received,
344/// and handles expected rejection patterns with appropriate log levels.
345///
346/// The client pool must contain `[1, 16]` clients.
347#[cfg_attr(feature = "python", pyo3::pyclass)]
348#[cfg_attr(
349    feature = "python",
350    pyo3_stub_gen::derive::gen_stub_pyclass(module = "nautilus_trader.adapters.bitmex")
351)]
352#[derive(Debug)]
353pub struct CancelBroadcaster {
354    config: CancelBroadcasterConfig,
355    transports: Arc<[TransportClient]>,
356    health_check_task: Arc<RwLock<Option<JoinHandle<()>>>>,
357    running: Arc<AtomicBool>,
358    total_cancels: Arc<AtomicU64>,
359    successful_cancels: Arc<AtomicU64>,
360    failed_cancels: Arc<AtomicU64>,
361    expected_rejects: Arc<AtomicU64>,
362    idempotent_successes: Arc<AtomicU64>,
363}
364
365impl CancelBroadcaster {
366    /// Creates a new [`CancelBroadcaster`] with internal HTTP client pool.
367    ///
368    /// # Errors
369    ///
370    /// Returns an error if `pool_size` is outside `[1, 16]` or any HTTP client fails to initialize.
371    pub fn new(config: CancelBroadcasterConfig) -> anyhow::Result<Self> {
372        validate_broadcaster_pool_size(config.pool_size, "pool_size")?;
373        let mut transports = Vec::with_capacity(config.pool_size);
374
375        let base_url = match config.environment {
376            BitmexEnvironment::Testnet if config.base_url.is_none() => {
377                Some(BITMEX_HTTP_TESTNET_URL.to_string())
378            }
379            _ => config.base_url.clone(),
380        };
381
382        for i in 0..config.pool_size {
383            // Assign proxy from config list, or None if index exceeds list length
384            let proxy_url = config
385                .proxy_urls
386                .get(i)
387                .and_then(|value| value.as_ref())
388                .map(|value| value.expose_secret().to_owned());
389
390            let client = BitmexHttpClient::with_credentials(
391                config.api_key.clone().map(SecretString::into_inner),
392                config.api_secret.clone().map(SecretString::into_inner),
393                base_url.clone(),
394                config.timeout_secs,
395                config.max_retries,
396                config.retry_delay_ms,
397                config.retry_delay_max_ms,
398                config.recv_window_ms,
399                config.max_requests_per_second,
400                config.max_requests_per_minute,
401                proxy_url,
402            )
403            .map_err(|e| anyhow::anyhow!("Failed to create HTTP client {i}: {e}"))?;
404
405            transports.push(TransportClient::new(client, format!("bitmex-cancel-{i}")));
406        }
407
408        Ok(Self {
409            config,
410            transports: Arc::from(transports),
411            health_check_task: Arc::new(RwLock::new(None)),
412            running: Arc::new(AtomicBool::new(false)),
413            total_cancels: Arc::new(AtomicU64::new(0)),
414            successful_cancels: Arc::new(AtomicU64::new(0)),
415            failed_cancels: Arc::new(AtomicU64::new(0)),
416            expected_rejects: Arc::new(AtomicU64::new(0)),
417            idempotent_successes: Arc::new(AtomicU64::new(0)),
418        })
419    }
420
421    /// Starts the broadcaster and health check loop.
422    ///
423    /// # Errors
424    ///
425    /// Returns an error if the broadcaster is already running.
426    pub async fn start(&self) -> anyhow::Result<()> {
427        if self.running.load(Ordering::Relaxed) {
428            return Ok(());
429        }
430
431        self.running.store(true, Ordering::Relaxed);
432
433        // Initial health check for all clients
434        self.run_health_checks().await;
435
436        // Start periodic health check task
437        let transports = Arc::clone(&self.transports);
438        let running = Arc::clone(&self.running);
439        let interval_secs = self.config.health_check_interval_secs;
440        let timeout_secs = self.config.health_check_timeout_secs;
441
442        let task = get_runtime().spawn(async move {
443            let mut ticker = interval(Duration::from_secs(interval_secs));
444            ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
445
446            loop {
447                ticker.tick().await;
448
449                if !running.load(Ordering::Relaxed) {
450                    break;
451                }
452
453                let tasks: Vec<_> = transports
454                    .iter()
455                    .map(|t| t.health_check(timeout_secs))
456                    .collect();
457
458                let results = future::join_all(tasks).await;
459                let healthy_count = results.iter().filter(|&&r| r).count();
460
461                log::debug!(
462                    "Health check complete: {}/{} clients healthy",
463                    healthy_count,
464                    results.len()
465                );
466            }
467        });
468
469        *self.health_check_task.write().await = Some(task);
470
471        log::debug!(
472            "CancelBroadcaster started with {} clients",
473            self.transports.len()
474        );
475
476        Ok(())
477    }
478
479    /// Stops the broadcaster and health check loop.
480    pub async fn stop(&self) {
481        if !self.running.load(Ordering::Relaxed) {
482            return;
483        }
484
485        self.running.store(false, Ordering::Relaxed);
486
487        if let Some(task) = self.health_check_task.write().await.take() {
488            task.abort();
489        }
490
491        log::debug!("CancelBroadcaster stopped");
492    }
493
494    async fn run_health_checks(&self) {
495        let tasks: Vec<_> = self
496            .transports
497            .iter()
498            .map(|t| t.health_check(self.config.health_check_timeout_secs))
499            .collect();
500
501        let results = future::join_all(tasks).await;
502        let healthy_count = results.iter().filter(|&&r| r).count();
503
504        log::debug!(
505            "Health check complete: {}/{} clients healthy",
506            healthy_count,
507            results.len()
508        );
509    }
510
511    fn is_expected_reject(&self, error_message: &str) -> bool {
512        self.config
513            .expected_reject_patterns
514            .iter()
515            .any(|pattern| error_message.contains(pattern))
516    }
517
518    fn is_idempotent_success(&self, error_message: &str) -> bool {
519        self.config
520            .idempotent_success_patterns
521            .iter()
522            .any(|pattern| error_message.contains(pattern))
523    }
524
525    /// Processes cancel request results, handling success, idempotent success, and failures.
526    ///
527    /// This method consolidates the error handling loop used across all broadcast methods.
528    async fn process_cancel_results<T>(
529        &self,
530        mut handles: Vec<JoinHandle<(String, anyhow::Result<T>)>>,
531        idempotent_result: impl FnOnce() -> anyhow::Result<T>,
532        operation: &str,
533        params: String,
534        idempotent_reason: &str,
535    ) -> anyhow::Result<T>
536    where
537        T: Send + 'static,
538    {
539        let mut errors = Vec::new();
540
541        while !handles.is_empty() {
542            let current_handles = std::mem::take(&mut handles);
543            let (result, _idx, remaining) = future::select_all(current_handles).await;
544            handles = remaining.into_iter().collect();
545
546            match result {
547                Ok((client_id, Ok(result))) => {
548                    // First success - abort remaining handles
549                    for handle in &handles {
550                        handle.abort();
551                    }
552                    self.successful_cancels.fetch_add(1, Ordering::Relaxed);
553
554                    log::debug!("{operation} broadcast succeeded [{client_id}] {params}");
555
556                    return Ok(result);
557                }
558                Ok((client_id, Err(e))) => {
559                    let error_msg = e.to_string();
560
561                    if self.is_idempotent_success(&error_msg) {
562                        // First idempotent success - abort remaining handles and return success
563                        for handle in &handles {
564                            handle.abort();
565                        }
566                        self.idempotent_successes.fetch_add(1, Ordering::Relaxed);
567
568                        log::debug!(
569                            "Idempotent success [{client_id}] - {idempotent_reason}: {error_msg} {params}",
570                        );
571
572                        return idempotent_result();
573                    }
574
575                    if self.is_expected_reject(&error_msg) {
576                        self.expected_rejects.fetch_add(1, Ordering::Relaxed);
577                        log::debug!(
578                            "Expected {} rejection [{}]: {} {}",
579                            operation.to_lowercase(),
580                            client_id,
581                            error_msg,
582                            params
583                        );
584                        errors.push(error_msg);
585                    } else {
586                        log::warn!(
587                            "{operation} request failed [{client_id}]: {error_msg} {params}"
588                        );
589                        errors.push(error_msg);
590                    }
591                }
592                Err(e) => {
593                    log::warn!("{operation} task join error: {e:?}");
594                    errors.push(format!("Task panicked: {e:?}"));
595                }
596            }
597        }
598
599        // All tasks failed
600        self.failed_cancels.fetch_add(1, Ordering::Relaxed);
601        log::error!(
602            "All {} requests failed: {errors:?} {params}",
603            operation.to_lowercase(),
604        );
605        Err(anyhow::anyhow!(
606            "All {} requests failed: {errors:?}",
607            operation.to_lowercase(),
608        ))
609    }
610
611    /// Broadcasts a single cancel request to all healthy clients in parallel.
612    ///
613    /// # Returns
614    ///
615    /// - `Ok(Some(report))` if successfully cancelled with a report.
616    /// - `Ok(None)` if the order was already cancelled (idempotent success).
617    /// - `Err` if all requests failed.
618    ///
619    /// # Errors
620    ///
621    /// Returns an error if all cancel requests fail or no healthy clients are available.
622    pub async fn broadcast_cancel(
623        &self,
624        instrument_id: InstrumentId,
625        client_order_id: Option<ClientOrderId>,
626        venue_order_id: Option<VenueOrderId>,
627    ) -> anyhow::Result<Option<OrderStatusReport>> {
628        self.total_cancels.fetch_add(1, Ordering::Relaxed);
629
630        let healthy_transports: Vec<TransportClient> = self
631            .transports
632            .iter()
633            .filter(|t| t.is_healthy())
634            .cloned()
635            .collect();
636
637        if healthy_transports.is_empty() {
638            self.failed_cancels.fetch_add(1, Ordering::Relaxed);
639            anyhow::bail!("No healthy transport clients available");
640        }
641
642        let mut handles = Vec::new();
643
644        for transport in healthy_transports {
645            let handle = get_runtime().spawn(async move {
646                let client_id = transport.client_id.clone();
647                let result = transport
648                    .cancel_order(instrument_id, client_order_id, venue_order_id)
649                    .await
650                    .map(Some); // Wrap success in Some for Option<OrderStatusReport>
651                (client_id, result)
652            });
653            handles.push(handle);
654        }
655
656        self.process_cancel_results(
657            handles,
658            || Ok(None),
659            "Cancel",
660            format!("(client_order_id={client_order_id:?}, venue_order_id={venue_order_id:?})"),
661            "order already cancelled/not found",
662        )
663        .await
664    }
665
666    /// Broadcasts a batch cancel request to all healthy clients in parallel.
667    ///
668    /// # Errors
669    ///
670    /// Returns an error if all cancel requests fail or no healthy clients are available.
671    pub async fn broadcast_batch_cancel(
672        &self,
673        instrument_id: InstrumentId,
674        client_order_ids: Option<Vec<ClientOrderId>>,
675        venue_order_ids: Option<Vec<VenueOrderId>>,
676    ) -> anyhow::Result<Vec<OrderStatusReport>> {
677        self.total_cancels.fetch_add(1, Ordering::Relaxed);
678
679        let healthy_transports: Vec<TransportClient> = self
680            .transports
681            .iter()
682            .filter(|t| t.is_healthy())
683            .cloned()
684            .collect();
685
686        if healthy_transports.is_empty() {
687            self.failed_cancels.fetch_add(1, Ordering::Relaxed);
688            anyhow::bail!("No healthy transport clients available");
689        }
690
691        let mut handles = Vec::new();
692
693        for transport in healthy_transports {
694            let client_order_ids_clone = client_order_ids.clone();
695            let venue_order_ids_clone = venue_order_ids.clone();
696            let handle = get_runtime().spawn(async move {
697                let client_id = transport.client_id.clone();
698                let result = transport
699                    .executor
700                    .cancel_orders(instrument_id, client_order_ids_clone, venue_order_ids_clone)
701                    .await;
702                (client_id, result)
703            });
704            handles.push(handle);
705        }
706
707        self.process_cancel_results(
708            handles,
709            || Ok(Vec::new()),
710            "Batch cancel",
711            format!("(client_order_ids={client_order_ids:?}, venue_order_ids={venue_order_ids:?})"),
712            "orders already cancelled/not found",
713        )
714        .await
715    }
716
717    /// Broadcasts a cancel all request to all healthy clients in parallel.
718    ///
719    /// # Errors
720    ///
721    /// Returns an error if all cancel requests fail or no healthy clients are available.
722    pub async fn broadcast_cancel_all(
723        &self,
724        instrument_id: InstrumentId,
725        order_side: Option<OrderSide>,
726    ) -> anyhow::Result<Vec<OrderStatusReport>> {
727        self.total_cancels.fetch_add(1, Ordering::Relaxed);
728
729        let healthy_transports: Vec<TransportClient> = self
730            .transports
731            .iter()
732            .filter(|t| t.is_healthy())
733            .cloned()
734            .collect();
735
736        if healthy_transports.is_empty() {
737            self.failed_cancels.fetch_add(1, Ordering::Relaxed);
738            anyhow::bail!("No healthy transport clients available");
739        }
740
741        let mut handles = Vec::new();
742
743        for transport in healthy_transports {
744            let handle = get_runtime().spawn(async move {
745                let client_id = transport.client_id.clone();
746                let result = transport
747                    .executor
748                    .cancel_all_orders(instrument_id, order_side)
749                    .await;
750                (client_id, result)
751            });
752            handles.push(handle);
753        }
754
755        self.process_cancel_results(
756            handles,
757            || Ok(Vec::new()),
758            "Cancel all",
759            format!("(instrument_id={instrument_id}, order_side={order_side:?})"),
760            "no orders to cancel",
761        )
762        .await
763    }
764
765    /// Gets broadcaster metrics.
766    pub fn get_metrics(&self) -> BroadcasterMetrics {
767        let healthy_clients = self.transports.iter().filter(|t| t.is_healthy()).count();
768        let total_clients = self.transports.len();
769
770        BroadcasterMetrics {
771            total_cancels: self.total_cancels.load(Ordering::Relaxed),
772            successful_cancels: self.successful_cancels.load(Ordering::Relaxed),
773            failed_cancels: self.failed_cancels.load(Ordering::Relaxed),
774            expected_rejects: self.expected_rejects.load(Ordering::Relaxed),
775            idempotent_successes: self.idempotent_successes.load(Ordering::Relaxed),
776            healthy_clients,
777            total_clients,
778        }
779    }
780
781    /// Gets broadcaster metrics (async version for use within async context).
782    pub async fn get_metrics_async(&self) -> BroadcasterMetrics {
783        self.get_metrics()
784    }
785
786    /// Gets per-client statistics.
787    pub fn get_client_stats(&self) -> Vec<ClientStats> {
788        self.transports
789            .iter()
790            .map(|t| ClientStats {
791                client_id: t.client_id.clone(),
792                healthy: t.is_healthy(),
793                cancel_count: t.get_cancel_count(),
794                error_count: t.get_error_count(),
795            })
796            .collect()
797    }
798
799    /// Gets per-client statistics (async version for use within async context).
800    pub async fn get_client_stats_async(&self) -> Vec<ClientStats> {
801        self.get_client_stats()
802    }
803
804    /// Caches an instrument in all HTTP clients in the pool.
805    pub fn cache_instrument(&self, instrument: &InstrumentAny) {
806        for transport in self.transports.iter() {
807            transport.executor.add_instrument(instrument.clone());
808        }
809    }
810
811    #[must_use]
812    pub fn clone_for_async(&self) -> Self {
813        Self {
814            config: self.config.clone(),
815            transports: Arc::clone(&self.transports),
816            health_check_task: Arc::clone(&self.health_check_task),
817            running: Arc::clone(&self.running),
818            total_cancels: Arc::clone(&self.total_cancels),
819            successful_cancels: Arc::clone(&self.successful_cancels),
820            failed_cancels: Arc::clone(&self.failed_cancels),
821            expected_rejects: Arc::clone(&self.expected_rejects),
822            idempotent_successes: Arc::clone(&self.idempotent_successes),
823        }
824    }
825
826    #[cfg(test)]
827    fn new_with_transports(
828        config: CancelBroadcasterConfig,
829        transports: Vec<TransportClient>,
830    ) -> Self {
831        Self {
832            config,
833            transports: Arc::from(transports),
834            health_check_task: Arc::new(RwLock::new(None)),
835            running: Arc::new(AtomicBool::new(false)),
836            total_cancels: Arc::new(AtomicU64::new(0)),
837            successful_cancels: Arc::new(AtomicU64::new(0)),
838            failed_cancels: Arc::new(AtomicU64::new(0)),
839            expected_rejects: Arc::new(AtomicU64::new(0)),
840            idempotent_successes: Arc::new(AtomicU64::new(0)),
841        }
842    }
843}
844
845/// Broadcaster metrics snapshot.
846#[derive(Debug, Clone)]
847pub struct BroadcasterMetrics {
848    pub total_cancels: u64,
849    pub successful_cancels: u64,
850    pub failed_cancels: u64,
851    pub expected_rejects: u64,
852    pub idempotent_successes: u64,
853    pub healthy_clients: usize,
854    pub total_clients: usize,
855}
856
857/// Per-client statistics.
858#[derive(Debug, Clone)]
859pub struct ClientStats {
860    pub client_id: String,
861    pub healthy: bool,
862    pub cancel_count: u64,
863    pub error_count: u64,
864}
865
866#[cfg(test)]
867mod tests {
868    use std::{str::FromStr, sync::atomic::Ordering, time::Duration};
869
870    use nautilus_core::UUID4;
871    use nautilus_model::{
872        enums::{OrderSide, OrderStatus, OrderType, TimeInForce},
873        identifiers::{AccountId, ClientOrderId, InstrumentId, VenueOrderId},
874        reports::OrderStatusReport,
875        types::{Price, Quantity},
876    };
877    use rstest::rstest;
878
879    use super::*;
880
881    #[rstest]
882    fn test_config_debug_redacts_credentials() {
883        let config = CancelBroadcasterConfig {
884            api_key: Some("cancel-key-sentinel".into()),
885            api_secret: Some("cancel-secret-sentinel".into()),
886            proxy_urls: vec![Some("http://cancel-user:cancel-password@localhost".into())],
887            ..Default::default()
888        };
889
890        let debug = format!("{config:?}");
891
892        assert!(!debug.contains("cancel-key-sentinel"));
893        assert!(!debug.contains("cancel-secret-sentinel"));
894        assert!(!debug.contains("cancel-password"));
895    }
896
897    /// Mock executor for testing.
898    #[derive(Clone)]
899    #[expect(clippy::type_complexity)]
900    struct MockExecutor {
901        handler: Arc<
902            dyn Fn(
903                    InstrumentId,
904                    Option<ClientOrderId>,
905                    Option<VenueOrderId>,
906                )
907                    -> Pin<Box<dyn Future<Output = anyhow::Result<OrderStatusReport>> + Send>>
908                + Send
909                + Sync,
910        >,
911    }
912
913    impl MockExecutor {
914        fn new<F, Fut>(handler: F) -> Self
915        where
916            F: Fn(InstrumentId, Option<ClientOrderId>, Option<VenueOrderId>) -> Fut
917                + Send
918                + Sync
919                + 'static,
920            Fut: Future<Output = anyhow::Result<OrderStatusReport>> + Send + 'static,
921        {
922            Self {
923                handler: Arc::new(move |id, cid, vid| Box::pin(handler(id, cid, vid))),
924            }
925        }
926    }
927
928    impl CancelExecutor for MockExecutor {
929        fn health_check(&self) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send + '_>> {
930            Box::pin(async { Ok(()) })
931        }
932
933        fn cancel_order(
934            &self,
935            instrument_id: InstrumentId,
936            client_order_id: Option<ClientOrderId>,
937            venue_order_id: Option<VenueOrderId>,
938        ) -> Pin<Box<dyn Future<Output = anyhow::Result<OrderStatusReport>> + Send + '_>> {
939            (self.handler)(instrument_id, client_order_id, venue_order_id)
940        }
941
942        fn cancel_orders(
943            &self,
944            _instrument_id: InstrumentId,
945            _client_order_ids: Option<Vec<ClientOrderId>>,
946            _venue_order_ids: Option<Vec<VenueOrderId>>,
947        ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>>
948        {
949            Box::pin(async { Ok(Vec::new()) })
950        }
951
952        fn cancel_all_orders(
953            &self,
954            instrument_id: InstrumentId,
955            _order_side: Option<OrderSide>,
956        ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>>
957        {
958            // Try to get result from the single-order handler to propagate errors
959            let handler = Arc::clone(&self.handler);
960            Box::pin(async move {
961                // Call the handler to check if it would fail
962                let result = handler(instrument_id, None, None).await;
963                match result {
964                    Ok(_) => Ok(Vec::new()),
965                    Err(e) => Err(e),
966                }
967            })
968        }
969
970        fn add_instrument(&self, _instrument: InstrumentAny) {
971            // No-op for mock
972        }
973    }
974
975    fn create_test_report(venue_order_id: &str) -> OrderStatusReport {
976        OrderStatusReport {
977            account_id: AccountId::from("BITMEX-001"),
978            instrument_id: InstrumentId::from_str("XBTUSD.BITMEX").unwrap(),
979            venue_order_id: VenueOrderId::from(venue_order_id),
980            order_side: OrderSide::Buy.into(),
981            order_type: OrderType::Limit,
982            time_in_force: TimeInForce::Gtc,
983            order_status: OrderStatus::Canceled,
984            price: Some(Price::new(50000.0, 2)),
985            quantity: Quantity::new(100.0, 0),
986            filled_qty: Quantity::new(0.0, 0),
987            report_id: UUID4::new(),
988            ts_accepted: 0.into(),
989            ts_last: 0.into(),
990            ts_init: 0.into(),
991            client_order_id: None,
992            avg_px: None,
993            activation_price: None,
994            trigger_price: None,
995            trigger_type: None,
996            contingency_type: None,
997            expire_time: None,
998            order_list_id: None,
999            venue_position_id: None,
1000            linked_order_ids: None,
1001            parent_order_id: None,
1002            display_qty: None,
1003            limit_offset: None,
1004            trailing_offset: None,
1005            trailing_offset_type: None,
1006            post_only: false,
1007            reduce_only: false,
1008            cancel_reason: None,
1009            ts_triggered: None,
1010        }
1011    }
1012
1013    fn create_stub_transport<F, Fut>(client_id: &str, handler: F) -> TransportClient
1014    where
1015        F: Fn(InstrumentId, Option<ClientOrderId>, Option<VenueOrderId>) -> Fut
1016            + Send
1017            + Sync
1018            + 'static,
1019        Fut: Future<Output = anyhow::Result<OrderStatusReport>> + Send + 'static,
1020    {
1021        let executor = MockExecutor::new(handler);
1022        TransportClient::new(executor, client_id.to_string())
1023    }
1024
1025    #[tokio::test]
1026    async fn test_broadcast_cancel_immediate_success() {
1027        let report = create_test_report("ORDER-1");
1028        let report_clone = report.clone();
1029
1030        let transports = vec![
1031            create_stub_transport("client-0", move |_, _, _| {
1032                let report = report_clone.clone();
1033                async move { Ok(report) }
1034            }),
1035            create_stub_transport("client-1", |_, _, _| async {
1036                tokio::time::sleep(Duration::from_secs(10)).await;
1037                anyhow::bail!("Should be aborted")
1038            }),
1039        ];
1040
1041        let config = CancelBroadcasterConfig::default();
1042        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1043
1044        let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1045        let result = broadcaster
1046            .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-123")), None)
1047            .await;
1048
1049        assert!(result.is_ok());
1050        let returned_report = result.unwrap();
1051        assert!(returned_report.is_some());
1052        assert_eq!(
1053            returned_report.unwrap().venue_order_id,
1054            report.venue_order_id
1055        );
1056
1057        let metrics = broadcaster.get_metrics_async().await;
1058        assert_eq!(metrics.successful_cancels, 1);
1059        assert_eq!(metrics.failed_cancels, 0);
1060        assert_eq!(metrics.total_cancels, 1);
1061    }
1062
1063    #[tokio::test]
1064    async fn test_broadcast_cancel_idempotent_success() {
1065        let transports = vec![
1066            create_stub_transport("client-0", |_, _, _| async {
1067                anyhow::bail!("AlreadyCanceled")
1068            }),
1069            create_stub_transport("client-1", |_, _, _| async {
1070                tokio::time::sleep(Duration::from_secs(10)).await;
1071                anyhow::bail!("Should be aborted")
1072            }),
1073        ];
1074
1075        let config = CancelBroadcasterConfig::default();
1076        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1077
1078        let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1079        let result = broadcaster
1080            .broadcast_cancel(instrument_id, None, Some(VenueOrderId::from("12345")))
1081            .await;
1082
1083        assert!(result.is_ok());
1084        assert!(result.unwrap().is_none());
1085
1086        let metrics = broadcaster.get_metrics_async().await;
1087        assert_eq!(metrics.idempotent_successes, 1);
1088        assert_eq!(metrics.successful_cancels, 0);
1089        assert_eq!(metrics.failed_cancels, 0);
1090    }
1091
1092    #[tokio::test]
1093    async fn test_broadcast_cancel_mixed_idempotent_and_failure() {
1094        let transports = vec![
1095            create_stub_transport("client-0", |_, _, _| async {
1096                anyhow::bail!("502 Bad Gateway")
1097            }),
1098            create_stub_transport("client-1", |_, _, _| async {
1099                anyhow::bail!("orderID not found")
1100            }),
1101        ];
1102
1103        let config = CancelBroadcasterConfig::default();
1104        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1105
1106        let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1107        let result = broadcaster
1108            .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-456")), None)
1109            .await;
1110
1111        assert!(result.is_ok());
1112        assert!(result.unwrap().is_none());
1113
1114        let metrics = broadcaster.get_metrics_async().await;
1115        assert_eq!(metrics.idempotent_successes, 1);
1116        assert_eq!(metrics.failed_cancels, 0);
1117    }
1118
1119    #[tokio::test]
1120    async fn test_broadcast_cancel_all_failures() {
1121        let transports = vec![
1122            create_stub_transport("client-0", |_, _, _| async {
1123                anyhow::bail!("502 Bad Gateway")
1124            }),
1125            create_stub_transport("client-1", |_, _, _| async {
1126                anyhow::bail!("Connection refused")
1127            }),
1128        ];
1129
1130        let config = CancelBroadcasterConfig::default();
1131        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1132
1133        let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1134        let result = broadcaster.broadcast_cancel_all(instrument_id, None).await;
1135
1136        assert!(result.is_err());
1137        assert!(
1138            result
1139                .unwrap_err()
1140                .to_string()
1141                .contains("All cancel all requests failed")
1142        );
1143
1144        let metrics = broadcaster.get_metrics_async().await;
1145        assert_eq!(metrics.failed_cancels, 1);
1146        assert_eq!(metrics.successful_cancels, 0);
1147        assert_eq!(metrics.idempotent_successes, 0);
1148    }
1149
1150    #[tokio::test]
1151    async fn test_broadcast_cancel_no_healthy_clients() {
1152        let transport = create_stub_transport("client-0", |_, _, _| async {
1153            Ok(create_test_report("ORDER-1"))
1154        });
1155        transport.healthy.store(false, Ordering::Relaxed);
1156
1157        let config = CancelBroadcasterConfig::default();
1158        let broadcaster = CancelBroadcaster::new_with_transports(config, vec![transport]);
1159
1160        let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1161        let result = broadcaster
1162            .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-789")), None)
1163            .await;
1164
1165        assert!(result.is_err());
1166        assert!(
1167            result
1168                .unwrap_err()
1169                .to_string()
1170                .contains("No healthy transport clients available")
1171        );
1172
1173        let metrics = broadcaster.get_metrics_async().await;
1174        assert_eq!(metrics.failed_cancels, 1);
1175    }
1176
1177    #[tokio::test]
1178    async fn test_broadcast_cancel_metrics_increment() {
1179        let report1 = create_test_report("ORDER-1");
1180        let report1_clone = report1.clone();
1181        let report2 = create_test_report("ORDER-2");
1182        let report2_clone = report2.clone();
1183
1184        let transports = vec![
1185            create_stub_transport("client-0", move |_, _, _| {
1186                let report = report1_clone.clone();
1187                async move { Ok(report) }
1188            }),
1189            create_stub_transport("client-1", move |_, _, _| {
1190                let report = report2_clone.clone();
1191                async move { Ok(report) }
1192            }),
1193        ];
1194
1195        let config = CancelBroadcasterConfig::default();
1196        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1197
1198        let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1199
1200        let _ = broadcaster
1201            .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-1")), None)
1202            .await;
1203
1204        let _ = broadcaster
1205            .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-2")), None)
1206            .await;
1207
1208        let metrics = broadcaster.get_metrics_async().await;
1209        assert_eq!(metrics.total_cancels, 2);
1210        assert_eq!(metrics.successful_cancels, 2);
1211    }
1212
1213    #[tokio::test]
1214    async fn test_broadcast_cancel_expected_reject_pattern() {
1215        let transports = vec![
1216            create_stub_transport("client-0", |_, _, _| async {
1217                anyhow::bail!("Order had execInst of ParticipateDoNotInitiate")
1218            }),
1219            create_stub_transport("client-1", |_, _, _| async {
1220                anyhow::bail!("Order had execInst of ParticipateDoNotInitiate")
1221            }),
1222        ];
1223
1224        let config = CancelBroadcasterConfig::default();
1225        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1226
1227        let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1228        let result = broadcaster
1229            .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-PDI")), None)
1230            .await;
1231
1232        assert!(result.is_err());
1233
1234        let metrics = broadcaster.get_metrics_async().await;
1235        assert_eq!(metrics.expected_rejects, 2);
1236        assert_eq!(metrics.failed_cancels, 1);
1237    }
1238
1239    #[tokio::test]
1240    async fn test_broadcaster_creation_with_pool() {
1241        let transports = vec![
1242            create_stub_transport("client-0", |_, _, _| async {
1243                Ok(create_test_report("ORDER-1"))
1244            }),
1245            create_stub_transport("client-1", |_, _, _| async {
1246                Ok(create_test_report("ORDER-1"))
1247            }),
1248            create_stub_transport("client-2", |_, _, _| async {
1249                Ok(create_test_report("ORDER-1"))
1250            }),
1251        ];
1252
1253        let config = CancelBroadcasterConfig::default();
1254        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1255        let metrics = broadcaster.get_metrics_async().await;
1256
1257        assert_eq!(metrics.total_clients, 3);
1258        assert_eq!(metrics.total_cancels, 0);
1259        assert_eq!(metrics.successful_cancels, 0);
1260        assert_eq!(metrics.failed_cancels, 0);
1261    }
1262
1263    #[tokio::test]
1264    async fn test_broadcaster_lifecycle() {
1265        let transports = vec![
1266            create_stub_transport("client-0", |_, _, _| async {
1267                Ok(create_test_report("ORDER-1"))
1268            }),
1269            create_stub_transport("client-1", |_, _, _| async {
1270                Ok(create_test_report("ORDER-1"))
1271            }),
1272        ];
1273
1274        let config = CancelBroadcasterConfig::default();
1275        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1276
1277        // Should not be running initially
1278        assert!(!broadcaster.running.load(Ordering::Relaxed));
1279
1280        // Start broadcaster
1281        let start_result = broadcaster.start().await;
1282        assert!(start_result.is_ok());
1283        assert!(broadcaster.running.load(Ordering::Relaxed));
1284
1285        // Starting again should be idempotent
1286        let start_again = broadcaster.start().await;
1287        assert!(start_again.is_ok());
1288
1289        // Stop broadcaster
1290        broadcaster.stop().await;
1291        assert!(!broadcaster.running.load(Ordering::Relaxed));
1292
1293        // Stopping again should be safe
1294        broadcaster.stop().await;
1295        assert!(!broadcaster.running.load(Ordering::Relaxed));
1296    }
1297
1298    #[tokio::test]
1299    async fn test_client_stats_collection() {
1300        let transports = vec![
1301            create_stub_transport("client-0", |_, _, _| async {
1302                Ok(create_test_report("ORDER-1"))
1303            }),
1304            create_stub_transport("client-1", |_, _, _| async {
1305                Ok(create_test_report("ORDER-1"))
1306            }),
1307        ];
1308
1309        let config = CancelBroadcasterConfig::default();
1310        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1311        let stats = broadcaster.get_client_stats_async().await;
1312
1313        assert_eq!(stats.len(), 2);
1314        assert_eq!(stats[0].client_id, "client-0");
1315        assert_eq!(stats[1].client_id, "client-1");
1316        assert!(stats[0].healthy); // Should be healthy initially
1317        assert!(stats[1].healthy);
1318        assert_eq!(stats[0].cancel_count, 0);
1319        assert_eq!(stats[1].cancel_count, 0);
1320        assert_eq!(stats[0].error_count, 0);
1321        assert_eq!(stats[1].error_count, 0);
1322    }
1323
1324    #[tokio::test]
1325    async fn test_testnet_config_sets_base_url() {
1326        let config = CancelBroadcasterConfig {
1327            pool_size: 1,
1328            api_key: Some("test_key".into()),
1329            api_secret: Some("test_secret".into()),
1330            base_url: None,
1331            environment: BitmexEnvironment::Testnet,
1332            timeout_secs: 5,
1333            max_retries: 3,
1334            retry_delay_ms: 1_000,
1335            retry_delay_max_ms: 5_000,
1336            recv_window_ms: 10_000,
1337            max_requests_per_second: 10,
1338            max_requests_per_minute: 120,
1339            health_check_interval_secs: 60,
1340            health_check_timeout_secs: 5,
1341            expected_reject_patterns: vec![],
1342            idempotent_success_patterns: vec![],
1343            proxy_urls: vec![],
1344        };
1345
1346        let broadcaster = CancelBroadcaster::new(config);
1347        assert!(broadcaster.is_ok());
1348    }
1349
1350    #[tokio::test]
1351    async fn test_constructor_honors_default_pool_size() {
1352        let config = CancelBroadcasterConfig {
1353            api_key: Some("test_key".into()),
1354            api_secret: Some("test_secret".into()),
1355            base_url: Some("http://127.0.0.1:19999".to_string()),
1356            ..Default::default()
1357        };
1358
1359        let expected_pool = config.pool_size;
1360        let broadcaster = CancelBroadcaster::new(config).unwrap();
1361        let metrics = broadcaster.get_metrics_async().await;
1362
1363        assert_eq!(metrics.total_clients, expected_pool);
1364    }
1365
1366    #[tokio::test]
1367    async fn test_constructor_accepts_maximum_pool_size() {
1368        let config = CancelBroadcasterConfig {
1369            pool_size: crate::config::MAX_BROADCASTER_POOL_SIZE,
1370            api_key: Some("test_key".into()),
1371            api_secret: Some("test_secret".into()),
1372            base_url: Some("http://127.0.0.1:19999".to_string()),
1373            ..Default::default()
1374        };
1375
1376        let broadcaster = CancelBroadcaster::new(config).unwrap();
1377        let metrics = broadcaster.get_metrics_async().await;
1378
1379        assert_eq!(
1380            metrics.total_clients,
1381            crate::config::MAX_BROADCASTER_POOL_SIZE
1382        );
1383    }
1384
1385    #[rstest]
1386    fn test_constructor_rejects_oversized_pool_before_allocation() {
1387        let config = CancelBroadcasterConfig {
1388            pool_size: crate::config::MAX_BROADCASTER_POOL_SIZE + 1,
1389            ..Default::default()
1390        };
1391
1392        let result = CancelBroadcaster::new(config);
1393        let err = result.expect_err("oversized pool must be rejected");
1394
1395        assert_eq!(
1396            err.to_string(),
1397            format!(
1398                "invalid usize for 'pool_size' not in range [1, {}], was {}",
1399                crate::config::MAX_BROADCASTER_POOL_SIZE,
1400                crate::config::MAX_BROADCASTER_POOL_SIZE + 1,
1401            ),
1402        );
1403    }
1404
1405    #[tokio::test]
1406    async fn test_default_config() {
1407        let transports = vec![
1408            create_stub_transport("client-0", |_, _, _| async {
1409                Ok(create_test_report("ORDER-1"))
1410            }),
1411            create_stub_transport("client-1", |_, _, _| async {
1412                Ok(create_test_report("ORDER-1"))
1413            }),
1414        ];
1415
1416        let config = CancelBroadcasterConfig::default();
1417        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1418        let metrics = broadcaster.get_metrics_async().await;
1419
1420        // Default pool_size is 2
1421        assert_eq!(metrics.total_clients, 2);
1422    }
1423
1424    #[tokio::test]
1425    async fn test_clone_for_async() {
1426        let transports = vec![create_stub_transport("client-0", |_, _, _| async {
1427            Ok(create_test_report("ORDER-1"))
1428        })];
1429
1430        let config = CancelBroadcasterConfig::default();
1431        let broadcaster1 = CancelBroadcaster::new_with_transports(config, transports);
1432
1433        // Increment a metric on original
1434        broadcaster1.total_cancels.fetch_add(1, Ordering::Relaxed);
1435
1436        // Clone should share the same atomic
1437        let broadcaster2 = broadcaster1.clone_for_async();
1438        let metrics2 = broadcaster2.get_metrics_async().await;
1439
1440        assert_eq!(metrics2.total_cancels, 1); // Should see the increment
1441
1442        // Modify through clone
1443        broadcaster2
1444            .successful_cancels
1445            .fetch_add(5, Ordering::Relaxed);
1446
1447        // Original should see the change
1448        let metrics1 = broadcaster1.get_metrics_async().await;
1449        assert_eq!(metrics1.successful_cancels, 5);
1450    }
1451
1452    #[tokio::test]
1453    async fn test_pattern_matching() {
1454        // Test that pattern matching works for expected rejects and idempotent successes
1455        let transports = vec![create_stub_transport("client-0", |_, _, _| async {
1456            Ok(create_test_report("ORDER-1"))
1457        })];
1458
1459        let config = CancelBroadcasterConfig {
1460            expected_reject_patterns: vec![
1461                "ParticipateDoNotInitiate".to_string(),
1462                "Close-only".to_string(),
1463            ],
1464            idempotent_success_patterns: vec![
1465                "AlreadyCanceled".to_string(),
1466                "orderID not found".to_string(),
1467                "Unable to cancel".to_string(),
1468            ],
1469            ..Default::default()
1470        };
1471
1472        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1473
1474        // Test expected reject patterns
1475        assert!(broadcaster.is_expected_reject("Order had execInst of ParticipateDoNotInitiate"));
1476        assert!(broadcaster.is_expected_reject("This is a Close-only order"));
1477        assert!(!broadcaster.is_expected_reject("Connection timeout"));
1478
1479        // Test idempotent success patterns
1480        assert!(broadcaster.is_idempotent_success("AlreadyCanceled"));
1481        assert!(broadcaster.is_idempotent_success("Error: orderID not found for this account"));
1482        assert!(broadcaster.is_idempotent_success("Unable to cancel order due to existing state"));
1483        assert!(!broadcaster.is_idempotent_success("502 Bad Gateway"));
1484    }
1485
1486    // Happy-path coverage for broadcast_batch_cancel and broadcast_cancel_all
1487    // Note: These use simplified stubs since batch/cancel-all bypass test_handler
1488    // Full HTTP mocking tested in integration tests
1489    #[tokio::test]
1490    async fn test_broadcast_batch_cancel_structure() {
1491        // Validates broadcaster structure and metric initialization
1492        let transports = vec![
1493            create_stub_transport("client-0", |_, _, _| async {
1494                Ok(create_test_report("ORDER-1"))
1495            }),
1496            create_stub_transport("client-1", |_, _, _| async {
1497                Ok(create_test_report("ORDER-1"))
1498            }),
1499        ];
1500
1501        let config = CancelBroadcasterConfig {
1502            idempotent_success_patterns: vec!["AlreadyCanceled".to_string()],
1503            ..Default::default()
1504        };
1505
1506        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1507        let metrics = broadcaster.get_metrics_async().await;
1508
1509        // Verify initial state
1510        assert_eq!(metrics.total_clients, 2);
1511        assert_eq!(metrics.total_cancels, 0);
1512        assert_eq!(metrics.successful_cancels, 0);
1513        assert_eq!(metrics.failed_cancels, 0);
1514    }
1515
1516    #[tokio::test]
1517    async fn test_broadcast_cancel_all_structure() {
1518        // Validates broadcaster structure for cancel_all operations
1519        let transports = vec![
1520            create_stub_transport("client-0", |_, _, _| async {
1521                Ok(create_test_report("ORDER-1"))
1522            }),
1523            create_stub_transport("client-1", |_, _, _| async {
1524                Ok(create_test_report("ORDER-1"))
1525            }),
1526            create_stub_transport("client-2", |_, _, _| async {
1527                Ok(create_test_report("ORDER-1"))
1528            }),
1529        ];
1530
1531        let config = CancelBroadcasterConfig {
1532            idempotent_success_patterns: vec!["orderID not found".to_string()],
1533            ..Default::default()
1534        };
1535
1536        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1537        let metrics = broadcaster.get_metrics_async().await;
1538
1539        // Verify pool size and initial metrics
1540        assert_eq!(metrics.total_clients, 3);
1541        assert_eq!(metrics.healthy_clients, 3);
1542        assert_eq!(metrics.total_cancels, 0);
1543    }
1544
1545    // Metric health tests - validates that idempotent successes don't increment failed_cancels
1546    #[tokio::test]
1547    async fn test_single_cancel_metrics_with_mixed_responses() {
1548        // Test similar to test_broadcast_cancel_mixed_idempotent_and_failure
1549        // but explicitly validates metric health
1550        let transports = vec![
1551            create_stub_transport("client-0", |_, _, _| async {
1552                anyhow::bail!("Connection timeout")
1553            }),
1554            create_stub_transport("client-1", |_, _, _| async {
1555                anyhow::bail!("AlreadyCanceled")
1556            }),
1557        ];
1558
1559        let config = CancelBroadcasterConfig::default();
1560        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1561
1562        let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1563        let result = broadcaster
1564            .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-123")), None)
1565            .await;
1566
1567        // Should succeed with idempotent
1568        assert!(result.is_ok());
1569        assert!(result.unwrap().is_none());
1570
1571        // Verify metrics: idempotent success doesn't count as failure
1572        let metrics = broadcaster.get_metrics_async().await;
1573        assert_eq!(
1574            metrics.failed_cancels, 0,
1575            "Idempotent success should not increment failed_cancels"
1576        );
1577        assert_eq!(metrics.idempotent_successes, 1);
1578        assert_eq!(metrics.successful_cancels, 0);
1579    }
1580
1581    #[tokio::test]
1582    async fn test_metrics_initialization_and_health() {
1583        // Validates that metrics start at zero and clients start healthy
1584        let transports = vec![
1585            create_stub_transport("client-0", |_, _, _| async {
1586                Ok(create_test_report("ORDER-1"))
1587            }),
1588            create_stub_transport("client-1", |_, _, _| async {
1589                Ok(create_test_report("ORDER-1"))
1590            }),
1591            create_stub_transport("client-2", |_, _, _| async {
1592                Ok(create_test_report("ORDER-1"))
1593            }),
1594            create_stub_transport("client-3", |_, _, _| async {
1595                Ok(create_test_report("ORDER-1"))
1596            }),
1597        ];
1598
1599        let config = CancelBroadcasterConfig::default();
1600        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1601        let metrics = broadcaster.get_metrics_async().await;
1602
1603        // All metrics should start at zero
1604        assert_eq!(metrics.total_cancels, 0);
1605        assert_eq!(metrics.successful_cancels, 0);
1606        assert_eq!(metrics.failed_cancels, 0);
1607        assert_eq!(metrics.expected_rejects, 0);
1608        assert_eq!(metrics.idempotent_successes, 0);
1609
1610        // All clients should start healthy
1611        assert_eq!(metrics.healthy_clients, 4);
1612        assert_eq!(metrics.total_clients, 4);
1613    }
1614
1615    // Health-check task lifecycle test
1616    #[tokio::test]
1617    async fn test_health_check_task_lifecycle() {
1618        let transports = vec![create_stub_transport("client-0", |_, _, _| async {
1619            Ok(create_test_report("ORDER-1"))
1620        })];
1621
1622        let config = CancelBroadcasterConfig {
1623            health_check_interval_secs: 1, // Very short interval
1624            health_check_timeout_secs: 1,
1625            ..Default::default()
1626        };
1627
1628        let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1629
1630        // Start the broadcaster
1631        broadcaster.start().await.unwrap();
1632        assert!(broadcaster.running.load(Ordering::Relaxed));
1633
1634        // Verify task handle exists
1635        {
1636            let task_guard = broadcaster.health_check_task.read().await;
1637            assert!(task_guard.is_some());
1638        }
1639
1640        // Stop the broadcaster
1641        broadcaster.stop().await;
1642        assert!(!broadcaster.running.load(Ordering::Relaxed));
1643
1644        // Verify task handle has been cleared
1645        {
1646            let task_guard = broadcaster.health_check_task.read().await;
1647            assert!(task_guard.is_none());
1648        }
1649    }
1650}