1use 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
65trait CancelExecutor: Send + Sync {
87 fn add_instrument(&self, instrument: InstrumentAny);
89
90 fn health_check(&self) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send + '_>>;
92
93 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 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 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#[derive(Debug, Clone)]
164pub struct CancelBroadcasterConfig {
165 pub pool_size: usize,
167 pub api_key: Option<SecretString>,
169 pub api_secret: Option<SecretString>,
171 pub base_url: Option<String>,
173 pub environment: BitmexEnvironment,
175 pub timeout_secs: u64,
177 pub max_retries: u32,
179 pub retry_delay_ms: u64,
181 pub retry_delay_max_ms: u64,
183 pub recv_window_ms: u64,
185 pub max_requests_per_second: u32,
187 pub max_requests_per_minute: u32,
189 pub health_check_interval_secs: u64,
191 pub health_check_timeout_secs: u64,
193 pub expected_reject_patterns: Vec<String>,
195 pub idempotent_success_patterns: Vec<String>,
197 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#[derive(Clone)]
237struct TransportClient {
238 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#[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 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 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 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 self.run_health_checks().await;
435
436 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 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 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 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 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 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 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); (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 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 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 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 pub async fn get_metrics_async(&self) -> BroadcasterMetrics {
783 self.get_metrics()
784 }
785
786 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 pub async fn get_client_stats_async(&self) -> Vec<ClientStats> {
801 self.get_client_stats()
802 }
803
804 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#[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#[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 #[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 let handler = Arc::clone(&self.handler);
960 Box::pin(async move {
961 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 }
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 assert!(!broadcaster.running.load(Ordering::Relaxed));
1279
1280 let start_result = broadcaster.start().await;
1282 assert!(start_result.is_ok());
1283 assert!(broadcaster.running.load(Ordering::Relaxed));
1284
1285 let start_again = broadcaster.start().await;
1287 assert!(start_again.is_ok());
1288
1289 broadcaster.stop().await;
1291 assert!(!broadcaster.running.load(Ordering::Relaxed));
1292
1293 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); 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 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 broadcaster1.total_cancels.fetch_add(1, Ordering::Relaxed);
1435
1436 let broadcaster2 = broadcaster1.clone_for_async();
1438 let metrics2 = broadcaster2.get_metrics_async().await;
1439
1440 assert_eq!(metrics2.total_cancels, 1); broadcaster2
1444 .successful_cancels
1445 .fetch_add(5, Ordering::Relaxed);
1446
1447 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 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 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 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 #[tokio::test]
1490 async fn test_broadcast_batch_cancel_structure() {
1491 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 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 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 assert_eq!(metrics.total_clients, 3);
1541 assert_eq!(metrics.healthy_clients, 3);
1542 assert_eq!(metrics.total_cancels, 0);
1543 }
1544
1545 #[tokio::test]
1547 async fn test_single_cancel_metrics_with_mixed_responses() {
1548 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 assert!(result.is_ok());
1569 assert!(result.unwrap().is_none());
1570
1571 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 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 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 assert_eq!(metrics.healthy_clients, 4);
1612 assert_eq!(metrics.total_clients, 4);
1613 }
1614
1615 #[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, health_check_timeout_secs: 1,
1625 ..Default::default()
1626 };
1627
1628 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1629
1630 broadcaster.start().await.unwrap();
1632 assert!(broadcaster.running.load(Ordering::Relaxed));
1633
1634 {
1636 let task_guard = broadcaster.health_check_task.read().await;
1637 assert!(task_guard.is_some());
1638 }
1639
1640 broadcaster.stop().await;
1642 assert!(!broadcaster.running.load(Ordering::Relaxed));
1643
1644 {
1646 let task_guard = broadcaster.health_check_task.read().await;
1647 assert!(task_guard.is_none());
1648 }
1649 }
1650}