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_model::{
47 enums::OrderSide,
48 identifiers::{ClientOrderId, InstrumentId, VenueOrderId},
49 instruments::InstrumentAny,
50 reports::OrderStatusReport,
51};
52use tokio::{sync::RwLock, task::JoinHandle, time::interval};
53
54use crate::{
55 common::{consts::BITMEX_HTTP_TESTNET_URL, enums::BitmexEnvironment},
56 http::client::BitmexHttpClient,
57};
58
59const IDEMPOTENT_ALREADY_CANCELED: &str = "AlreadyCanceled";
60const IDEMPOTENT_ORDER_NOT_FOUND: &str = "orderID not found";
61const IDEMPOTENT_UNABLE_DUE_TO_STATE: &str = "Unable to cancel order due to existing state";
62
63trait CancelExecutor: Send + Sync {
85 fn add_instrument(&self, instrument: InstrumentAny);
87
88 fn health_check(&self) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send + '_>>;
90
91 fn cancel_order(
93 &self,
94 instrument_id: InstrumentId,
95 client_order_id: Option<ClientOrderId>,
96 venue_order_id: Option<VenueOrderId>,
97 ) -> Pin<Box<dyn Future<Output = anyhow::Result<OrderStatusReport>> + Send + '_>>;
98
99 fn cancel_orders(
101 &self,
102 instrument_id: InstrumentId,
103 client_order_ids: Option<Vec<ClientOrderId>>,
104 venue_order_ids: Option<Vec<VenueOrderId>>,
105 ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>>;
106
107 fn cancel_all_orders(
109 &self,
110 instrument_id: InstrumentId,
111 order_side: Option<OrderSide>,
112 ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>>;
113}
114
115impl CancelExecutor for BitmexHttpClient {
116 fn add_instrument(&self, instrument: InstrumentAny) {
117 Self::cache_instrument(self, instrument);
118 }
119
120 fn health_check(&self) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send + '_>> {
121 Box::pin(async move {
122 Self::get_server_time(self)
123 .await
124 .map(|_| ())
125 .map_err(|e| anyhow::anyhow!("{e}"))
126 })
127 }
128
129 fn cancel_order(
130 &self,
131 instrument_id: InstrumentId,
132 client_order_id: Option<ClientOrderId>,
133 venue_order_id: Option<VenueOrderId>,
134 ) -> Pin<Box<dyn Future<Output = anyhow::Result<OrderStatusReport>> + Send + '_>> {
135 Box::pin(async move {
136 Self::cancel_order(self, instrument_id, client_order_id, venue_order_id).await
137 })
138 }
139
140 fn cancel_orders(
141 &self,
142 instrument_id: InstrumentId,
143 client_order_ids: Option<Vec<ClientOrderId>>,
144 venue_order_ids: Option<Vec<VenueOrderId>>,
145 ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>> {
146 Box::pin(async move {
147 Self::cancel_orders(self, instrument_id, client_order_ids, venue_order_ids).await
148 })
149 }
150
151 fn cancel_all_orders(
152 &self,
153 instrument_id: InstrumentId,
154 order_side: Option<OrderSide>,
155 ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>> {
156 Box::pin(async move { Self::cancel_all_orders(self, instrument_id, order_side).await })
157 }
158}
159
160#[derive(Debug, Clone)]
162pub struct CancelBroadcasterConfig {
163 pub pool_size: usize,
165 pub api_key: Option<String>,
167 pub api_secret: Option<String>,
169 pub base_url: Option<String>,
171 pub environment: BitmexEnvironment,
173 pub timeout_secs: u64,
175 pub max_retries: u32,
177 pub retry_delay_ms: u64,
179 pub retry_delay_max_ms: u64,
181 pub recv_window_ms: u64,
183 pub max_requests_per_second: u32,
185 pub max_requests_per_minute: u32,
187 pub health_check_interval_secs: u64,
189 pub health_check_timeout_secs: u64,
191 pub expected_reject_patterns: Vec<String>,
193 pub idempotent_success_patterns: Vec<String>,
195 pub proxy_urls: Vec<Option<String>>,
201}
202
203impl Default for CancelBroadcasterConfig {
204 fn default() -> Self {
205 Self {
206 pool_size: 2,
207 api_key: None,
208 api_secret: None,
209 base_url: None,
210 environment: BitmexEnvironment::Mainnet,
211 timeout_secs: 60,
212 max_retries: 3,
213 retry_delay_ms: 1_000,
214 retry_delay_max_ms: 5_000,
215 recv_window_ms: 10_000,
216 max_requests_per_second: 10,
217 max_requests_per_minute: 120,
218 health_check_interval_secs: 30,
219 health_check_timeout_secs: 5,
220 expected_reject_patterns: vec![
221 "Order had execInst of ParticipateDoNotInitiate".to_string(),
222 ],
223 idempotent_success_patterns: vec![
224 IDEMPOTENT_ALREADY_CANCELED.to_string(),
225 IDEMPOTENT_ORDER_NOT_FOUND.to_string(),
226 IDEMPOTENT_UNABLE_DUE_TO_STATE.to_string(),
227 ],
228 proxy_urls: vec![],
229 }
230 }
231}
232
233#[derive(Clone)]
235struct TransportClient {
236 executor: Arc<dyn CancelExecutor>,
241 client_id: String,
242 healthy: Arc<AtomicBool>,
243 cancel_count: Arc<AtomicU64>,
244 error_count: Arc<AtomicU64>,
245}
246
247impl Debug for TransportClient {
248 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
249 f.debug_struct(stringify!(TransportClient))
250 .field("client_id", &self.client_id)
251 .field("healthy", &self.healthy)
252 .field("cancel_count", &self.cancel_count)
253 .field("error_count", &self.error_count)
254 .finish()
255 }
256}
257
258impl TransportClient {
259 fn new<E: CancelExecutor + 'static>(executor: E, client_id: String) -> Self {
260 Self {
261 executor: Arc::new(executor),
262 client_id,
263 healthy: Arc::new(AtomicBool::new(true)),
264 cancel_count: Arc::new(AtomicU64::new(0)),
265 error_count: Arc::new(AtomicU64::new(0)),
266 }
267 }
268
269 fn is_healthy(&self) -> bool {
270 self.healthy.load(Ordering::Relaxed)
271 }
272
273 fn mark_healthy(&self) {
274 self.healthy.store(true, Ordering::Relaxed);
275 }
276
277 fn mark_unhealthy(&self) {
278 self.healthy.store(false, Ordering::Relaxed);
279 }
280
281 fn get_cancel_count(&self) -> u64 {
282 self.cancel_count.load(Ordering::Relaxed)
283 }
284
285 fn get_error_count(&self) -> u64 {
286 self.error_count.load(Ordering::Relaxed)
287 }
288
289 async fn health_check(&self, timeout_secs: u64) -> bool {
290 match tokio::time::timeout(
291 Duration::from_secs(timeout_secs),
292 self.executor.health_check(),
293 )
294 .await
295 {
296 Ok(Ok(())) => {
297 self.mark_healthy();
298 true
299 }
300 Ok(Err(e)) => {
301 log::warn!("Health check failed for client {}: {e:?}", self.client_id);
302 self.mark_unhealthy();
303 false
304 }
305 Err(_) => {
306 log::warn!("Health check timeout for client {}", self.client_id);
307 self.mark_unhealthy();
308 false
309 }
310 }
311 }
312
313 async fn cancel_order(
314 &self,
315 instrument_id: InstrumentId,
316 client_order_id: Option<ClientOrderId>,
317 venue_order_id: Option<VenueOrderId>,
318 ) -> anyhow::Result<OrderStatusReport> {
319 self.cancel_count.fetch_add(1, Ordering::Relaxed);
320
321 match self
322 .executor
323 .cancel_order(instrument_id, client_order_id, venue_order_id)
324 .await
325 {
326 Ok(report) => {
327 self.mark_healthy();
328 Ok(report)
329 }
330 Err(e) => {
331 self.error_count.fetch_add(1, Ordering::Relaxed);
332 Err(e)
333 }
334 }
335 }
336}
337
338#[cfg_attr(feature = "python", pyo3::pyclass)]
344#[cfg_attr(
345 feature = "python",
346 pyo3_stub_gen::derive::gen_stub_pyclass(module = "nautilus_trader.adapters.bitmex")
347)]
348#[derive(Debug)]
349pub struct CancelBroadcaster {
350 config: CancelBroadcasterConfig,
351 transports: Arc<[TransportClient]>,
352 health_check_task: Arc<RwLock<Option<JoinHandle<()>>>>,
353 running: Arc<AtomicBool>,
354 total_cancels: Arc<AtomicU64>,
355 successful_cancels: Arc<AtomicU64>,
356 failed_cancels: Arc<AtomicU64>,
357 expected_rejects: Arc<AtomicU64>,
358 idempotent_successes: Arc<AtomicU64>,
359}
360
361impl CancelBroadcaster {
362 pub fn new(config: CancelBroadcasterConfig) -> anyhow::Result<Self> {
368 let mut transports = Vec::with_capacity(config.pool_size);
369
370 let base_url = match config.environment {
371 BitmexEnvironment::Testnet if config.base_url.is_none() => {
372 Some(BITMEX_HTTP_TESTNET_URL.to_string())
373 }
374 _ => config.base_url.clone(),
375 };
376
377 for i in 0..config.pool_size {
378 let proxy_url = config.proxy_urls.get(i).and_then(|p| p.clone());
380
381 let client = BitmexHttpClient::with_credentials(
382 config.api_key.clone(),
383 config.api_secret.clone(),
384 base_url.clone(),
385 config.timeout_secs,
386 config.max_retries,
387 config.retry_delay_ms,
388 config.retry_delay_max_ms,
389 config.recv_window_ms,
390 config.max_requests_per_second,
391 config.max_requests_per_minute,
392 proxy_url,
393 )
394 .map_err(|e| anyhow::anyhow!("Failed to create HTTP client {i}: {e}"))?;
395
396 transports.push(TransportClient::new(client, format!("bitmex-cancel-{i}")));
397 }
398
399 Ok(Self {
400 config,
401 transports: Arc::from(transports),
402 health_check_task: Arc::new(RwLock::new(None)),
403 running: Arc::new(AtomicBool::new(false)),
404 total_cancels: Arc::new(AtomicU64::new(0)),
405 successful_cancels: Arc::new(AtomicU64::new(0)),
406 failed_cancels: Arc::new(AtomicU64::new(0)),
407 expected_rejects: Arc::new(AtomicU64::new(0)),
408 idempotent_successes: Arc::new(AtomicU64::new(0)),
409 })
410 }
411
412 pub async fn start(&self) -> anyhow::Result<()> {
418 if self.running.load(Ordering::Relaxed) {
419 return Ok(());
420 }
421
422 self.running.store(true, Ordering::Relaxed);
423
424 self.run_health_checks().await;
426
427 let transports = Arc::clone(&self.transports);
429 let running = Arc::clone(&self.running);
430 let interval_secs = self.config.health_check_interval_secs;
431 let timeout_secs = self.config.health_check_timeout_secs;
432
433 let task = get_runtime().spawn(async move {
434 let mut ticker = interval(Duration::from_secs(interval_secs));
435 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
436
437 loop {
438 ticker.tick().await;
439
440 if !running.load(Ordering::Relaxed) {
441 break;
442 }
443
444 let tasks: Vec<_> = transports
445 .iter()
446 .map(|t| t.health_check(timeout_secs))
447 .collect();
448
449 let results = future::join_all(tasks).await;
450 let healthy_count = results.iter().filter(|&&r| r).count();
451
452 log::debug!(
453 "Health check complete: {}/{} clients healthy",
454 healthy_count,
455 results.len()
456 );
457 }
458 });
459
460 *self.health_check_task.write().await = Some(task);
461
462 log::debug!(
463 "CancelBroadcaster started with {} clients",
464 self.transports.len()
465 );
466
467 Ok(())
468 }
469
470 pub async fn stop(&self) {
472 if !self.running.load(Ordering::Relaxed) {
473 return;
474 }
475
476 self.running.store(false, Ordering::Relaxed);
477
478 if let Some(task) = self.health_check_task.write().await.take() {
479 task.abort();
480 }
481
482 log::debug!("CancelBroadcaster stopped");
483 }
484
485 async fn run_health_checks(&self) {
486 let tasks: Vec<_> = self
487 .transports
488 .iter()
489 .map(|t| t.health_check(self.config.health_check_timeout_secs))
490 .collect();
491
492 let results = future::join_all(tasks).await;
493 let healthy_count = results.iter().filter(|&&r| r).count();
494
495 log::debug!(
496 "Health check complete: {}/{} clients healthy",
497 healthy_count,
498 results.len()
499 );
500 }
501
502 fn is_expected_reject(&self, error_message: &str) -> bool {
503 self.config
504 .expected_reject_patterns
505 .iter()
506 .any(|pattern| error_message.contains(pattern))
507 }
508
509 fn is_idempotent_success(&self, error_message: &str) -> bool {
510 self.config
511 .idempotent_success_patterns
512 .iter()
513 .any(|pattern| error_message.contains(pattern))
514 }
515
516 async fn process_cancel_results<T>(
520 &self,
521 mut handles: Vec<JoinHandle<(String, anyhow::Result<T>)>>,
522 idempotent_result: impl FnOnce() -> anyhow::Result<T>,
523 operation: &str,
524 params: String,
525 idempotent_reason: &str,
526 ) -> anyhow::Result<T>
527 where
528 T: Send + 'static,
529 {
530 let mut errors = Vec::new();
531
532 while !handles.is_empty() {
533 let current_handles = std::mem::take(&mut handles);
534 let (result, _idx, remaining) = future::select_all(current_handles).await;
535 handles = remaining.into_iter().collect();
536
537 match result {
538 Ok((client_id, Ok(result))) => {
539 for handle in &handles {
541 handle.abort();
542 }
543 self.successful_cancels.fetch_add(1, Ordering::Relaxed);
544
545 log::debug!("{operation} broadcast succeeded [{client_id}] {params}");
546
547 return Ok(result);
548 }
549 Ok((client_id, Err(e))) => {
550 let error_msg = e.to_string();
551
552 if self.is_idempotent_success(&error_msg) {
553 for handle in &handles {
555 handle.abort();
556 }
557 self.idempotent_successes.fetch_add(1, Ordering::Relaxed);
558
559 log::debug!(
560 "Idempotent success [{client_id}] - {idempotent_reason}: {error_msg} {params}",
561 );
562
563 return idempotent_result();
564 }
565
566 if self.is_expected_reject(&error_msg) {
567 self.expected_rejects.fetch_add(1, Ordering::Relaxed);
568 log::debug!(
569 "Expected {} rejection [{}]: {} {}",
570 operation.to_lowercase(),
571 client_id,
572 error_msg,
573 params
574 );
575 errors.push(error_msg);
576 } else {
577 log::warn!(
578 "{operation} request failed [{client_id}]: {error_msg} {params}"
579 );
580 errors.push(error_msg);
581 }
582 }
583 Err(e) => {
584 log::warn!("{operation} task join error: {e:?}");
585 errors.push(format!("Task panicked: {e:?}"));
586 }
587 }
588 }
589
590 self.failed_cancels.fetch_add(1, Ordering::Relaxed);
592 log::error!(
593 "All {} requests failed: {errors:?} {params}",
594 operation.to_lowercase(),
595 );
596 Err(anyhow::anyhow!(
597 "All {} requests failed: {errors:?}",
598 operation.to_lowercase(),
599 ))
600 }
601
602 pub async fn broadcast_cancel(
614 &self,
615 instrument_id: InstrumentId,
616 client_order_id: Option<ClientOrderId>,
617 venue_order_id: Option<VenueOrderId>,
618 ) -> anyhow::Result<Option<OrderStatusReport>> {
619 self.total_cancels.fetch_add(1, Ordering::Relaxed);
620
621 let healthy_transports: Vec<TransportClient> = self
622 .transports
623 .iter()
624 .filter(|t| t.is_healthy())
625 .cloned()
626 .collect();
627
628 if healthy_transports.is_empty() {
629 self.failed_cancels.fetch_add(1, Ordering::Relaxed);
630 anyhow::bail!("No healthy transport clients available");
631 }
632
633 let mut handles = Vec::new();
634
635 for transport in healthy_transports {
636 let handle = get_runtime().spawn(async move {
637 let client_id = transport.client_id.clone();
638 let result = transport
639 .cancel_order(instrument_id, client_order_id, venue_order_id)
640 .await
641 .map(Some); (client_id, result)
643 });
644 handles.push(handle);
645 }
646
647 self.process_cancel_results(
648 handles,
649 || Ok(None),
650 "Cancel",
651 format!("(client_order_id={client_order_id:?}, venue_order_id={venue_order_id:?})"),
652 "order already cancelled/not found",
653 )
654 .await
655 }
656
657 pub async fn broadcast_batch_cancel(
663 &self,
664 instrument_id: InstrumentId,
665 client_order_ids: Option<Vec<ClientOrderId>>,
666 venue_order_ids: Option<Vec<VenueOrderId>>,
667 ) -> anyhow::Result<Vec<OrderStatusReport>> {
668 self.total_cancels.fetch_add(1, Ordering::Relaxed);
669
670 let healthy_transports: Vec<TransportClient> = self
671 .transports
672 .iter()
673 .filter(|t| t.is_healthy())
674 .cloned()
675 .collect();
676
677 if healthy_transports.is_empty() {
678 self.failed_cancels.fetch_add(1, Ordering::Relaxed);
679 anyhow::bail!("No healthy transport clients available");
680 }
681
682 let mut handles = Vec::new();
683
684 for transport in healthy_transports {
685 let client_order_ids_clone = client_order_ids.clone();
686 let venue_order_ids_clone = venue_order_ids.clone();
687 let handle = get_runtime().spawn(async move {
688 let client_id = transport.client_id.clone();
689 let result = transport
690 .executor
691 .cancel_orders(instrument_id, client_order_ids_clone, venue_order_ids_clone)
692 .await;
693 (client_id, result)
694 });
695 handles.push(handle);
696 }
697
698 self.process_cancel_results(
699 handles,
700 || Ok(Vec::new()),
701 "Batch cancel",
702 format!("(client_order_ids={client_order_ids:?}, venue_order_ids={venue_order_ids:?})"),
703 "orders already cancelled/not found",
704 )
705 .await
706 }
707
708 pub async fn broadcast_cancel_all(
714 &self,
715 instrument_id: InstrumentId,
716 order_side: Option<OrderSide>,
717 ) -> anyhow::Result<Vec<OrderStatusReport>> {
718 self.total_cancels.fetch_add(1, Ordering::Relaxed);
719
720 let healthy_transports: Vec<TransportClient> = self
721 .transports
722 .iter()
723 .filter(|t| t.is_healthy())
724 .cloned()
725 .collect();
726
727 if healthy_transports.is_empty() {
728 self.failed_cancels.fetch_add(1, Ordering::Relaxed);
729 anyhow::bail!("No healthy transport clients available");
730 }
731
732 let mut handles = Vec::new();
733
734 for transport in healthy_transports {
735 let handle = get_runtime().spawn(async move {
736 let client_id = transport.client_id.clone();
737 let result = transport
738 .executor
739 .cancel_all_orders(instrument_id, order_side)
740 .await;
741 (client_id, result)
742 });
743 handles.push(handle);
744 }
745
746 self.process_cancel_results(
747 handles,
748 || Ok(Vec::new()),
749 "Cancel all",
750 format!("(instrument_id={instrument_id}, order_side={order_side:?})"),
751 "no orders to cancel",
752 )
753 .await
754 }
755
756 pub fn get_metrics(&self) -> BroadcasterMetrics {
758 let healthy_clients = self.transports.iter().filter(|t| t.is_healthy()).count();
759 let total_clients = self.transports.len();
760
761 BroadcasterMetrics {
762 total_cancels: self.total_cancels.load(Ordering::Relaxed),
763 successful_cancels: self.successful_cancels.load(Ordering::Relaxed),
764 failed_cancels: self.failed_cancels.load(Ordering::Relaxed),
765 expected_rejects: self.expected_rejects.load(Ordering::Relaxed),
766 idempotent_successes: self.idempotent_successes.load(Ordering::Relaxed),
767 healthy_clients,
768 total_clients,
769 }
770 }
771
772 pub async fn get_metrics_async(&self) -> BroadcasterMetrics {
774 self.get_metrics()
775 }
776
777 pub fn get_client_stats(&self) -> Vec<ClientStats> {
779 self.transports
780 .iter()
781 .map(|t| ClientStats {
782 client_id: t.client_id.clone(),
783 healthy: t.is_healthy(),
784 cancel_count: t.get_cancel_count(),
785 error_count: t.get_error_count(),
786 })
787 .collect()
788 }
789
790 pub async fn get_client_stats_async(&self) -> Vec<ClientStats> {
792 self.get_client_stats()
793 }
794
795 pub fn cache_instrument(&self, instrument: &InstrumentAny) {
797 for transport in self.transports.iter() {
798 transport.executor.add_instrument(instrument.clone());
799 }
800 }
801
802 #[must_use]
803 pub fn clone_for_async(&self) -> Self {
804 Self {
805 config: self.config.clone(),
806 transports: Arc::clone(&self.transports),
807 health_check_task: Arc::clone(&self.health_check_task),
808 running: Arc::clone(&self.running),
809 total_cancels: Arc::clone(&self.total_cancels),
810 successful_cancels: Arc::clone(&self.successful_cancels),
811 failed_cancels: Arc::clone(&self.failed_cancels),
812 expected_rejects: Arc::clone(&self.expected_rejects),
813 idempotent_successes: Arc::clone(&self.idempotent_successes),
814 }
815 }
816
817 #[cfg(test)]
818 fn new_with_transports(
819 config: CancelBroadcasterConfig,
820 transports: Vec<TransportClient>,
821 ) -> Self {
822 Self {
823 config,
824 transports: Arc::from(transports),
825 health_check_task: Arc::new(RwLock::new(None)),
826 running: Arc::new(AtomicBool::new(false)),
827 total_cancels: Arc::new(AtomicU64::new(0)),
828 successful_cancels: Arc::new(AtomicU64::new(0)),
829 failed_cancels: Arc::new(AtomicU64::new(0)),
830 expected_rejects: Arc::new(AtomicU64::new(0)),
831 idempotent_successes: Arc::new(AtomicU64::new(0)),
832 }
833 }
834}
835
836#[derive(Debug, Clone)]
838pub struct BroadcasterMetrics {
839 pub total_cancels: u64,
840 pub successful_cancels: u64,
841 pub failed_cancels: u64,
842 pub expected_rejects: u64,
843 pub idempotent_successes: u64,
844 pub healthy_clients: usize,
845 pub total_clients: usize,
846}
847
848#[derive(Debug, Clone)]
850pub struct ClientStats {
851 pub client_id: String,
852 pub healthy: bool,
853 pub cancel_count: u64,
854 pub error_count: u64,
855}
856
857#[cfg(test)]
858mod tests {
859 use std::{str::FromStr, sync::atomic::Ordering, time::Duration};
860
861 use nautilus_core::UUID4;
862 use nautilus_model::{
863 enums::{OrderSide, OrderStatus, OrderType, TimeInForce},
864 identifiers::{AccountId, ClientOrderId, InstrumentId, VenueOrderId},
865 reports::OrderStatusReport,
866 types::{Price, Quantity},
867 };
868
869 use super::*;
870
871 #[derive(Clone)]
873 #[expect(clippy::type_complexity)]
874 struct MockExecutor {
875 handler: Arc<
876 dyn Fn(
877 InstrumentId,
878 Option<ClientOrderId>,
879 Option<VenueOrderId>,
880 )
881 -> Pin<Box<dyn Future<Output = anyhow::Result<OrderStatusReport>> + Send>>
882 + Send
883 + Sync,
884 >,
885 }
886
887 impl MockExecutor {
888 fn new<F, Fut>(handler: F) -> Self
889 where
890 F: Fn(InstrumentId, Option<ClientOrderId>, Option<VenueOrderId>) -> Fut
891 + Send
892 + Sync
893 + 'static,
894 Fut: Future<Output = anyhow::Result<OrderStatusReport>> + Send + 'static,
895 {
896 Self {
897 handler: Arc::new(move |id, cid, vid| Box::pin(handler(id, cid, vid))),
898 }
899 }
900 }
901
902 impl CancelExecutor for MockExecutor {
903 fn health_check(&self) -> Pin<Box<dyn Future<Output = anyhow::Result<()>> + Send + '_>> {
904 Box::pin(async { Ok(()) })
905 }
906
907 fn cancel_order(
908 &self,
909 instrument_id: InstrumentId,
910 client_order_id: Option<ClientOrderId>,
911 venue_order_id: Option<VenueOrderId>,
912 ) -> Pin<Box<dyn Future<Output = anyhow::Result<OrderStatusReport>> + Send + '_>> {
913 (self.handler)(instrument_id, client_order_id, venue_order_id)
914 }
915
916 fn cancel_orders(
917 &self,
918 _instrument_id: InstrumentId,
919 _client_order_ids: Option<Vec<ClientOrderId>>,
920 _venue_order_ids: Option<Vec<VenueOrderId>>,
921 ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>>
922 {
923 Box::pin(async { Ok(Vec::new()) })
924 }
925
926 fn cancel_all_orders(
927 &self,
928 instrument_id: InstrumentId,
929 _order_side: Option<OrderSide>,
930 ) -> Pin<Box<dyn Future<Output = anyhow::Result<Vec<OrderStatusReport>>> + Send + '_>>
931 {
932 let handler = Arc::clone(&self.handler);
934 Box::pin(async move {
935 let result = handler(instrument_id, None, None).await;
937 match result {
938 Ok(_) => Ok(Vec::new()),
939 Err(e) => Err(e),
940 }
941 })
942 }
943
944 fn add_instrument(&self, _instrument: InstrumentAny) {
945 }
947 }
948
949 fn create_test_report(venue_order_id: &str) -> OrderStatusReport {
950 OrderStatusReport {
951 account_id: AccountId::from("BITMEX-001"),
952 instrument_id: InstrumentId::from_str("XBTUSD.BITMEX").unwrap(),
953 venue_order_id: VenueOrderId::from(venue_order_id),
954 order_side: OrderSide::Buy.into(),
955 order_type: OrderType::Limit,
956 time_in_force: TimeInForce::Gtc,
957 order_status: OrderStatus::Canceled,
958 price: Some(Price::new(50000.0, 2)),
959 quantity: Quantity::new(100.0, 0),
960 filled_qty: Quantity::new(0.0, 0),
961 report_id: UUID4::new(),
962 ts_accepted: 0.into(),
963 ts_last: 0.into(),
964 ts_init: 0.into(),
965 client_order_id: None,
966 avg_px: None,
967 activation_price: None,
968 trigger_price: None,
969 trigger_type: None,
970 contingency_type: None,
971 expire_time: None,
972 order_list_id: None,
973 venue_position_id: None,
974 linked_order_ids: None,
975 parent_order_id: None,
976 display_qty: None,
977 limit_offset: None,
978 trailing_offset: None,
979 trailing_offset_type: None,
980 post_only: false,
981 reduce_only: false,
982 cancel_reason: None,
983 ts_triggered: None,
984 }
985 }
986
987 fn create_stub_transport<F, Fut>(client_id: &str, handler: F) -> TransportClient
988 where
989 F: Fn(InstrumentId, Option<ClientOrderId>, Option<VenueOrderId>) -> Fut
990 + Send
991 + Sync
992 + 'static,
993 Fut: Future<Output = anyhow::Result<OrderStatusReport>> + Send + 'static,
994 {
995 let executor = MockExecutor::new(handler);
996 TransportClient::new(executor, client_id.to_string())
997 }
998
999 #[tokio::test]
1000 async fn test_broadcast_cancel_immediate_success() {
1001 let report = create_test_report("ORDER-1");
1002 let report_clone = report.clone();
1003
1004 let transports = vec![
1005 create_stub_transport("client-0", move |_, _, _| {
1006 let report = report_clone.clone();
1007 async move { Ok(report) }
1008 }),
1009 create_stub_transport("client-1", |_, _, _| async {
1010 tokio::time::sleep(Duration::from_secs(10)).await;
1011 anyhow::bail!("Should be aborted")
1012 }),
1013 ];
1014
1015 let config = CancelBroadcasterConfig::default();
1016 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1017
1018 let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1019 let result = broadcaster
1020 .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-123")), None)
1021 .await;
1022
1023 assert!(result.is_ok());
1024 let returned_report = result.unwrap();
1025 assert!(returned_report.is_some());
1026 assert_eq!(
1027 returned_report.unwrap().venue_order_id,
1028 report.venue_order_id
1029 );
1030
1031 let metrics = broadcaster.get_metrics_async().await;
1032 assert_eq!(metrics.successful_cancels, 1);
1033 assert_eq!(metrics.failed_cancels, 0);
1034 assert_eq!(metrics.total_cancels, 1);
1035 }
1036
1037 #[tokio::test]
1038 async fn test_broadcast_cancel_idempotent_success() {
1039 let transports = vec![
1040 create_stub_transport("client-0", |_, _, _| async {
1041 anyhow::bail!("AlreadyCanceled")
1042 }),
1043 create_stub_transport("client-1", |_, _, _| async {
1044 tokio::time::sleep(Duration::from_secs(10)).await;
1045 anyhow::bail!("Should be aborted")
1046 }),
1047 ];
1048
1049 let config = CancelBroadcasterConfig::default();
1050 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1051
1052 let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1053 let result = broadcaster
1054 .broadcast_cancel(instrument_id, None, Some(VenueOrderId::from("12345")))
1055 .await;
1056
1057 assert!(result.is_ok());
1058 assert!(result.unwrap().is_none());
1059
1060 let metrics = broadcaster.get_metrics_async().await;
1061 assert_eq!(metrics.idempotent_successes, 1);
1062 assert_eq!(metrics.successful_cancels, 0);
1063 assert_eq!(metrics.failed_cancels, 0);
1064 }
1065
1066 #[tokio::test]
1067 async fn test_broadcast_cancel_mixed_idempotent_and_failure() {
1068 let transports = vec![
1069 create_stub_transport("client-0", |_, _, _| async {
1070 anyhow::bail!("502 Bad Gateway")
1071 }),
1072 create_stub_transport("client-1", |_, _, _| async {
1073 anyhow::bail!("orderID not found")
1074 }),
1075 ];
1076
1077 let config = CancelBroadcasterConfig::default();
1078 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1079
1080 let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1081 let result = broadcaster
1082 .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-456")), None)
1083 .await;
1084
1085 assert!(result.is_ok());
1086 assert!(result.unwrap().is_none());
1087
1088 let metrics = broadcaster.get_metrics_async().await;
1089 assert_eq!(metrics.idempotent_successes, 1);
1090 assert_eq!(metrics.failed_cancels, 0);
1091 }
1092
1093 #[tokio::test]
1094 async fn test_broadcast_cancel_all_failures() {
1095 let transports = vec![
1096 create_stub_transport("client-0", |_, _, _| async {
1097 anyhow::bail!("502 Bad Gateway")
1098 }),
1099 create_stub_transport("client-1", |_, _, _| async {
1100 anyhow::bail!("Connection refused")
1101 }),
1102 ];
1103
1104 let config = CancelBroadcasterConfig::default();
1105 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1106
1107 let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1108 let result = broadcaster.broadcast_cancel_all(instrument_id, None).await;
1109
1110 assert!(result.is_err());
1111 assert!(
1112 result
1113 .unwrap_err()
1114 .to_string()
1115 .contains("All cancel all requests failed")
1116 );
1117
1118 let metrics = broadcaster.get_metrics_async().await;
1119 assert_eq!(metrics.failed_cancels, 1);
1120 assert_eq!(metrics.successful_cancels, 0);
1121 assert_eq!(metrics.idempotent_successes, 0);
1122 }
1123
1124 #[tokio::test]
1125 async fn test_broadcast_cancel_no_healthy_clients() {
1126 let transport = create_stub_transport("client-0", |_, _, _| async {
1127 Ok(create_test_report("ORDER-1"))
1128 });
1129 transport.healthy.store(false, Ordering::Relaxed);
1130
1131 let config = CancelBroadcasterConfig::default();
1132 let broadcaster = CancelBroadcaster::new_with_transports(config, vec![transport]);
1133
1134 let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1135 let result = broadcaster
1136 .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-789")), None)
1137 .await;
1138
1139 assert!(result.is_err());
1140 assert!(
1141 result
1142 .unwrap_err()
1143 .to_string()
1144 .contains("No healthy transport clients available")
1145 );
1146
1147 let metrics = broadcaster.get_metrics_async().await;
1148 assert_eq!(metrics.failed_cancels, 1);
1149 }
1150
1151 #[tokio::test]
1152 async fn test_broadcast_cancel_metrics_increment() {
1153 let report1 = create_test_report("ORDER-1");
1154 let report1_clone = report1.clone();
1155 let report2 = create_test_report("ORDER-2");
1156 let report2_clone = report2.clone();
1157
1158 let transports = vec![
1159 create_stub_transport("client-0", move |_, _, _| {
1160 let report = report1_clone.clone();
1161 async move { Ok(report) }
1162 }),
1163 create_stub_transport("client-1", move |_, _, _| {
1164 let report = report2_clone.clone();
1165 async move { Ok(report) }
1166 }),
1167 ];
1168
1169 let config = CancelBroadcasterConfig::default();
1170 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1171
1172 let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1173
1174 let _ = broadcaster
1175 .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-1")), None)
1176 .await;
1177
1178 let _ = broadcaster
1179 .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-2")), None)
1180 .await;
1181
1182 let metrics = broadcaster.get_metrics_async().await;
1183 assert_eq!(metrics.total_cancels, 2);
1184 assert_eq!(metrics.successful_cancels, 2);
1185 }
1186
1187 #[tokio::test]
1188 async fn test_broadcast_cancel_expected_reject_pattern() {
1189 let transports = vec![
1190 create_stub_transport("client-0", |_, _, _| async {
1191 anyhow::bail!("Order had execInst of ParticipateDoNotInitiate")
1192 }),
1193 create_stub_transport("client-1", |_, _, _| async {
1194 anyhow::bail!("Order had execInst of ParticipateDoNotInitiate")
1195 }),
1196 ];
1197
1198 let config = CancelBroadcasterConfig::default();
1199 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1200
1201 let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1202 let result = broadcaster
1203 .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-PDI")), None)
1204 .await;
1205
1206 assert!(result.is_err());
1207
1208 let metrics = broadcaster.get_metrics_async().await;
1209 assert_eq!(metrics.expected_rejects, 2);
1210 assert_eq!(metrics.failed_cancels, 1);
1211 }
1212
1213 #[tokio::test]
1214 async fn test_broadcaster_creation_with_pool() {
1215 let transports = vec![
1216 create_stub_transport("client-0", |_, _, _| async {
1217 Ok(create_test_report("ORDER-1"))
1218 }),
1219 create_stub_transport("client-1", |_, _, _| async {
1220 Ok(create_test_report("ORDER-1"))
1221 }),
1222 create_stub_transport("client-2", |_, _, _| async {
1223 Ok(create_test_report("ORDER-1"))
1224 }),
1225 ];
1226
1227 let config = CancelBroadcasterConfig::default();
1228 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1229 let metrics = broadcaster.get_metrics_async().await;
1230
1231 assert_eq!(metrics.total_clients, 3);
1232 assert_eq!(metrics.total_cancels, 0);
1233 assert_eq!(metrics.successful_cancels, 0);
1234 assert_eq!(metrics.failed_cancels, 0);
1235 }
1236
1237 #[tokio::test]
1238 async fn test_broadcaster_lifecycle() {
1239 let transports = vec![
1240 create_stub_transport("client-0", |_, _, _| async {
1241 Ok(create_test_report("ORDER-1"))
1242 }),
1243 create_stub_transport("client-1", |_, _, _| async {
1244 Ok(create_test_report("ORDER-1"))
1245 }),
1246 ];
1247
1248 let config = CancelBroadcasterConfig::default();
1249 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1250
1251 assert!(!broadcaster.running.load(Ordering::Relaxed));
1253
1254 let start_result = broadcaster.start().await;
1256 assert!(start_result.is_ok());
1257 assert!(broadcaster.running.load(Ordering::Relaxed));
1258
1259 let start_again = broadcaster.start().await;
1261 assert!(start_again.is_ok());
1262
1263 broadcaster.stop().await;
1265 assert!(!broadcaster.running.load(Ordering::Relaxed));
1266
1267 broadcaster.stop().await;
1269 assert!(!broadcaster.running.load(Ordering::Relaxed));
1270 }
1271
1272 #[tokio::test]
1273 async fn test_client_stats_collection() {
1274 let transports = vec![
1275 create_stub_transport("client-0", |_, _, _| async {
1276 Ok(create_test_report("ORDER-1"))
1277 }),
1278 create_stub_transport("client-1", |_, _, _| async {
1279 Ok(create_test_report("ORDER-1"))
1280 }),
1281 ];
1282
1283 let config = CancelBroadcasterConfig::default();
1284 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1285 let stats = broadcaster.get_client_stats_async().await;
1286
1287 assert_eq!(stats.len(), 2);
1288 assert_eq!(stats[0].client_id, "client-0");
1289 assert_eq!(stats[1].client_id, "client-1");
1290 assert!(stats[0].healthy); assert!(stats[1].healthy);
1292 assert_eq!(stats[0].cancel_count, 0);
1293 assert_eq!(stats[1].cancel_count, 0);
1294 assert_eq!(stats[0].error_count, 0);
1295 assert_eq!(stats[1].error_count, 0);
1296 }
1297
1298 #[tokio::test]
1299 async fn test_testnet_config_sets_base_url() {
1300 let config = CancelBroadcasterConfig {
1301 pool_size: 1,
1302 api_key: Some("test_key".to_string()),
1303 api_secret: Some("test_secret".to_string()),
1304 base_url: None,
1305 environment: BitmexEnvironment::Testnet,
1306 timeout_secs: 5,
1307 max_retries: 3,
1308 retry_delay_ms: 1_000,
1309 retry_delay_max_ms: 5_000,
1310 recv_window_ms: 10_000,
1311 max_requests_per_second: 10,
1312 max_requests_per_minute: 120,
1313 health_check_interval_secs: 60,
1314 health_check_timeout_secs: 5,
1315 expected_reject_patterns: vec![],
1316 idempotent_success_patterns: vec![],
1317 proxy_urls: vec![],
1318 };
1319
1320 let broadcaster = CancelBroadcaster::new(config);
1321 assert!(broadcaster.is_ok());
1322 }
1323
1324 #[tokio::test]
1325 async fn test_constructor_honors_default_pool_size() {
1326 let config = CancelBroadcasterConfig {
1327 api_key: Some("test_key".to_string()),
1328 api_secret: Some("test_secret".to_string()),
1329 base_url: Some("http://127.0.0.1:19999".to_string()),
1330 ..Default::default()
1331 };
1332
1333 let expected_pool = config.pool_size;
1334 let broadcaster = CancelBroadcaster::new(config).unwrap();
1335 let metrics = broadcaster.get_metrics_async().await;
1336
1337 assert_eq!(metrics.total_clients, expected_pool);
1338 }
1339
1340 #[tokio::test]
1341 async fn test_default_config() {
1342 let transports = vec![
1343 create_stub_transport("client-0", |_, _, _| async {
1344 Ok(create_test_report("ORDER-1"))
1345 }),
1346 create_stub_transport("client-1", |_, _, _| async {
1347 Ok(create_test_report("ORDER-1"))
1348 }),
1349 ];
1350
1351 let config = CancelBroadcasterConfig::default();
1352 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1353 let metrics = broadcaster.get_metrics_async().await;
1354
1355 assert_eq!(metrics.total_clients, 2);
1357 }
1358
1359 #[tokio::test]
1360 async fn test_clone_for_async() {
1361 let transports = vec![create_stub_transport("client-0", |_, _, _| async {
1362 Ok(create_test_report("ORDER-1"))
1363 })];
1364
1365 let config = CancelBroadcasterConfig::default();
1366 let broadcaster1 = CancelBroadcaster::new_with_transports(config, transports);
1367
1368 broadcaster1.total_cancels.fetch_add(1, Ordering::Relaxed);
1370
1371 let broadcaster2 = broadcaster1.clone_for_async();
1373 let metrics2 = broadcaster2.get_metrics_async().await;
1374
1375 assert_eq!(metrics2.total_cancels, 1); broadcaster2
1379 .successful_cancels
1380 .fetch_add(5, Ordering::Relaxed);
1381
1382 let metrics1 = broadcaster1.get_metrics_async().await;
1384 assert_eq!(metrics1.successful_cancels, 5);
1385 }
1386
1387 #[tokio::test]
1388 async fn test_pattern_matching() {
1389 let transports = vec![create_stub_transport("client-0", |_, _, _| async {
1391 Ok(create_test_report("ORDER-1"))
1392 })];
1393
1394 let config = CancelBroadcasterConfig {
1395 expected_reject_patterns: vec![
1396 "ParticipateDoNotInitiate".to_string(),
1397 "Close-only".to_string(),
1398 ],
1399 idempotent_success_patterns: vec![
1400 "AlreadyCanceled".to_string(),
1401 "orderID not found".to_string(),
1402 "Unable to cancel".to_string(),
1403 ],
1404 ..Default::default()
1405 };
1406
1407 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1408
1409 assert!(broadcaster.is_expected_reject("Order had execInst of ParticipateDoNotInitiate"));
1411 assert!(broadcaster.is_expected_reject("This is a Close-only order"));
1412 assert!(!broadcaster.is_expected_reject("Connection timeout"));
1413
1414 assert!(broadcaster.is_idempotent_success("AlreadyCanceled"));
1416 assert!(broadcaster.is_idempotent_success("Error: orderID not found for this account"));
1417 assert!(broadcaster.is_idempotent_success("Unable to cancel order due to existing state"));
1418 assert!(!broadcaster.is_idempotent_success("502 Bad Gateway"));
1419 }
1420
1421 #[tokio::test]
1425 async fn test_broadcast_batch_cancel_structure() {
1426 let transports = vec![
1428 create_stub_transport("client-0", |_, _, _| async {
1429 Ok(create_test_report("ORDER-1"))
1430 }),
1431 create_stub_transport("client-1", |_, _, _| async {
1432 Ok(create_test_report("ORDER-1"))
1433 }),
1434 ];
1435
1436 let config = CancelBroadcasterConfig {
1437 idempotent_success_patterns: vec!["AlreadyCanceled".to_string()],
1438 ..Default::default()
1439 };
1440
1441 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1442 let metrics = broadcaster.get_metrics_async().await;
1443
1444 assert_eq!(metrics.total_clients, 2);
1446 assert_eq!(metrics.total_cancels, 0);
1447 assert_eq!(metrics.successful_cancels, 0);
1448 assert_eq!(metrics.failed_cancels, 0);
1449 }
1450
1451 #[tokio::test]
1452 async fn test_broadcast_cancel_all_structure() {
1453 let transports = vec![
1455 create_stub_transport("client-0", |_, _, _| async {
1456 Ok(create_test_report("ORDER-1"))
1457 }),
1458 create_stub_transport("client-1", |_, _, _| async {
1459 Ok(create_test_report("ORDER-1"))
1460 }),
1461 create_stub_transport("client-2", |_, _, _| async {
1462 Ok(create_test_report("ORDER-1"))
1463 }),
1464 ];
1465
1466 let config = CancelBroadcasterConfig {
1467 idempotent_success_patterns: vec!["orderID not found".to_string()],
1468 ..Default::default()
1469 };
1470
1471 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1472 let metrics = broadcaster.get_metrics_async().await;
1473
1474 assert_eq!(metrics.total_clients, 3);
1476 assert_eq!(metrics.healthy_clients, 3);
1477 assert_eq!(metrics.total_cancels, 0);
1478 }
1479
1480 #[tokio::test]
1482 async fn test_single_cancel_metrics_with_mixed_responses() {
1483 let transports = vec![
1486 create_stub_transport("client-0", |_, _, _| async {
1487 anyhow::bail!("Connection timeout")
1488 }),
1489 create_stub_transport("client-1", |_, _, _| async {
1490 anyhow::bail!("AlreadyCanceled")
1491 }),
1492 ];
1493
1494 let config = CancelBroadcasterConfig::default();
1495 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1496
1497 let instrument_id = InstrumentId::from_str("XBTUSD.BITMEX").unwrap();
1498 let result = broadcaster
1499 .broadcast_cancel(instrument_id, Some(ClientOrderId::from("O-123")), None)
1500 .await;
1501
1502 assert!(result.is_ok());
1504 assert!(result.unwrap().is_none());
1505
1506 let metrics = broadcaster.get_metrics_async().await;
1508 assert_eq!(
1509 metrics.failed_cancels, 0,
1510 "Idempotent success should not increment failed_cancels"
1511 );
1512 assert_eq!(metrics.idempotent_successes, 1);
1513 assert_eq!(metrics.successful_cancels, 0);
1514 }
1515
1516 #[tokio::test]
1517 async fn test_metrics_initialization_and_health() {
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 create_stub_transport("client-3", |_, _, _| async {
1530 Ok(create_test_report("ORDER-1"))
1531 }),
1532 ];
1533
1534 let config = CancelBroadcasterConfig::default();
1535 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1536 let metrics = broadcaster.get_metrics_async().await;
1537
1538 assert_eq!(metrics.total_cancels, 0);
1540 assert_eq!(metrics.successful_cancels, 0);
1541 assert_eq!(metrics.failed_cancels, 0);
1542 assert_eq!(metrics.expected_rejects, 0);
1543 assert_eq!(metrics.idempotent_successes, 0);
1544
1545 assert_eq!(metrics.healthy_clients, 4);
1547 assert_eq!(metrics.total_clients, 4);
1548 }
1549
1550 #[tokio::test]
1552 async fn test_health_check_task_lifecycle() {
1553 let transports = vec![create_stub_transport("client-0", |_, _, _| async {
1554 Ok(create_test_report("ORDER-1"))
1555 })];
1556
1557 let config = CancelBroadcasterConfig {
1558 health_check_interval_secs: 1, health_check_timeout_secs: 1,
1560 ..Default::default()
1561 };
1562
1563 let broadcaster = CancelBroadcaster::new_with_transports(config, transports);
1564
1565 broadcaster.start().await.unwrap();
1567 assert!(broadcaster.running.load(Ordering::Relaxed));
1568
1569 {
1571 let task_guard = broadcaster.health_check_task.read().await;
1572 assert!(task_guard.is_some());
1573 }
1574
1575 broadcaster.stop().await;
1577 assert!(!broadcaster.running.load(Ordering::Relaxed));
1578
1579 {
1581 let task_guard = broadcaster.health_check_task.read().await;
1582 assert!(task_guard.is_none());
1583 }
1584 }
1585}