1use std::collections::HashMap;
7use std::sync::Arc;
8use std::time::Duration;
9
10use futures::future::join_all;
11use tokio::sync::{Mutex, OwnedSemaphorePermit, RwLock, Semaphore};
12use tokio::time::Instant;
13use tokio_util::sync::CancellationToken;
14use tracing::{debug, info, warn};
15
16use aranet_types::{CurrentReading, DeviceInfo, DeviceType};
17
18use crate::connector::{ConnectFn, SensorLink, ble_connector, log_failed_release, release_link};
19use crate::device::Device;
20use crate::error::{ConnectionFailureReason, Error, Result};
21use crate::events::{DeviceEvent, DeviceId, DisconnectReason, EventDispatcher};
22use crate::passive::{PassiveMonitor, PassiveMonitorOptions, PassiveReading};
23use crate::reconnect::ReconnectOptions;
24use crate::scan::{DiscoveredDevice, ScanOptions, scan_with_options};
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default)]
35pub enum DevicePriority {
36 Low,
38 #[default]
40 Normal,
41 High,
43 Critical,
45}
46
47#[derive(Debug, Clone)]
52pub struct AdaptiveInterval {
53 pub base: Duration,
55 current: Duration,
57 pub min: Duration,
59 pub max: Duration,
61 consecutive_successes: u32,
63 consecutive_failures: u32,
65 success_threshold: u32,
67 failure_threshold: u32,
69}
70
71impl Default for AdaptiveInterval {
72 fn default() -> Self {
73 Self {
74 base: Duration::from_secs(30),
75 current: Duration::from_secs(30),
76 min: Duration::from_secs(5),
77 max: Duration::from_secs(120),
78 consecutive_successes: 0,
79 consecutive_failures: 0,
80 success_threshold: 3,
81 failure_threshold: 1,
82 }
83 }
84}
85
86impl AdaptiveInterval {
87 pub fn new(base: Duration, min: Duration, max: Duration) -> Self {
89 Self {
90 base,
91 current: base,
92 min,
93 max,
94 ..Default::default()
95 }
96 }
97
98 pub fn current(&self) -> Duration {
100 self.current
101 }
102
103 pub fn on_success(&mut self) {
108 self.consecutive_failures = 0;
109 self.consecutive_successes += 1;
110
111 if self.consecutive_successes >= self.success_threshold {
112 let new_interval = self.current.saturating_mul(2);
114 self.current = new_interval.min(self.max);
115 self.consecutive_successes = 0;
116 debug!(
117 "Health check stable, increasing interval to {:?}",
118 self.current
119 );
120 }
121 }
122
123 pub fn on_failure(&mut self) {
128 self.consecutive_successes = 0;
129 self.consecutive_failures += 1;
130
131 if self.consecutive_failures >= self.failure_threshold {
132 let new_interval = self.current / 2;
134 self.current = new_interval.max(self.min);
135 self.consecutive_failures = 0;
136 debug!(
137 "Health check unstable, decreasing interval to {:?}",
138 self.current
139 );
140 }
141 }
142
143 pub fn reset(&mut self) {
145 self.current = self.base;
146 self.consecutive_successes = 0;
147 self.consecutive_failures = 0;
148 }
149}
150
151#[derive(Debug)]
157pub struct ManagedDevice {
158 pub id: String,
160 pub name: Option<String>,
162 pub device_type: Option<DeviceType>,
164 device: Option<Arc<Device>>,
167 pub auto_reconnect: bool,
169 pub last_reading: Option<CurrentReading>,
171 pub info: Option<DeviceInfo>,
173 pub reconnect_options: ReconnectOptions,
175 pub priority: DevicePriority,
177 pub consecutive_failures: u32,
179 pub last_success: Option<u64>,
181}
182
183impl ManagedDevice {
184 pub fn new(id: &str) -> Self {
186 Self {
187 id: id.to_string(),
188 name: None,
189 device_type: None,
190 device: None,
191 auto_reconnect: true,
192 last_reading: None,
193 info: None,
194 reconnect_options: ReconnectOptions::default(),
195 priority: DevicePriority::default(),
196 consecutive_failures: 0,
197 last_success: None,
198 }
199 }
200
201 pub fn with_reconnect_options(id: &str, options: ReconnectOptions) -> Self {
203 Self {
204 reconnect_options: options,
205 ..Self::new(id)
206 }
207 }
208
209 pub fn with_priority(id: &str, priority: DevicePriority) -> Self {
211 Self {
212 priority,
213 ..Self::new(id)
214 }
215 }
216
217 pub fn with_options(id: &str, options: ReconnectOptions, priority: DevicePriority) -> Self {
219 Self {
220 reconnect_options: options,
221 priority,
222 ..Self::new(id)
223 }
224 }
225
226 pub fn record_success(&mut self) {
228 self.consecutive_failures = 0;
229 self.last_success = Some(
230 std::time::SystemTime::now()
231 .duration_since(std::time::UNIX_EPOCH)
232 .unwrap_or_default()
233 .as_millis() as u64,
234 );
235 }
236
237 pub fn record_failure(&mut self) {
239 self.consecutive_failures += 1;
240 }
241
242 pub fn has_device(&self) -> bool {
244 self.device.is_some()
245 }
246
247 pub async fn is_connected(&self) -> bool {
249 if let Some(device) = &self.device {
250 device.is_connected().await
251 } else {
252 false
253 }
254 }
255
256 pub fn device(&self) -> Option<&Arc<Device>> {
258 self.device.as_ref()
259 }
260
261 pub fn device_arc(&self) -> Option<Arc<Device>> {
263 self.device.clone()
264 }
265}
266
267#[derive(Debug, Clone)]
269pub struct ManagerConfig {
270 pub scan_options: ScanOptions,
272 pub default_reconnect_options: ReconnectOptions,
286 pub event_capacity: usize,
288 pub health_check_interval: Duration,
290 pub max_concurrent_connections: usize,
296 pub use_adaptive_interval: bool,
302 pub min_health_check_interval: Duration,
304 pub max_health_check_interval: Duration,
306 pub default_priority: DevicePriority,
308 pub use_connection_validation: bool,
319}
320
321impl Default for ManagerConfig {
322 fn default() -> Self {
323 let platform_config = crate::platform::PlatformConfig::for_current_platform();
325
326 Self {
327 scan_options: ScanOptions::default(),
328 default_reconnect_options: ReconnectOptions::unlimited(),
329 event_capacity: 100,
330 health_check_interval: Duration::from_secs(30),
331 max_concurrent_connections: platform_config.max_concurrent_connections,
332 use_adaptive_interval: true,
333 min_health_check_interval: Duration::from_secs(5),
334 max_health_check_interval: Duration::from_secs(120),
335 default_priority: DevicePriority::Normal,
336 use_connection_validation: true,
337 }
338 }
339}
340
341impl ManagerConfig {
342 pub fn with_max_connections(mut self, max: usize) -> Self {
344 self.max_concurrent_connections = max;
345 self
346 }
347
348 pub fn unlimited_connections(mut self) -> Self {
350 self.max_concurrent_connections = 0;
351 self
352 }
353
354 pub fn adaptive_interval(mut self, enabled: bool) -> Self {
356 self.use_adaptive_interval = enabled;
357 self
358 }
359
360 pub fn health_check_interval(mut self, interval: Duration) -> Self {
362 self.health_check_interval = interval;
363 self
364 }
365
366 pub fn default_priority(mut self, priority: DevicePriority) -> Self {
368 self.default_priority = priority;
369 self
370 }
371
372 pub fn connection_validation(mut self, enabled: bool) -> Self {
374 self.use_connection_validation = enabled;
375 self
376 }
377}
378
379const MAX_RETRY_WAIT: Duration = Duration::from_secs(365 * 24 * 60 * 60);
383
384struct Entry<L> {
386 name: Option<String>,
387 device_type: Option<DeviceType>,
388 info: Option<DeviceInfo>,
389 last_reading: Option<CurrentReading>,
390 priority: DevicePriority,
391 reconnect_options: ReconnectOptions,
392 auto_reconnect: bool,
393 failures: u32,
395 wanted: bool,
402 retry_at: Option<Instant>,
405 gave_up: bool,
407 ever_connected: bool,
411 link: Option<Arc<L>>,
413 slot: Option<OwnedSemaphorePermit>,
416 op: Arc<Mutex<()>>,
419}
420
421impl<L> Entry<L> {
422 fn new(reconnect_options: ReconnectOptions, priority: DevicePriority) -> Self {
423 Self {
424 name: None,
425 device_type: None,
426 info: None,
427 last_reading: None,
428 priority,
429 reconnect_options,
430 auto_reconnect: true,
431 failures: 0,
432 wanted: true,
433 retry_at: None,
434 gave_up: false,
435 ever_connected: false,
436 link: None,
437 slot: None,
438 op: Arc::new(Mutex::new(())),
439 }
440 }
441
442 fn is_due_for_repair(&self) -> bool {
444 self.auto_reconnect
445 && self.wanted
446 && !self.gave_up
447 && self.link.is_none()
448 && self.retry_at.is_none_or(|at| at <= Instant::now())
449 }
450
451 fn give_up_if_attempts_used_up(&mut self) -> Option<u32> {
454 let used_up = self
455 .reconnect_options
456 .max_attempts
457 .is_some_and(|max| self.failures >= max);
458 if used_up && !self.gave_up {
459 self.gave_up = true;
460 return Some(self.failures);
461 }
462 None
463 }
464
465 fn give_up_before_reconnect(&mut self) -> Option<u32> {
472 if self.ever_connected {
473 self.give_up_if_attempts_used_up()
474 } else {
475 None
476 }
477 }
478
479 fn record_repair_failure(&mut self, now: Instant) -> Option<u32> {
482 self.failures = self.failures.saturating_add(1);
483 let wait = self
484 .reconnect_options
485 .delay_for_attempt(self.failures - 1)
486 .min(MAX_RETRY_WAIT);
487 self.retry_at = Some(now + wait);
488 self.give_up_if_attempts_used_up()
489 }
490
491 fn reset_backoff(&mut self) {
495 self.failures = 0;
496 self.retry_at = None;
497 self.gave_up = false;
498 }
499}
500
501#[cfg(test)]
503#[derive(Debug, Clone, PartialEq, Eq)]
504pub(crate) struct EntrySnapshot {
505 pub(crate) has_link: bool,
506 pub(crate) failures: u32,
507 pub(crate) wanted: bool,
508 pub(crate) gave_up: bool,
509 pub(crate) retry_at: Option<Instant>,
510}
511
512struct NewLink<L: SensorLink> {
518 identifier: String,
520 parts: Option<(Arc<L>, Option<OwnedSemaphorePermit>)>,
521 held: Arc<tokio::sync::OwnedMutexGuard<()>>,
522}
523
524impl<L: SensorLink> NewLink<L> {
525 fn new(
526 identifier: &str,
527 link: Arc<L>,
528 slot: Option<OwnedSemaphorePermit>,
529 held: &Arc<tokio::sync::OwnedMutexGuard<()>>,
530 ) -> Self {
531 Self {
532 identifier: identifier.to_owned(),
533 parts: Some((link, slot)),
534 held: Arc::clone(held),
535 }
536 }
537
538 fn into_parts(mut self) -> (Arc<L>, Option<OwnedSemaphorePermit>) {
541 self.parts
542 .take()
543 .expect("a NewLink is taken apart only once")
544 }
545}
546
547impl<L: SensorLink> Drop for NewLink<L> {
548 fn drop(&mut self) {
549 let Some((link, slot)) = self.parts.take() else {
550 return;
551 };
552 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
554 let held = Arc::clone(&self.held);
555 let identifier = std::mem::take(&mut self.identifier);
556 runtime.spawn(async move {
557 if let Err(e) = release_link(link, (slot, held)).await {
558 log_failed_release("Disconnecting an abandoned new link", &identifier, &e);
559 }
560 });
561 }
562 }
563}
564
565async fn finished(task: tokio::task::JoinHandle<Result<()>>) -> Result<()> {
568 task.await.map_err(std::io::Error::from)?
569}
570
571#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
573pub(crate) struct TickOutcome {
574 pub(crate) healthy: usize,
576 pub(crate) repaired: usize,
578 pub(crate) failed: usize,
580}
581
582pub(crate) struct ManagerCore<L: SensorLink> {
586 devices: RwLock<HashMap<String, Entry<L>>>,
587 events: EventDispatcher,
588 config: ManagerConfig,
589 connect: ConnectFn<L>,
590 slots: Option<Arc<Semaphore>>,
592}
593
594impl<L: SensorLink> ManagerCore<L> {
595 pub(crate) fn new(config: ManagerConfig, connect: ConnectFn<L>) -> Self {
596 let max = config.max_concurrent_connections;
597 let slots = if max == 0 {
598 None
599 } else if max > Semaphore::MAX_PERMITS {
600 warn!(
601 "max_concurrent_connections ({max}) is above {}; connections are not limited",
602 Semaphore::MAX_PERMITS
603 );
604 None
605 } else {
606 Some(Arc::new(Semaphore::new(max)))
607 };
608 Self {
609 devices: RwLock::new(HashMap::new()),
610 events: EventDispatcher::new(config.event_capacity),
611 config,
612 connect,
613 slots,
614 }
615 }
616
617 fn limit_error(&self, identifier: &str) -> Error {
619 let max = self.config.max_concurrent_connections;
620 let free = self
621 .slots
622 .as_ref()
623 .map_or(0, |slots| slots.available_permits());
624 let used = max.saturating_sub(free);
625 warn!("Connection limit reached ({used}/{max}), cannot connect to {identifier}");
626 Error::connection_failed(
627 Some(identifier.to_string()),
628 ConnectionFailureReason::Other(format!("Connection limit reached ({used}/{max})")),
629 )
630 }
631
632 fn take_slot(&self, identifier: &str) -> Result<Option<OwnedSemaphorePermit>> {
635 match &self.slots {
636 Some(slots) => Arc::clone(slots)
637 .try_acquire_owned()
638 .map(Some)
639 .map_err(|_| self.limit_error(identifier)),
640 None => Ok(None),
641 }
642 }
643
644 pub(crate) async fn add_device(&self, identifier: &str) -> Result<()> {
645 self.add_device_with_options(identifier, self.config.default_reconnect_options.clone())
646 .await
647 }
648
649 pub(crate) async fn add_device_with_options(
650 &self,
651 identifier: &str,
652 reconnect_options: ReconnectOptions,
653 ) -> Result<()> {
654 reconnect_options.validate()?;
655 let mut devices = self.devices.write().await;
656
657 if devices.contains_key(identifier) {
658 return Ok(()); }
660
661 devices.insert(
662 identifier.to_string(),
663 Entry::new(reconnect_options, DevicePriority::default()),
664 );
665
666 info!("Added device to manager: {}", identifier);
667 Ok(())
668 }
669
670 pub(crate) async fn connect(&self, identifier: &str) -> Result<()> {
671 let (op, reserved) = {
672 let mut devices = self.devices.write().await;
673 let reserved = if devices.contains_key(identifier) {
677 None
678 } else {
679 self.take_slot(identifier)?
680 };
681 let entry = devices.entry(identifier.to_string()).or_insert_with(|| {
682 info!("Adding device to manager: {identifier}");
683 Entry::new(
684 self.config.default_reconnect_options.clone(),
685 DevicePriority::default(),
686 )
687 });
688 entry.wanted = true;
690 (Arc::clone(&entry.op), reserved)
691 };
692 let held = Arc::new(Arc::clone(&op).lock_owned().await);
697 if let Some(entry) = self.devices.write().await.get_mut(identifier)
701 && Arc::ptr_eq(&entry.op, &op)
702 {
703 entry.reset_backoff();
704 }
705 self.connect_locked(identifier, &held, reserved).await
706 }
707
708 async fn connect_locked(
716 &self,
717 identifier: &str,
718 held: &Arc<tokio::sync::OwnedMutexGuard<()>>,
719 reserved: Option<OwnedSemaphorePermit>,
720 ) -> Result<()> {
721 let op = tokio::sync::OwnedMutexGuard::mutex(held);
722 let existing = {
723 let devices = self.devices.read().await;
724 match devices.get(identifier) {
725 Some(entry) if Arc::ptr_eq(&entry.op, op) && entry.wanted => entry.link.clone(),
726 _ => return Err(Error::Cancelled),
728 }
729 };
730 if let Some(link) = existing {
731 if link.is_connected().await {
732 debug!("Device {identifier} is already connected");
733 return Ok(());
734 }
735 info!("The connection to {identifier} was lost; reconnecting");
736 self.detach(identifier, &link, held, DisconnectReason::Unknown)
737 .await;
738 }
739
740 let slot = match reserved {
743 Some(slot) => Some(slot),
744 None => self.take_slot(identifier)?,
745 };
746
747 let new_link = NewLink::new(
748 identifier,
749 Arc::new((self.connect)(identifier).await?),
750 slot,
751 held,
752 );
753
754 let stored = {
758 let mut devices = self.devices.write().await;
759 match devices.get_mut(identifier) {
760 Some(entry) if Arc::ptr_eq(&entry.op, op) && entry.wanted => {
761 let (link, slot) = new_link.into_parts();
762 entry.link = Some(Arc::clone(&link));
763 entry.slot = slot;
764 entry.ever_connected = true;
765 entry.reset_backoff();
768 Ok(link)
769 }
770 _ => Err(new_link),
771 }
772 };
773 let link = match stored {
774 Ok(link) => link,
775 Err(new_link) => {
776 debug!("{identifier} was disconnected or removed; closing the new link");
777 let (link, slot) = new_link.into_parts();
778 if let Err(e) = release_link(link, (slot, Arc::clone(held))).await {
779 log_failed_release("Disconnecting the unused link", identifier, &e);
780 }
781 return Err(Error::Cancelled);
782 }
783 };
784
785 let info = link.read_device_info().await.ok();
786 let device_type = link.device_type();
787 let name = link.name().map(|s| s.to_string());
788 {
789 let mut devices = self.devices.write().await;
790 match devices.get_mut(identifier) {
791 Some(entry) if Arc::ptr_eq(&entry.op, op) => {
792 entry.info = info.clone();
793 entry.device_type = device_type;
794 entry.name = name.clone();
795 }
796 _ => return Err(Error::Cancelled),
798 }
799 }
800
801 self.events.send(DeviceEvent::Connected {
803 device: DeviceId {
804 id: identifier.to_string(),
805 name,
806 device_type,
807 },
808 info,
809 });
810
811 info!("Connected to device: {identifier}");
812 Ok(())
813 }
814
815 async fn withdraw(&self, identifier: &str) -> Option<Arc<Mutex<()>>> {
820 let mut devices = self.devices.write().await;
821 let entry = devices.get_mut(identifier)?;
822 entry.wanted = false;
823 Some(Arc::clone(&entry.op))
824 }
825
826 pub(crate) async fn disconnect(self: &Arc<Self>, identifier: &str) -> Result<()> {
827 let Some(op) = self.withdraw(identifier).await else {
828 return Ok(());
829 };
830 finished(self.spawn_disconnect(identifier, op)).await
831 }
832
833 fn spawn_disconnect(
841 self: &Arc<Self>,
842 identifier: &str,
843 op: Arc<Mutex<()>>,
844 ) -> tokio::task::JoinHandle<Result<()>> {
845 let core = Arc::clone(self);
846 let identifier = identifier.to_owned();
847 tokio::spawn(async move {
848 let held = Arc::new(op.lock_owned().await);
849 core.disconnect_locked(&identifier, &held, true).await
850 })
851 }
852
853 async fn disconnect_locked(
866 &self,
867 identifier: &str,
868 held: &Arc<tokio::sync::OwnedMutexGuard<()>>,
869 unless_rearmed: bool,
870 ) -> Result<()> {
871 let op = tokio::sync::OwnedMutexGuard::mutex(held);
872 let taken = {
873 let mut devices = self.devices.write().await;
874 match devices.get_mut(identifier) {
875 Some(entry) if Arc::ptr_eq(&entry.op, op) && !(unless_rearmed && entry.wanted) => {
876 entry.link.take().map(|link| (link, entry.slot.take()))
877 }
878 _ => None,
879 }
880 };
881 let Some((link, slot)) = taken else {
882 return Ok(());
883 };
884
885 self.events.send(DeviceEvent::Disconnected {
890 device: DeviceId::new(identifier),
891 reason: DisconnectReason::UserRequested,
892 });
893 release_link(link, (slot, Arc::clone(held))).await
894 }
895
896 pub(crate) async fn remove_device(self: &Arc<Self>, identifier: &str) -> Result<()> {
897 let Some(op) = self.withdraw(identifier).await else {
898 return Ok(());
899 };
900 let core = Arc::clone(self);
903 let identifier = identifier.to_owned();
904 finished(tokio::spawn(async move {
905 let held = Arc::new(Arc::clone(&op).lock_owned().await);
906 let result = core.disconnect_locked(&identifier, &held, false).await;
909 let mut devices = core.devices.write().await;
912 if devices
913 .get(&identifier)
914 .is_some_and(|entry| Arc::ptr_eq(&entry.op, &op))
915 {
916 devices.remove(&identifier);
917 info!("Removed device from manager: {identifier}");
918 }
919 result
920 }))
921 .await
922 }
923
924 pub(crate) async fn device_ids(&self) -> Vec<String> {
925 self.devices.read().await.keys().cloned().collect()
926 }
927
928 pub(crate) async fn device_count(&self) -> usize {
929 self.devices.read().await.len()
930 }
931
932 pub(crate) async fn connected_count(&self) -> usize {
933 let devices = self.devices.read().await;
934 devices
935 .values()
936 .filter(|entry| entry.link.is_some())
937 .count()
938 }
939
940 pub(crate) async fn can_connect(&self) -> bool {
941 self.slots
942 .as_ref()
943 .is_none_or(|slots| slots.available_permits() > 0)
944 }
945
946 pub(crate) async fn connection_status(&self) -> (usize, usize) {
947 (
948 self.connected_count().await,
949 self.config.max_concurrent_connections,
950 )
951 }
952
953 pub(crate) async fn available_connections(&self) -> Option<usize> {
954 self.slots.as_ref().map(|slots| slots.available_permits())
955 }
956
957 pub(crate) async fn connected_count_verified(&self) -> usize {
958 let links: Vec<Arc<L>> = {
960 let devices = self.devices.read().await;
961 devices
962 .values()
963 .filter_map(|entry| entry.link.clone())
964 .collect()
965 };
966 let results = join_all(links.iter().map(|link| link.is_connected())).await;
970
971 results.into_iter().filter(|&connected| connected).count()
972 }
973
974 pub(crate) async fn read_current(&self, identifier: &str) -> Result<CurrentReading> {
975 let link = {
977 let devices = self.devices.read().await;
978 let entry = devices
979 .get(identifier)
980 .ok_or_else(|| Error::device_not_found(identifier))?;
981 entry.link.clone().ok_or(Error::NotConnected)?
982 };
983 let reading = link.read_current().await?;
986
987 self.events.send(DeviceEvent::Reading {
989 device: DeviceId::new(identifier),
990 reading,
991 });
992
993 {
995 let mut devices = self.devices.write().await;
996 if let Some(entry) = devices.get_mut(identifier) {
997 entry.last_reading = Some(reading);
998 }
999 }
1000
1001 Ok(reading)
1002 }
1003
1004 pub(crate) async fn read_all(&self) -> HashMap<String, Result<CurrentReading>> {
1005 let links: Vec<(String, Arc<L>)> = {
1007 let devices = self.devices.read().await;
1008 devices
1009 .iter()
1010 .filter_map(|(id, entry)| entry.link.clone().map(|link| (id.clone(), link)))
1011 .collect()
1012 };
1013 let read_futures = links.into_iter().map(|(id, link)| async move {
1017 let result = link.read_current().await;
1018 (id, result)
1019 });
1020
1021 let read_results: Vec<(String, Result<CurrentReading>)> = join_all(read_futures).await;
1022
1023 for (id, result) in &read_results {
1025 if let Ok(reading) = result {
1026 self.events.send(DeviceEvent::Reading {
1027 device: DeviceId::new(id),
1028 reading: *reading,
1029 });
1030 }
1031 }
1032
1033 {
1035 let mut devices = self.devices.write().await;
1036 for (id, result) in &read_results {
1037 if let Ok(reading) = result
1038 && let Some(entry) = devices.get_mut(id)
1039 {
1040 entry.last_reading = Some(*reading);
1041 }
1042 }
1043 }
1044
1045 read_results.into_iter().collect()
1046 }
1047
1048 pub(crate) async fn connect_all(&self) -> HashMap<String, Result<()>> {
1049 let ids: Vec<_> = self.devices.read().await.keys().cloned().collect();
1050
1051 let connect_futures = ids.into_iter().map(|id| async move {
1054 let result = self.connect(&id).await;
1055 (id, result)
1056 });
1057
1058 join_all(connect_futures).await.into_iter().collect()
1059 }
1060
1061 pub(crate) async fn disconnect_all(self: &Arc<Self>) -> HashMap<String, Result<()>> {
1062 let withdrawn: Vec<(String, Arc<Mutex<()>>)> = {
1063 let mut devices = self.devices.write().await;
1064 for entry in devices.values_mut() {
1068 entry.wanted = false;
1069 }
1070 devices
1071 .iter()
1072 .filter(|(_, entry)| entry.link.is_some())
1073 .map(|(id, entry)| (id.clone(), Arc::clone(&entry.op)))
1074 .collect()
1075 };
1076 let disconnects: Vec<_> = withdrawn
1079 .into_iter()
1080 .map(|(id, op)| {
1081 let task = self.spawn_disconnect(&id, op);
1082 async move { (id, finished(task).await) }
1083 })
1084 .collect();
1085 join_all(disconnects).await.into_iter().collect()
1086 }
1087
1088 pub(crate) fn try_is_connected(&self, identifier: &str) -> Option<bool> {
1089 match self.devices.try_read() {
1091 Ok(devices) => Some(
1092 devices
1093 .get(identifier)
1094 .is_some_and(|entry| entry.link.is_some()),
1095 ),
1096 Err(_) => None, }
1098 }
1099
1100 pub(crate) async fn is_connected(&self, identifier: &str) -> bool {
1101 let link = {
1102 let devices = self.devices.read().await;
1103 devices.get(identifier).and_then(|entry| entry.link.clone())
1104 };
1105
1106 if let Some(link) = link {
1107 link.is_connected().await
1108 } else {
1109 false
1110 }
1111 }
1112
1113 pub(crate) async fn get_device_info(&self, identifier: &str) -> Option<DeviceInfo> {
1114 let devices = self.devices.read().await;
1115 devices.get(identifier).and_then(|entry| entry.info.clone())
1116 }
1117
1118 pub(crate) async fn get_last_reading(&self, identifier: &str) -> Option<CurrentReading> {
1119 let devices = self.devices.read().await;
1120 devices.get(identifier).and_then(|entry| entry.last_reading)
1121 }
1122
1123 pub(crate) async fn health_tick(&self) -> TickOutcome {
1130 let mut outcome = TickOutcome::default();
1131
1132 let targets: Vec<_> = {
1135 let devices = self.devices.read().await;
1136 devices
1137 .iter()
1138 .filter(|(_, entry)| entry.wanted)
1140 .filter_map(|(id, entry)| {
1141 let link = entry.link.as_ref()?;
1142 Some((id.clone(), Arc::clone(link), Arc::clone(&entry.op)))
1143 })
1144 .collect()
1145 };
1146 let mut probes = Vec::with_capacity(targets.len());
1147 for (id, link, op) in targets {
1148 probes.push(self.probe(id, link, op));
1149 }
1150 outcome.healthy = join_all(probes)
1151 .await
1152 .into_iter()
1153 .filter(|&healthy| healthy)
1154 .count();
1155
1156 let mut due: Vec<_> = {
1159 let devices = self.devices.read().await;
1160 devices
1161 .iter()
1162 .filter(|(_, entry)| entry.is_due_for_repair())
1163 .map(|(id, entry)| (id.clone(), entry.priority, Arc::clone(&entry.op)))
1164 .collect()
1165 };
1166 due.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
1167 for (id, _priority, op) in due {
1168 let Ok(guard) = Arc::clone(&op).try_lock_owned() else {
1170 continue;
1171 };
1172 let held = Arc::new(guard);
1173 let next = match self.devices.write().await.get_mut(&id) {
1178 Some(entry) if Arc::ptr_eq(&entry.op, &op) && entry.is_due_for_repair() => {
1179 match entry.give_up_before_reconnect() {
1180 Some(failures) => Err(failures),
1181 None => Ok(entry.failures.saturating_add(1)),
1182 }
1183 }
1184 _ => continue,
1185 };
1186 let attempt = match next {
1187 Ok(attempt) => attempt,
1188 Err(failures) => {
1189 self.report_give_up(&id, failures);
1190 continue;
1191 }
1192 };
1193 if self
1197 .slots
1198 .as_ref()
1199 .is_some_and(|slots| slots.available_permits() == 0)
1200 {
1201 debug!("Health monitor: no free connection slot for {id}");
1202 continue;
1203 }
1204 debug!("Health monitor: reconnecting {id} (attempt {attempt})");
1205 self.events.send(DeviceEvent::ReconnectStarted {
1206 device: DeviceId::new(&id),
1207 attempt,
1208 });
1209 match self.connect_locked(&id, &held, None).await {
1210 Ok(()) => {
1213 info!("Health monitor: reconnected {id}");
1214 self.events.send(DeviceEvent::ReconnectSucceeded {
1215 device: DeviceId::new(&id),
1216 attempts: attempt,
1217 });
1218 outcome.repaired += 1;
1219 }
1220 Err(Error::Cancelled) => debug!("Health monitor: reconnect of {id} cancelled"),
1222 Err(e) => {
1223 warn!("Health monitor: reconnect {attempt} of {id} failed: {e}");
1224 let gave_up_after = self
1225 .devices
1226 .write()
1227 .await
1228 .get_mut(&id)
1229 .and_then(|entry| entry.record_repair_failure(Instant::now()));
1230 if let Some(failures) = gave_up_after {
1231 self.report_give_up(&id, failures);
1232 }
1233 outcome.failed += 1;
1234 }
1235 }
1236 }
1237 outcome
1238 }
1239
1240 fn report_give_up(&self, id: &str, failures: u32) {
1243 let attempts = if failures == 1 { "attempt" } else { "attempts" };
1244 warn!(
1245 "Health monitor: giving up on {id} after {failures} failed {attempts}; \
1246 connect() starts over"
1247 );
1248 self.events.send(DeviceEvent::Error {
1249 device: DeviceId::new(id),
1250 error: format!("auto-reconnect gave up after {failures} {attempts}"),
1251 });
1252 }
1253
1254 async fn probe(&self, id: String, link: Arc<L>, op: Arc<Mutex<()>>) -> bool {
1257 let Ok(guard) = op.try_lock_owned() else {
1259 return false;
1260 };
1261 let held = Arc::new(guard);
1262 let alive = if self.config.use_connection_validation {
1263 link.is_alive().await
1264 } else {
1265 link.is_connected().await
1266 };
1267 if !alive {
1268 warn!("Health monitor: the connection to {id} is dead; closing it");
1269 self.detach(&id, &link, &held, DisconnectReason::Unknown)
1270 .await;
1271 }
1272 alive
1273 }
1274
1275 async fn detach(
1283 &self,
1284 identifier: &str,
1285 link: &Arc<L>,
1286 held: &Arc<tokio::sync::OwnedMutexGuard<()>>,
1287 reason: DisconnectReason,
1288 ) {
1289 let taken = {
1290 let mut devices = self.devices.write().await;
1291 match devices.get_mut(identifier) {
1292 Some(entry)
1293 if entry
1294 .link
1295 .as_ref()
1296 .is_some_and(|stored| Arc::ptr_eq(stored, link)) =>
1297 {
1298 entry.retry_at = Some(Instant::now());
1300 entry.link.take().map(|stored| (stored, entry.slot.take()))
1301 }
1302 _ => None,
1303 }
1304 };
1305 let Some((link, slot)) = taken else {
1306 return;
1307 };
1308 self.events.send(DeviceEvent::Disconnected {
1312 device: DeviceId::new(identifier),
1313 reason,
1314 });
1315 if let Err(e) = release_link(link, (slot, Arc::clone(held))).await {
1318 log_failed_release("Closing the dead connection", identifier, &e);
1319 }
1320 }
1321
1322 pub(crate) fn spawn_health_monitor(
1323 self: Arc<Self>,
1324 cancel: CancellationToken,
1325 ) -> tokio::task::JoinHandle<()> {
1326 tokio::spawn(async move {
1327 let mut adaptive = self.config.use_adaptive_interval.then(|| {
1328 AdaptiveInterval::new(
1329 self.config.health_check_interval,
1330 self.config.min_health_check_interval,
1331 self.config.max_health_check_interval,
1332 )
1333 });
1334 loop {
1335 let interval = adaptive
1336 .as_ref()
1337 .map_or(self.config.health_check_interval, AdaptiveInterval::current);
1338 if cancel
1339 .run_until_cancelled(tokio::time::sleep(interval))
1340 .await
1341 .is_none()
1342 {
1343 break;
1344 }
1345 let Some(outcome) = cancel.run_until_cancelled(self.health_tick()).await else {
1350 break;
1351 };
1352 if let Some(adaptive) = adaptive.as_mut() {
1353 if outcome.failed > 0 {
1354 adaptive.on_failure();
1355 } else if outcome.healthy > 0 && outcome.repaired == 0 {
1356 adaptive.on_success();
1357 }
1358 }
1359 }
1360 info!("Health monitor cancelled, shutting down");
1361 })
1362 }
1363
1364 pub(crate) async fn add_device_with_priority(
1365 &self,
1366 identifier: &str,
1367 priority: DevicePriority,
1368 ) -> Result<()> {
1369 self.config.default_reconnect_options.validate()?;
1370 let mut devices = self.devices.write().await;
1371
1372 if let Some(entry) = devices.get_mut(identifier) {
1373 entry.priority = priority;
1375 return Ok(());
1376 }
1377
1378 devices.insert(
1379 identifier.to_string(),
1380 Entry::new(self.config.default_reconnect_options.clone(), priority),
1381 );
1382
1383 info!(
1384 "Added device to manager with priority {:?}: {}",
1385 priority, identifier
1386 );
1387 Ok(())
1388 }
1389
1390 pub(crate) async fn lowest_priority_connected(&self) -> Option<String> {
1391 let devices = self.devices.read().await;
1392 devices
1393 .iter()
1394 .filter(|(_, entry)| entry.link.is_some() && entry.priority != DevicePriority::Critical)
1395 .min_by_key(|(_, entry)| entry.priority)
1396 .map(|(id, _)| id.clone())
1397 }
1398
1399 pub(crate) async fn evict_lowest_priority(self: &Arc<Self>) -> Result<bool> {
1400 if let Some(id) = self.lowest_priority_connected().await {
1401 info!("Evicting lowest priority device: {}", id);
1402 self.disconnect(&id).await?;
1403 Ok(true)
1404 } else {
1405 Ok(false)
1406 }
1407 }
1408
1409 pub(crate) fn spawn_hybrid_monitor(
1410 self: Arc<Self>,
1411 cancel: CancellationToken,
1412 options: PassiveMonitorOptions,
1413 ) -> tokio::task::JoinHandle<()> {
1414 tokio::spawn(async move {
1415 info!("Starting hybrid monitor (passive + active)");
1416
1417 let passive_monitor = Arc::new(PassiveMonitor::new(options));
1419 let mut passive_rx = passive_monitor.subscribe();
1420
1421 let passive_cancel = cancel.clone();
1423 let _passive_handle = passive_monitor.start(passive_cancel);
1424
1425 loop {
1426 tokio::select! {
1427 _ = cancel.cancelled() => {
1428 info!("Hybrid monitor cancelled");
1429 break;
1430 }
1431 result = passive_rx.recv() => {
1432 match result {
1433 Ok(passive_reading) => {
1434 if let Some(reading) = passive_reading_to_current(&passive_reading) {
1436 if let Some(entry) = self.devices.write().await.get_mut(&passive_reading.device_id) {
1438 entry.last_reading = Some(reading);
1439 }
1440
1441 self.events.send(DeviceEvent::Reading {
1443 device: DeviceId {
1444 id: passive_reading.device_id.clone(),
1445 name: passive_reading.device_name.clone(),
1446 device_type: Some(passive_reading.data.device_type),
1447 },
1448 reading,
1449 });
1450 }
1451 }
1452 Err(tokio::sync::broadcast::error::RecvError::Lagged(n)) => {
1453 warn!("Hybrid monitor lagged {} messages", n);
1454 }
1455 Err(tokio::sync::broadcast::error::RecvError::Closed) => {
1456 info!("Passive monitor channel closed");
1457 break;
1458 }
1459 }
1460 }
1461 }
1462 }
1463 })
1464 }
1465
1466 pub(crate) async fn read_hybrid(
1467 &self,
1468 identifier: &str,
1469 max_passive_age: Option<Duration>,
1470 ) -> Result<CurrentReading> {
1471 let max_age = max_passive_age.unwrap_or(Duration::from_secs(60));
1472
1473 {
1475 let devices = self.devices.read().await;
1476 if let Some(entry) = devices.get(identifier)
1477 && let Some(reading) = entry.last_reading
1478 {
1479 if let Some(captured) = reading.captured_at {
1481 let age = time::OffsetDateTime::now_utc() - captured;
1482 if age
1483 < time::Duration::try_from(max_age).unwrap_or(time::Duration::seconds(60))
1484 {
1485 debug!("Using cached passive reading for {}", identifier);
1486 return Ok(reading);
1487 }
1488 }
1489 }
1490 }
1491
1492 debug!(
1494 "No recent passive reading, using active connection for {}",
1495 identifier
1496 );
1497 self.read_current(identifier).await
1498 }
1499
1500 #[cfg(test)]
1502 pub(crate) async fn snapshot(&self, identifier: &str) -> Option<EntrySnapshot> {
1503 let devices = self.devices.read().await;
1504 devices.get(identifier).map(|entry| EntrySnapshot {
1505 has_link: entry.link.is_some(),
1506 failures: entry.failures,
1507 wanted: entry.wanted,
1508 gave_up: entry.gave_up,
1509 retry_at: entry.retry_at,
1510 })
1511 }
1512}
1513
1514pub struct DeviceManager {
1516 core: Arc<ManagerCore<Device>>,
1517}
1518
1519impl DeviceManager {
1520 pub fn new() -> Self {
1522 Self::with_config(ManagerConfig::default())
1523 }
1524
1525 pub fn with_event_capacity(capacity: usize) -> Self {
1527 Self::with_config(ManagerConfig {
1528 event_capacity: capacity,
1529 ..Default::default()
1530 })
1531 }
1532
1533 pub fn with_config(config: ManagerConfig) -> Self {
1535 Self {
1536 core: Arc::new(ManagerCore::new(config, ble_connector())),
1537 }
1538 }
1539
1540 pub fn events(&self) -> &EventDispatcher {
1542 &self.core.events
1543 }
1544
1545 pub fn config(&self) -> &ManagerConfig {
1547 &self.core.config
1548 }
1549
1550 pub async fn scan(&self) -> Result<Vec<DiscoveredDevice>> {
1552 scan_with_options(self.core.config.scan_options.clone()).await
1553 }
1554
1555 pub async fn scan_with_options(&self, options: ScanOptions) -> Result<Vec<DiscoveredDevice>> {
1557 let devices = scan_with_options(options).await?;
1558
1559 for device in &devices {
1561 self.core.events.send(DeviceEvent::Discovered {
1562 device: DeviceId {
1563 id: device.identifier.clone(),
1564 name: device.name.clone(),
1565 device_type: device.device_type,
1566 },
1567 rssi: device.rssi,
1568 });
1569 }
1570
1571 Ok(devices)
1572 }
1573
1574 pub async fn add_device(&self, identifier: &str) -> Result<()> {
1584 self.core.add_device(identifier).await
1585 }
1586
1587 pub async fn add_device_with_options(
1601 &self,
1602 identifier: &str,
1603 reconnect_options: ReconnectOptions,
1604 ) -> Result<()> {
1605 self.core
1606 .add_device_with_options(identifier, reconnect_options)
1607 .await
1608 }
1609
1610 pub async fn connect(&self, identifier: &str) -> Result<()> {
1655 self.core.connect(identifier).await
1656 }
1657
1658 pub async fn disconnect(&self, identifier: &str) -> Result<()> {
1684 self.core.disconnect(identifier).await
1685 }
1686
1687 pub async fn remove_device(&self, identifier: &str) -> Result<()> {
1704 self.core.remove_device(identifier).await
1705 }
1706
1707 pub async fn device_ids(&self) -> Vec<String> {
1709 self.core.device_ids().await
1710 }
1711
1712 pub async fn device_count(&self) -> usize {
1714 self.core.device_count().await
1715 }
1716
1717 pub async fn connected_count(&self) -> usize {
1723 self.core.connected_count().await
1724 }
1725
1726 pub async fn can_connect(&self) -> bool {
1732 self.core.can_connect().await
1733 }
1734
1735 pub async fn connection_status(&self) -> (usize, usize) {
1744 self.core.connection_status().await
1745 }
1746
1747 pub async fn available_connections(&self) -> Option<usize> {
1752 self.core.available_connections().await
1753 }
1754
1755 pub async fn connected_count_verified(&self) -> usize {
1760 self.core.connected_count_verified().await
1761 }
1762
1763 pub async fn read_current(&self, identifier: &str) -> Result<CurrentReading> {
1765 self.core.read_current(identifier).await
1766 }
1767
1768 pub async fn read_all(&self) -> HashMap<String, Result<CurrentReading>> {
1774 self.core.read_all().await
1775 }
1776
1777 pub async fn connect_all(&self) -> HashMap<String, Result<()>> {
1781 self.core.connect_all().await
1782 }
1783
1784 pub async fn disconnect_all(&self) -> HashMap<String, Result<()>> {
1792 self.core.disconnect_all().await
1793 }
1794
1795 pub fn try_is_connected(&self, identifier: &str) -> Option<bool> {
1805 self.core.try_is_connected(identifier)
1806 }
1807
1808 pub async fn is_connected(&self, identifier: &str) -> bool {
1812 self.core.is_connected(identifier).await
1813 }
1814
1815 pub async fn get_device_info(&self, identifier: &str) -> Option<DeviceInfo> {
1817 self.core.get_device_info(identifier).await
1818 }
1819
1820 pub async fn get_last_reading(&self, identifier: &str) -> Option<CurrentReading> {
1822 self.core.get_last_reading(identifier).await
1823 }
1824
1825 pub fn start_health_monitor(
1896 self: &Arc<Self>,
1897 cancel_token: CancellationToken,
1898 ) -> tokio::task::JoinHandle<()> {
1899 Arc::clone(&self.core).spawn_health_monitor(cancel_token)
1900 }
1901
1902 pub async fn add_device_with_priority(
1912 &self,
1913 identifier: &str,
1914 priority: DevicePriority,
1915 ) -> Result<()> {
1916 self.core
1917 .add_device_with_priority(identifier, priority)
1918 .await
1919 }
1920
1921 pub async fn lowest_priority_connected(&self) -> Option<String> {
1925 self.core.lowest_priority_connected().await
1926 }
1927
1928 pub async fn evict_lowest_priority(&self) -> Result<bool> {
1938 self.core.evict_lowest_priority().await
1939 }
1940
1941 pub fn start_hybrid_monitor(
1970 self: &Arc<Self>,
1971 cancel_token: CancellationToken,
1972 passive_options: Option<PassiveMonitorOptions>,
1973 ) -> tokio::task::JoinHandle<()> {
1974 Arc::clone(&self.core)
1975 .spawn_hybrid_monitor(cancel_token, passive_options.unwrap_or_default())
1976 }
1977
1978 pub async fn read_hybrid(
1988 &self,
1989 identifier: &str,
1990 max_passive_age: Option<Duration>,
1991 ) -> Result<CurrentReading> {
1992 self.core.read_hybrid(identifier, max_passive_age).await
1993 }
1994
1995 pub async fn supports_passive_monitoring(&self, identifier: &str) -> bool {
2009 let options = PassiveMonitorOptions::default()
2011 .scan_duration(Duration::from_secs(5))
2012 .filter_devices(vec![identifier.to_string()]);
2013
2014 let monitor = Arc::new(PassiveMonitor::new(options));
2015 let readings = monitor.subscribe();
2016
2017 receives_within(readings, Duration::from_secs(15), |cancel| {
2020 monitor.start(cancel)
2021 })
2022 .await
2023 }
2024}
2025
2026async fn receives_within<T: Clone>(
2031 mut readings: tokio::sync::broadcast::Receiver<T>,
2032 limit: Duration,
2033 start: impl FnOnce(CancellationToken) -> tokio::task::JoinHandle<()>,
2034) -> bool {
2035 let cancel = CancellationToken::new();
2036 let _stop_on_drop = cancel.clone().drop_guard();
2039 let _monitor = start(cancel);
2040 matches!(
2041 tokio::time::timeout(limit, readings.recv()).await,
2042 Ok(Ok(_))
2043 )
2044}
2045
2046fn passive_reading_to_current(passive: &PassiveReading) -> Option<CurrentReading> {
2048 let data = &passive.data;
2049
2050 if data.co2.is_none()
2052 && data.temperature.is_none()
2053 && data.humidity.is_none()
2054 && data.radon.is_none()
2055 && data.radiation_dose_rate.is_none()
2056 {
2057 return None;
2058 }
2059
2060 Some(CurrentReading {
2061 co2: data.co2.unwrap_or(0),
2062 temperature: data.temperature.unwrap_or(0.0),
2063 pressure: data.pressure.unwrap_or(0.0),
2064 humidity: data.humidity.unwrap_or(0),
2065 battery: data.battery,
2066 status: data.status,
2067 interval: data.interval,
2068 age: data.age,
2069 captured_at: Some(time::OffsetDateTime::now_utc()),
2070 radon: data.radon,
2071 radon_avg_24h: None,
2072 radon_avg_7d: None,
2073 radon_avg_30d: None,
2074 radiation_rate: data.radiation_dose_rate,
2075 radiation_total: None, })
2077}
2078
2079impl Default for DeviceManager {
2080 fn default() -> Self {
2081 Self::new()
2082 }
2083}
2084
2085#[cfg(test)]
2086mod tests {
2087 use super::*;
2088 use crate::test_support::within;
2089
2090 #[tokio::test]
2091 async fn test_manager_add_device() {
2092 let manager = DeviceManager::new();
2093 manager.add_device("test-device").await.unwrap();
2094
2095 assert_eq!(manager.device_count().await, 1);
2096 assert!(
2097 manager
2098 .device_ids()
2099 .await
2100 .contains(&"test-device".to_string())
2101 );
2102 }
2103
2104 #[tokio::test]
2105 async fn test_manager_remove_device() {
2106 let manager = DeviceManager::new();
2107 manager.add_device("test-device").await.unwrap();
2108 manager.remove_device("test-device").await.unwrap();
2109
2110 assert_eq!(manager.device_count().await, 0);
2111 }
2112
2113 #[tokio::test]
2114 async fn test_manager_not_connected_by_default() {
2115 let manager = DeviceManager::new();
2116 manager.add_device("test-device").await.unwrap();
2117
2118 assert!(!manager.is_connected("test-device").await);
2119 assert_eq!(manager.connected_count().await, 0);
2120 }
2121
2122 #[tokio::test]
2123 async fn test_manager_events() {
2124 let manager = DeviceManager::new();
2125 let _rx = manager.events().subscribe();
2126
2127 manager.add_device("test-device").await.unwrap();
2128
2129 assert_eq!(manager.events().receiver_count(), 1);
2131 }
2132
2133 fn fake_monitor(
2137 cancel: CancellationToken,
2138 sender: tokio::sync::broadcast::Sender<()>,
2139 reading: bool,
2140 stopped: tokio::sync::oneshot::Sender<()>,
2141 ) -> tokio::task::JoinHandle<()> {
2142 tokio::spawn(async move {
2143 if reading {
2144 tokio::time::sleep(Duration::from_secs(3)).await;
2145 let _ = sender.send(());
2146 }
2147 cancel.cancelled().await;
2148 let _ = stopped.send(());
2149 })
2150 }
2151
2152 #[tokio::test(start_paused = true)]
2155 async fn the_passive_check_stops_its_monitor_when_it_returns() {
2156 within(Duration::from_secs(600), async {
2157 for (reading, after) in [(true, 3), (false, 15)] {
2158 let (sender, readings) = tokio::sync::broadcast::channel(1);
2159 let (stopped, monitor_stopped) = tokio::sync::oneshot::channel();
2160 let start = tokio::time::Instant::now();
2161
2162 let found = receives_within(readings, Duration::from_secs(15), |cancel| {
2163 fake_monitor(cancel, sender, reading, stopped)
2164 })
2165 .await;
2166
2167 assert_eq!(found, reading);
2168 assert_eq!(start.elapsed(), Duration::from_secs(after));
2169 within(Duration::from_secs(10), monitor_stopped)
2170 .await
2171 .expect("the monitor ended without being stopped");
2172 }
2173 })
2174 .await;
2175 }
2176
2177 #[tokio::test(start_paused = true)]
2181 async fn a_dropped_passive_check_stops_its_monitor() {
2182 within(Duration::from_secs(600), async {
2183 let (sender, readings) = tokio::sync::broadcast::channel(1);
2184 let (stopped, monitor_stopped) = tokio::sync::oneshot::channel();
2185 let check = receives_within(readings, Duration::from_secs(15), |cancel| {
2186 fake_monitor(cancel, sender, false, stopped)
2187 });
2188
2189 let given_up = tokio::time::timeout(Duration::from_secs(1), check).await;
2190 assert!(given_up.is_err(), "the check returned {given_up:?}");
2191 within(Duration::from_secs(10), monitor_stopped)
2192 .await
2193 .expect("the monitor ended without being stopped");
2194 })
2195 .await;
2196 }
2197}
2198
2199#[cfg(test)]
2200mod lifecycle_tests {
2201 use std::sync::Arc;
2202 use std::time::Duration;
2203
2204 use tokio::time::timeout;
2205
2206 use super::{DevicePriority, EntrySnapshot, ManagerConfig, ManagerCore, TickOutcome};
2207 use crate::error::{ConnectionFailureReason, Error, Result};
2208 use crate::events::{DeviceEvent, EventReceiver};
2209 use crate::test_support::{FakeConn, FakeEvent, FakeRadio, within};
2210
2211 const TEST_LIMIT: Duration = Duration::from_secs(600);
2212
2213 fn core(radio: &FakeRadio, config: ManagerConfig) -> Arc<ManagerCore<FakeConn>> {
2214 Arc::new(ManagerCore::new(config, radio.connector()))
2215 }
2216
2217 fn drain(events: &mut EventReceiver) -> Vec<String> {
2220 let mut lines = Vec::new();
2221 while let Ok(event) = events.try_recv() {
2222 lines.push(match event {
2223 DeviceEvent::Connected { device, .. } => format!("Connected {}", device.id),
2224 DeviceEvent::Disconnected { device, reason } => {
2225 format!("Disconnected {} {reason:?}", device.id)
2226 }
2227 DeviceEvent::ReconnectStarted { device, attempt } => {
2228 format!("ReconnectStarted {} {attempt}", device.id)
2229 }
2230 DeviceEvent::ReconnectSucceeded { device, attempts } => {
2231 format!("ReconnectSucceeded {} {attempts}", device.id)
2232 }
2233 DeviceEvent::Error { device, error } => format!("Error {}: {error}", device.id),
2234 other => format!("{other:?}"),
2235 });
2236 }
2237 lines
2238 }
2239
2240 async fn assert_no_orphans(core: &ManagerCore<FakeConn>, radio: &FakeRadio) {
2242 for id in radio.up_ids() {
2243 assert!(
2244 core.snapshot(&id).await.is_some_and(|entry| entry.has_link),
2245 "{id} is connected but the manager holds no link for it: {:#?}",
2246 radio.events()
2247 );
2248 }
2249 }
2250
2251 fn is_limit_error(result: &Result<()>) -> bool {
2252 matches!(
2253 result,
2254 Err(Error::ConnectionFailed {
2255 reason: ConnectionFailureReason::Other(message),
2256 ..
2257 }) if message.starts_with("Connection limit reached")
2258 )
2259 }
2260
2261 #[tokio::test(start_paused = true)]
2264 async fn timed_out_connect_does_not_wedge_the_device() {
2265 within(TEST_LIMIT, async {
2266 let radio = FakeRadio::new();
2267 let core = core(&radio, ManagerConfig::default());
2268
2269 radio.set_connect_delay("A", Duration::from_secs(10));
2270 assert!(
2271 timeout(Duration::from_secs(1), core.connect("A"))
2272 .await
2273 .is_err()
2274 );
2275
2276 radio.set_connect_delay("A", Duration::ZERO);
2277 core.connect("A").await.expect("second connect");
2278 assert!(
2279 radio.link_up("A"),
2280 "the second connect returned Ok without connecting"
2281 );
2282 assert_eq!(radio.connect_count("A"), 2);
2283
2284 assert_no_orphans(&core, &radio).await;
2285 radio.assert_no_drop_teardown();
2286 })
2287 .await;
2288 }
2289
2290 #[tokio::test(start_paused = true)]
2293 async fn concurrent_connects_to_one_device_share_the_real_result() {
2294 within(TEST_LIMIT, async {
2295 let radio = FakeRadio::new();
2296 let core = core(&radio, ManagerConfig::default());
2297 radio.set_connect_delay("A", Duration::from_secs(5));
2298 radio.script_connects("A", [false, false]);
2299
2300 let (first, second) = tokio::join!(core.connect("A"), core.connect("A"));
2301 assert!(first.is_err(), "first connect: {first:?}");
2302 assert!(
2303 second.is_err(),
2304 "the second connect returned {second:?} although no connect succeeded"
2305 );
2306
2307 let a = || "A".to_string();
2309 let events: Vec<FakeEvent> =
2310 radio.events().into_iter().map(|(_, event)| event).collect();
2311 assert_eq!(
2312 events,
2313 [
2314 FakeEvent::ConnectStarted { id: a() },
2315 FakeEvent::ConnectFailed { id: a() },
2316 FakeEvent::ConnectStarted { id: a() },
2317 FakeEvent::ConnectFailed { id: a() },
2318 ]
2319 );
2320
2321 assert_no_orphans(&core, &radio).await;
2322 radio.assert_no_drop_teardown();
2323 })
2324 .await;
2325 }
2326
2327 #[tokio::test(start_paused = true)]
2330 async fn connection_limit_counts_in_flight_connects() {
2331 within(TEST_LIMIT, async {
2332 let radio = FakeRadio::new();
2333 let core = core(&radio, ManagerConfig::default().with_max_connections(1));
2334 radio.set_connect_delay("A", Duration::from_secs(5));
2335 radio.set_connect_delay("B", Duration::from_secs(5));
2336
2337 let (a, b) = tokio::join!(core.connect("A"), core.connect("B"));
2338 let results = [a, b];
2339 assert_eq!(
2340 results.iter().filter(|result| result.is_ok()).count(),
2341 1,
2342 "{results:?}"
2343 );
2344 assert_eq!(
2345 results
2346 .iter()
2347 .filter(|result| is_limit_error(result))
2348 .count(),
2349 1,
2350 "{results:?}"
2351 );
2352 assert_eq!(core.connected_count().await, 1);
2353 assert_eq!(core.device_count().await, 1);
2355 assert_eq!(radio.up_ids().len(), 1, "links up: {:?}", radio.up_ids());
2356
2357 assert_no_orphans(&core, &radio).await;
2358 radio.assert_no_drop_teardown();
2359 })
2360 .await;
2361 }
2362
2363 #[tokio::test(start_paused = true)]
2366 async fn connect_all_respects_the_connection_limit() {
2367 within(TEST_LIMIT, async {
2368 let radio = FakeRadio::new();
2369 let core = core(&radio, ManagerConfig::default().with_max_connections(2));
2370 for id in ["A", "B", "C", "D"] {
2371 radio.set_connect_delay(id, Duration::from_secs(5));
2372 core.add_device(id).await.expect("add_device");
2373 }
2374
2375 let results = core.connect_all().await;
2376 assert_eq!(results.len(), 4);
2377 assert_eq!(
2378 results.values().filter(|result| result.is_ok()).count(),
2379 2,
2380 "{results:?}"
2381 );
2382 assert_eq!(
2383 results
2384 .values()
2385 .filter(|result| is_limit_error(result))
2386 .count(),
2387 2,
2388 "{results:?}"
2389 );
2390 assert_eq!(core.connected_count().await, 2);
2391 assert_eq!(radio.up_ids().len(), 2, "links up: {:?}", radio.up_ids());
2392 assert_eq!(core.available_connections().await, Some(0));
2393
2394 assert_no_orphans(&core, &radio).await;
2395 radio.assert_no_drop_teardown();
2396 })
2397 .await;
2398 }
2399
2400 #[tokio::test(start_paused = true)]
2402 async fn failed_or_cancelled_connect_releases_its_slot() {
2403 within(TEST_LIMIT, async {
2404 let radio = FakeRadio::new();
2405 let core = core(&radio, ManagerConfig::default().with_max_connections(1));
2406
2407 radio.script_connects("A", [false]);
2408 assert!(core.connect("A").await.is_err());
2409 assert_eq!(core.available_connections().await, Some(1));
2410
2411 radio.set_connect_delay("B", Duration::from_secs(10));
2412 assert!(
2413 timeout(Duration::from_secs(1), core.connect("B"))
2414 .await
2415 .is_err()
2416 );
2417 assert_eq!(core.available_connections().await, Some(1));
2418
2419 core.connect("C").await.expect("C gets the free slot");
2420 assert!(radio.link_up("C"));
2421 assert_eq!(core.available_connections().await, Some(0));
2422 assert_eq!(
2423 core.snapshot("C").await,
2424 Some(EntrySnapshot {
2425 has_link: true,
2426 failures: 0,
2427 wanted: true,
2428 gave_up: false,
2429 retry_at: None,
2430 })
2431 );
2432
2433 assert_no_orphans(&core, &radio).await;
2434 radio.assert_no_drop_teardown();
2435 })
2436 .await;
2437 }
2438
2439 #[tokio::test(start_paused = true)]
2442 async fn disconnect_releases_slot_only_after_the_link_is_down() {
2443 within(TEST_LIMIT, async {
2444 let radio = FakeRadio::new();
2445 let core = core(&radio, ManagerConfig::default().with_max_connections(1));
2446 core.connect("A").await.expect("connect");
2447 radio.set_disconnect_delay("A", Duration::from_secs(2));
2448
2449 let disconnect = tokio::spawn({
2450 let core = Arc::clone(&core);
2451 async move { core.disconnect("A").await }
2452 });
2453 tokio::time::sleep(Duration::from_secs(1)).await;
2454 assert!(radio.link_up("A"), "the disconnect is still running");
2455 assert_eq!(
2456 core.available_connections().await,
2457 Some(0),
2458 "the slot was freed while the sensor was still connected"
2459 );
2460
2461 disconnect
2462 .await
2463 .expect("disconnect task")
2464 .expect("disconnect");
2465 assert!(!radio.link_up("A"));
2466 assert_eq!(core.available_connections().await, Some(1));
2467
2468 assert_no_orphans(&core, &radio).await;
2469 radio.assert_no_drop_teardown();
2470 })
2471 .await;
2472 }
2473
2474 #[tokio::test(start_paused = true)]
2477 async fn cancel_during_device_info_read_leaves_no_dropped_link() {
2478 within(TEST_LIMIT, async {
2479 let radio = FakeRadio::new();
2480 let core = core(&radio, ManagerConfig::default().with_max_connections(1));
2481 radio.set_info_delay("A", Duration::from_secs(5));
2482 let mut events = core.events.subscribe();
2483
2484 assert!(
2485 timeout(Duration::from_secs(1), core.connect("A"))
2486 .await
2487 .is_err()
2488 );
2489
2490 assert!(
2491 core.snapshot("A").await.is_some_and(|entry| entry.has_link),
2492 "the manager holds no link after the cancel: {:#?}",
2493 radio.events()
2494 );
2495 assert!(radio.link_up("A"));
2496 assert_eq!(core.available_connections().await, Some(0));
2497 assert!(core.get_device_info("A").await.is_none());
2500 let sent = events.try_recv();
2501 assert!(sent.is_err(), "unexpected event: {sent:?}");
2502
2503 assert_no_orphans(&core, &radio).await;
2504 radio.assert_no_drop_teardown();
2505 })
2506 .await;
2507 }
2508
2509 #[tokio::test(start_paused = true)]
2513 async fn rejected_connect_adds_no_device_while_the_map_is_busy() {
2514 within(TEST_LIMIT, async {
2515 let radio = FakeRadio::new();
2516 let core = core(&radio, ManagerConfig::default().with_max_connections(1));
2517 radio.set_connect_delay("A", Duration::from_secs(5));
2518 radio.set_connect_delay("B", Duration::from_secs(5));
2519
2520 let busy = core.devices.read().await;
2523 let spawn_connect = |id: &'static str| {
2524 let core = Arc::clone(&core);
2525 tokio::spawn(async move { core.connect(id).await })
2526 };
2527 let a = spawn_connect("A");
2528 let b = spawn_connect("B");
2529 tokio::time::sleep(Duration::from_secs(1)).await;
2530 drop(busy);
2531
2532 let a = a.await.expect("connect A task");
2533 let b = b.await.expect("connect B task");
2534 assert!(a.is_ok(), "A asked first and gets the only slot: {a:?}");
2535 assert!(is_limit_error(&b), "B is over the limit: {b:?}");
2536 assert_eq!(
2537 core.device_ids().await,
2538 ["A"],
2539 "a connect that the limit rejects doesn't add the device"
2540 );
2541 assert_eq!(radio.up_ids(), ["A"]);
2542
2543 assert_no_orphans(&core, &radio).await;
2544 radio.assert_no_drop_teardown();
2545 })
2546 .await;
2547 }
2548
2549 #[tokio::test(start_paused = true)]
2553 async fn connect_cancelled_while_storing_its_link_disconnects_it() {
2554 within(TEST_LIMIT, async {
2555 let radio = FakeRadio::new();
2556 let core = core(&radio, ManagerConfig::default().with_max_connections(1));
2557 radio.set_connect_delay("A", Duration::from_secs(2));
2558 radio.set_disconnect_delay("A", Duration::from_secs(2));
2559
2560 let connect = tokio::spawn({
2561 let core = Arc::clone(&core);
2562 async move { timeout(Duration::from_secs(3), core.connect("A")).await }
2563 });
2564 tokio::time::sleep(Duration::from_secs(1)).await;
2568 let busy = core.devices.read().await;
2569 assert!(connect.await.expect("connect task").is_err());
2570
2571 assert!(
2572 radio.link_up("A"),
2573 "the cancelled connect dropped its new link: {:#?}",
2574 radio.events()
2575 );
2576 assert_eq!(
2577 core.available_connections().await,
2578 Some(0),
2579 "the slot was freed while the sensor was still connected"
2580 );
2581 tokio::time::sleep(Duration::from_secs(3)).await;
2582 assert!(!radio.link_up("A"), "the new link was not disconnected");
2583 assert_eq!(core.available_connections().await, Some(1));
2584 drop(busy);
2585
2586 assert!(
2587 core.snapshot("A")
2588 .await
2589 .is_some_and(|entry| !entry.has_link)
2590 );
2591 assert_no_orphans(&core, &radio).await;
2592 radio.assert_no_drop_teardown();
2593 })
2594 .await;
2595 }
2596
2597 #[tokio::test(start_paused = true)]
2602 async fn connect_after_a_cancelled_disconnect_waits_for_the_link_to_go_down() {
2603 within(TEST_LIMIT, async {
2604 let radio = FakeRadio::new();
2605 let core = core(&radio, ManagerConfig::default());
2606 core.connect("A").await.expect("first connect");
2607 radio.set_disconnect_delay("A", Duration::from_secs(2));
2608
2609 assert!(
2610 timeout(Duration::from_millis(500), core.disconnect("A"))
2611 .await
2612 .is_err()
2613 );
2614 core.connect("A").await.expect("second connect");
2615 tokio::time::sleep(Duration::from_secs(3)).await;
2617
2618 assert!(
2619 radio.link_up("A"),
2620 "the old link's disconnect took down the new link: {:#?}",
2621 radio.events()
2622 );
2623 let a = || "A".to_string();
2625 let events: Vec<FakeEvent> =
2626 radio.events().into_iter().map(|(_, event)| event).collect();
2627 assert_eq!(
2628 events,
2629 [
2630 FakeEvent::ConnectStarted { id: a() },
2631 FakeEvent::Connected { id: a(), handle: 1 },
2632 FakeEvent::Disconnect { id: a(), handle: 1 },
2633 FakeEvent::ConnectStarted { id: a() },
2634 FakeEvent::Connected { id: a(), handle: 2 },
2635 ]
2636 );
2637
2638 assert_no_orphans(&core, &radio).await;
2639 radio.assert_no_drop_teardown();
2640 })
2641 .await;
2642 }
2643
2644 #[tokio::test(start_paused = true)]
2649 async fn connect_after_a_cancelled_connect_waits_for_its_link_to_go_down() {
2650 within(TEST_LIMIT, async {
2651 let radio = FakeRadio::new();
2652 let core = core(&radio, ManagerConfig::default());
2653 radio.set_connect_delay("A", Duration::from_secs(2));
2654 radio.set_disconnect_delay("A", Duration::from_secs(2));
2655
2656 let connect = tokio::spawn({
2659 let core = Arc::clone(&core);
2660 async move { timeout(Duration::from_secs(3), core.connect("A")).await }
2661 });
2662 tokio::time::sleep(Duration::from_secs(1)).await;
2663 let busy = core.devices.read().await;
2664 assert!(connect.await.expect("connect task").is_err());
2665 drop(busy);
2666
2667 radio.set_connect_delay("A", Duration::ZERO);
2668 core.connect("A").await.expect("second connect");
2669 tokio::time::sleep(Duration::from_secs(3)).await;
2671
2672 assert!(
2673 radio.link_up("A"),
2674 "A is down after the second connect: {:#?}",
2675 radio.events()
2676 );
2677 let a = || "A".to_string();
2679 let events: Vec<FakeEvent> =
2680 radio.events().into_iter().map(|(_, event)| event).collect();
2681 assert_eq!(
2682 events,
2683 [
2684 FakeEvent::ConnectStarted { id: a() },
2685 FakeEvent::Connected { id: a(), handle: 1 },
2686 FakeEvent::Disconnect { id: a(), handle: 1 },
2687 FakeEvent::ConnectStarted { id: a() },
2688 FakeEvent::Connected { id: a(), handle: 2 },
2689 ]
2690 );
2691
2692 assert_no_orphans(&core, &radio).await;
2693 radio.assert_no_drop_teardown();
2694 })
2695 .await;
2696 }
2697
2698 #[tokio::test(start_paused = true)]
2704 async fn disconnects_send_disconnected_even_if_the_close_fails() {
2705 within(TEST_LIMIT, async {
2706 let radio = FakeRadio::new();
2707 let core = core(&radio, ManagerConfig::default());
2708 core.add_device_with_priority("B", DevicePriority::Low)
2709 .await
2710 .expect("add_device_with_priority");
2711 for id in ["A", "B", "C"] {
2712 core.connect(id).await.expect("connect");
2713 radio.fail_disconnects(id);
2714 }
2715 let mut events = core.events.subscribe();
2716
2717 let result = core.disconnect("A").await;
2718 assert!(matches!(result, Err(Error::Timeout { .. })), "{result:?}");
2719 assert_eq!(drain(&mut events), ["Disconnected A UserRequested"]);
2720
2721 let result = core.evict_lowest_priority().await;
2723 assert!(matches!(result, Err(Error::Timeout { .. })), "{result:?}");
2724 assert_eq!(drain(&mut events), ["Disconnected B UserRequested"]);
2725
2726 let results = core.disconnect_all().await;
2728 assert!(
2729 results.len() == 1 && matches!(results.get("C"), Some(Err(Error::Timeout { .. }))),
2730 "{results:?}"
2731 );
2732 assert_eq!(drain(&mut events), ["Disconnected C UserRequested"]);
2733
2734 for id in ["A", "B", "C"] {
2735 assert_eq!(core.try_is_connected(id), Some(false), "{id}");
2736 }
2737 assert_no_orphans(&core, &radio).await;
2738 radio.assert_no_drop_teardown();
2739 })
2740 .await;
2741 }
2742
2743 #[tokio::test(start_paused = true)]
2747 async fn remove_device_removes_the_device_even_if_its_disconnect_fails() {
2748 within(TEST_LIMIT, async {
2749 let radio = FakeRadio::new();
2750 let core = core(&radio, ManagerConfig::default());
2751 core.connect("A").await.expect("connect");
2752 radio.fail_disconnects("A");
2753 let mut events = core.events.subscribe();
2754
2755 let result = core.remove_device("A").await;
2756
2757 assert!(matches!(result, Err(Error::Timeout { .. })), "{result:?}");
2758 assert_eq!(core.device_count().await, 0);
2759 assert_eq!(drain(&mut events), ["Disconnected A UserRequested"]);
2760 assert_eq!(core.health_tick().await, TickOutcome::default());
2762 assert_eq!(radio.connect_count("A"), 1);
2763 assert!(!radio.link_up("A"));
2764
2765 assert_no_orphans(&core, &radio).await;
2766 radio.assert_no_drop_teardown();
2767 })
2768 .await;
2769 }
2770
2771 mod health {
2774 use std::sync::Arc;
2775 use std::time::Duration;
2776
2777 use tokio::time::{Instant, timeout};
2778 use tokio_util::sync::CancellationToken;
2779
2780 use super::{assert_no_orphans, core, drain};
2781 use crate::error::Error;
2782 use crate::manager::{DevicePriority, EntrySnapshot, ManagerConfig, TickOutcome};
2783 use crate::reconnect::ReconnectOptions;
2784 use crate::test_support::{FakeEvent, FakeRadio, within};
2785
2786 const LIMIT: Duration = Duration::from_secs(600);
2787 const LONG_LIMIT: Duration = Duration::from_secs(2 * 60 * 60);
2789
2790 fn radio_log(radio: &FakeRadio) -> Vec<FakeEvent> {
2792 radio
2793 .events()
2794 .into_iter()
2795 .map(|(_, event)| event)
2796 .filter(|event| !matches!(event, FakeEvent::Op { .. }))
2797 .collect()
2798 }
2799
2800 fn connect_starts(radio: &FakeRadio, id: &str) -> Vec<Duration> {
2802 radio
2803 .events()
2804 .into_iter()
2805 .filter_map(|(at, event)| {
2806 matches!(&event, FakeEvent::ConnectStarted { id: started } if started == id)
2807 .then_some(at)
2808 })
2809 .collect()
2810 }
2811
2812 fn first_disconnect(radio: &FakeRadio, id: &str) -> Option<Duration> {
2814 radio.events().into_iter().find_map(|(at, event)| {
2815 matches!(&event, FakeEvent::Disconnect { id: closed, .. } if closed == id)
2816 .then_some(at)
2817 })
2818 }
2819
2820 #[tokio::test(start_paused = true)]
2823 async fn health_tick_replaces_a_dead_handle() {
2824 within(LIMIT, async {
2825 let radio = FakeRadio::new();
2826 let core = core(&radio, ManagerConfig::default());
2827 core.connect("A").await.unwrap();
2828 radio.lose_link("A");
2829 let mut events = core.events.subscribe();
2830
2831 let outcome = core.health_tick().await;
2832
2833 assert_eq!(
2834 outcome,
2835 TickOutcome {
2836 healthy: 0,
2837 repaired: 1,
2838 failed: 0
2839 }
2840 );
2841 let log = radio_log(&radio);
2842 let [
2843 FakeEvent::ConnectStarted { .. },
2844 FakeEvent::Connected { handle: first, .. },
2845 FakeEvent::Disconnect { handle: closed, .. },
2846 FakeEvent::ConnectStarted { .. },
2847 FakeEvent::Connected { handle: second, .. },
2848 ] = log.as_slice()
2849 else {
2850 panic!("unexpected radio log: {log:?}");
2851 };
2852 assert_eq!(closed, first, "the dead handle was not the one closed");
2853 assert_ne!(second, first);
2854 assert!(radio.link_up("A"));
2855 assert_eq!(
2856 drain(&mut events),
2857 [
2858 "Disconnected A Unknown",
2859 "ReconnectStarted A 1",
2860 "Connected A",
2861 "ReconnectSucceeded A 1"
2862 ]
2863 );
2864 core.read_current("A").await.unwrap();
2865
2866 assert_no_orphans(&core, &radio).await;
2867 radio.assert_no_drop_teardown();
2868 })
2869 .await;
2870 }
2871
2872 #[tokio::test(start_paused = true)]
2873 async fn connect_replaces_a_dead_link() {
2874 within(LIMIT, async {
2875 let radio = FakeRadio::new();
2876 let core = core(&radio, ManagerConfig::default());
2877 core.connect("A").await.unwrap();
2878 core.connect("A").await.unwrap();
2880 assert_eq!(radio.connect_count("A"), 1);
2881
2882 radio.lose_link("A");
2883 let mut events = core.events.subscribe();
2884 core.connect("A").await.unwrap();
2885
2886 assert_eq!(radio.connect_count("A"), 2, "connect() kept the dead link");
2887 let log = radio_log(&radio);
2888 let [
2889 FakeEvent::ConnectStarted { .. },
2890 FakeEvent::Connected { handle: first, .. },
2891 FakeEvent::Disconnect { handle: closed, .. },
2892 FakeEvent::ConnectStarted { .. },
2893 FakeEvent::Connected { handle: second, .. },
2894 ] = log.as_slice()
2895 else {
2896 panic!("unexpected radio log: {log:?}");
2897 };
2898 assert_eq!(closed, first, "the dead handle was not the one closed");
2899 assert_ne!(second, first);
2900 assert!(radio.link_up("A"));
2901 assert_eq!(
2902 drain(&mut events),
2903 ["Disconnected A Unknown", "Connected A"]
2904 );
2905 core.read_current("A").await.unwrap();
2906
2907 assert_no_orphans(&core, &radio).await;
2908 radio.assert_no_drop_teardown();
2909 })
2910 .await;
2911 }
2912
2913 #[tokio::test(start_paused = true)]
2914 async fn health_tick_replaces_a_zombie_handle_when_validation_is_on() {
2915 within(LIMIT, async {
2916 let radio = FakeRadio::new();
2917 let validating = core(&radio, ManagerConfig::default().connection_validation(true));
2918 validating.connect("A").await.unwrap();
2919 radio.make_zombie("A");
2920
2921 assert_eq!(
2922 validating.health_tick().await,
2923 TickOutcome {
2924 healthy: 0,
2925 repaired: 1,
2926 failed: 0
2927 }
2928 );
2929 assert_eq!(radio.connect_count("A"), 2);
2930 assert!(radio.link_up("A"));
2931 assert_eq!(
2933 validating.health_tick().await,
2934 TickOutcome {
2935 healthy: 1,
2936 repaired: 0,
2937 failed: 0
2938 }
2939 );
2940 validating.read_current("A").await.unwrap();
2941 assert_no_orphans(&validating, &radio).await;
2942 radio.assert_no_drop_teardown();
2943
2944 let radio = FakeRadio::new();
2947 let asking = core(
2948 &radio,
2949 ManagerConfig::default().connection_validation(false),
2950 );
2951 asking.connect("A").await.unwrap();
2952 radio.make_zombie("A");
2953 assert_eq!(
2954 asking.health_tick().await,
2955 TickOutcome {
2956 healthy: 1,
2957 repaired: 0,
2958 failed: 0
2959 }
2960 );
2961 assert_eq!(radio.connect_count("A"), 1);
2962 assert!(
2963 asking.read_current("A").await.is_err(),
2964 "the zombie's reads should time out"
2965 );
2966 assert_no_orphans(&asking, &radio).await;
2967 radio.assert_no_drop_teardown();
2968 })
2969 .await;
2970 }
2971
2972 #[tokio::test(start_paused = true)]
2973 async fn health_tick_counts_a_failed_repair_as_failure() {
2974 within(LIMIT, async {
2975 let radio = FakeRadio::new();
2976 let core = core(&radio, ManagerConfig::default());
2977 core.connect("A").await.unwrap();
2978 radio.lose_link("A");
2979 radio.script_connects("A", [false]);
2980
2981 assert_eq!(
2982 core.health_tick().await,
2983 TickOutcome {
2984 healthy: 0,
2985 repaired: 0,
2986 failed: 1
2987 }
2988 );
2989 let snapshot = core.snapshot("A").await.unwrap();
2990 assert!(!snapshot.has_link);
2991 assert_eq!(snapshot.failures, 1);
2992 assert!(!radio.link_up("A"));
2993
2994 assert_no_orphans(&core, &radio).await;
2995 radio.assert_no_drop_teardown();
2996 })
2997 .await;
2998 }
2999
3000 #[tokio::test(start_paused = true)]
3001 async fn health_monitor_tightens_interval_while_repairs_fail() {
3002 within(LIMIT, async {
3003 let radio = FakeRadio::new();
3004 let core = core(&radio, ManagerConfig::default());
3006 core.add_device_with_options(
3007 "A",
3008 ReconnectOptions {
3009 max_attempts: None,
3010 ..ReconnectOptions::fixed_delay(Duration::from_millis(100))
3011 },
3012 )
3013 .await
3014 .unwrap();
3015 core.connect("A").await.unwrap();
3016 core.connect("B").await.unwrap();
3019 radio.lose_link("A");
3020 radio.script_connects("A", [false; 10]);
3021
3022 let cancel = CancellationToken::new();
3023 let monitor = Arc::clone(&core).spawn_health_monitor(cancel.clone());
3024 tokio::time::sleep(Duration::from_secs(70)).await;
3025 cancel.cancel();
3026 monitor.await.unwrap();
3027
3028 let starts = connect_starts(&radio, "A");
3031 assert_eq!(
3032 starts[1..],
3033 [30_000, 45_000, 52_500, 57_500, 62_500, 67_500].map(Duration::from_millis)
3034 );
3035
3036 assert_no_orphans(&core, &radio).await;
3037 radio.assert_no_drop_teardown();
3038 })
3039 .await;
3040 }
3041
3042 #[tokio::test(start_paused = true)]
3043 async fn probes_run_concurrently_so_a_slow_device_does_not_delay_others() {
3044 within(LIMIT, async {
3045 let radio = FakeRadio::new();
3046 let core = core(&radio, ManagerConfig::default().connection_validation(true));
3047 core.add_device_with_priority("A", DevicePriority::High)
3048 .await
3049 .unwrap();
3050 core.add_device_with_priority("B", DevicePriority::Normal)
3051 .await
3052 .unwrap();
3053 core.connect("A").await.unwrap();
3054 core.connect("B").await.unwrap();
3055 radio.make_zombie("A");
3058 radio.set_probe_delay("A", Duration::from_secs(3));
3059 radio.lose_link("B");
3060
3061 assert_eq!(
3062 core.health_tick().await,
3063 TickOutcome {
3064 healthy: 0,
3065 repaired: 2,
3066 failed: 0
3067 }
3068 );
3069 let b = first_disconnect(&radio, "B").expect("B's dead link was not closed");
3070 assert!(
3071 b < Duration::from_secs(1),
3072 "B was closed at {b:?}, after A's check"
3073 );
3074 let a = first_disconnect(&radio, "A").expect("A's zombie link was not closed");
3075 assert!(
3076 a >= Duration::from_secs(3),
3077 "A was closed at {a:?}, before its check ended"
3078 );
3079 let started: Vec<FakeEvent> = radio_log(&radio)
3081 .into_iter()
3082 .filter(|event| matches!(event, FakeEvent::ConnectStarted { .. }))
3083 .collect();
3084 assert_eq!(
3085 started[2..],
3086 [
3087 FakeEvent::ConnectStarted { id: "A".into() },
3088 FakeEvent::ConnectStarted { id: "B".into() }
3089 ]
3090 );
3091
3092 assert_no_orphans(&core, &radio).await;
3093 radio.assert_no_drop_teardown();
3094 })
3095 .await;
3096 }
3097
3098 #[tokio::test(start_paused = true)]
3099 async fn health_monitor_stops_promptly_when_cancelled_mid_reconnect() {
3100 within(LIMIT, async {
3101 let start = Instant::now();
3102 let radio = FakeRadio::new();
3103 let core = core(
3104 &radio,
3105 ManagerConfig::default()
3106 .health_check_interval(Duration::from_secs(5))
3107 .adaptive_interval(false),
3108 );
3109 core.add_device("A").await.unwrap();
3112 radio.set_connect_delay("A", Duration::from_secs(60));
3113
3114 let cancel = CancellationToken::new();
3115 let monitor = Arc::clone(&core).spawn_health_monitor(cancel.clone());
3116 tokio::time::sleep(Duration::from_secs(6)).await;
3117 assert_eq!(radio.connect_count("A"), 1, "the tick started no connect");
3118
3119 cancel.cancel();
3120 monitor.await.unwrap();
3121 let stopped = start.elapsed();
3122 assert!(
3123 stopped < Duration::from_secs(10),
3124 "the monitor stopped at {stopped:?}"
3125 );
3126
3127 tokio::time::sleep(Duration::from_secs(120)).await;
3129 assert!(!radio.link_up("A"));
3130 assert_eq!(
3131 core.available_connections().await,
3132 Some(core.config.max_concurrent_connections)
3133 );
3134
3135 assert_no_orphans(&core, &radio).await;
3136 radio.assert_no_drop_teardown();
3137 })
3138 .await;
3139 }
3140
3141 #[tokio::test(start_paused = true)]
3146 async fn connect_after_a_cancelled_connect_waits_for_the_dead_link_to_go_down() {
3147 within(LIMIT, async {
3148 let radio = FakeRadio::new();
3149 let core = core(&radio, ManagerConfig::default());
3150 core.connect("A").await.unwrap();
3151 radio.set_disconnect_delay("A", Duration::from_secs(2));
3152 radio.lose_link("A");
3153
3154 let first = timeout(Duration::from_millis(500), core.connect("A")).await;
3156 assert!(
3157 first.is_err(),
3158 "connect() returned {first:?} without closing the dead link"
3159 );
3160 core.connect("A").await.expect("second connect");
3161 tokio::time::sleep(Duration::from_secs(3)).await;
3163
3164 assert!(
3165 radio.link_up("A"),
3166 "the old link's disconnect took down the new link: {:#?}",
3167 radio.events()
3168 );
3169 assert_eq!(
3171 connect_starts(&radio, "A"),
3172 [Duration::ZERO, Duration::from_secs(2)]
3173 );
3174
3175 assert_no_orphans(&core, &radio).await;
3176 radio.assert_no_drop_teardown();
3177 })
3178 .await;
3179 }
3180
3181 #[tokio::test(start_paused = true)]
3184 async fn connect_after_a_cancelled_health_tick_waits_for_the_dead_link_to_go_down() {
3185 within(LIMIT, async {
3186 let radio = FakeRadio::new();
3187 let core = core(
3188 &radio,
3189 ManagerConfig::default()
3190 .health_check_interval(Duration::from_secs(5))
3191 .adaptive_interval(false),
3192 );
3193 core.connect("A").await.unwrap();
3194 radio.set_disconnect_delay("A", Duration::from_secs(2));
3195 radio.lose_link("A");
3196
3197 let cancel = CancellationToken::new();
3200 let monitor = Arc::clone(&core).spawn_health_monitor(cancel.clone());
3201 tokio::time::sleep(Duration::from_millis(5500)).await;
3202 cancel.cancel();
3203 monitor.await.unwrap();
3204 core.connect("A").await.expect("connect after the cancel");
3205 assert_eq!(radio.connect_count("A"), 2, "the dead link was kept");
3206 tokio::time::sleep(Duration::from_secs(3)).await;
3208
3209 assert!(
3210 radio.link_up("A"),
3211 "the old link's disconnect took down the new link: {:#?}",
3212 radio.events()
3213 );
3214 assert_eq!(
3216 connect_starts(&radio, "A"),
3217 [Duration::ZERO, Duration::from_secs(7)]
3218 );
3219
3220 assert_no_orphans(&core, &radio).await;
3221 radio.assert_no_drop_teardown();
3222 })
3223 .await;
3224 }
3225
3226 #[tokio::test(start_paused = true)]
3229 async fn health_tick_skips_user_disconnected_device() {
3230 within(LIMIT, async {
3231 let radio = FakeRadio::new();
3232 let core = core(&radio, ManagerConfig::default());
3233 core.connect("A").await.unwrap();
3234 core.disconnect("A").await.unwrap();
3235 let mut events = core.events.subscribe();
3236
3237 assert_eq!(core.health_tick().await, TickOutcome::default());
3238 assert_eq!(radio.connect_count("A"), 1);
3239 assert!(!radio.link_up("A"));
3240 assert!(!core.snapshot("A").await.unwrap().wanted);
3241 assert_eq!(drain(&mut events), Vec::<String>::new());
3242
3243 core.add_device("B").await.unwrap();
3246 assert!(core.snapshot("B").await.unwrap().wanted);
3247 assert!(core.disconnect_all().await.is_empty());
3248 assert_eq!(core.health_tick().await, TickOutcome::default());
3249 assert_eq!(radio.connect_count("B"), 0);
3250
3251 assert_no_orphans(&core, &radio).await;
3252 radio.assert_no_drop_teardown();
3253 })
3254 .await;
3255 }
3256
3257 #[tokio::test(start_paused = true)]
3258 async fn health_tick_skips_evicted_device() {
3259 within(LIMIT, async {
3260 let radio = FakeRadio::new();
3261 let core = core(&radio, ManagerConfig::default().with_max_connections(2));
3262 core.add_device_with_priority("A", DevicePriority::Low)
3263 .await
3264 .unwrap();
3265 core.add_device_with_priority("B", DevicePriority::High)
3266 .await
3267 .unwrap();
3268 core.connect("A").await.unwrap();
3269 core.connect("B").await.unwrap();
3270
3271 assert!(core.evict_lowest_priority().await.unwrap());
3272 assert_eq!(
3273 core.health_tick().await,
3274 TickOutcome {
3275 healthy: 1,
3276 repaired: 0,
3277 failed: 0
3278 }
3279 );
3280 assert_eq!(radio.connect_count("A"), 1);
3281 assert!(!radio.link_up("A"));
3282 assert!(radio.link_up("B"));
3283
3284 assert_no_orphans(&core, &radio).await;
3285 radio.assert_no_drop_teardown();
3286 })
3287 .await;
3288 }
3289
3290 #[tokio::test(start_paused = true)]
3291 async fn health_tick_honours_reconnect_backoff() {
3292 within(LONG_LIMIT, async {
3293 let start = Instant::now();
3294 let radio = FakeRadio::new();
3295 let core = core(&radio, ManagerConfig::default());
3296 core.add_device_with_options(
3297 "A",
3298 ReconnectOptions {
3299 max_attempts: None,
3300 initial_delay: Duration::from_secs(60),
3301 max_delay: Duration::from_secs(600),
3302 backoff_multiplier: 2.0,
3303 use_exponential_backoff: true,
3304 },
3305 )
3306 .await
3307 .unwrap();
3308 core.connect("A").await.unwrap();
3309 radio.lose_link("A");
3310 radio.script_connects("A", [false; 20]);
3311
3312 for _ in 0..800 {
3313 tokio::time::sleep(Duration::from_secs(5)).await;
3314 core.health_tick().await;
3315 }
3316
3317 let starts = connect_starts(&radio, "A");
3321 let gaps: Vec<u64> = starts[1..]
3322 .windows(2)
3323 .map(|pair| (pair[1] - pair[0]).as_secs())
3324 .collect();
3325 assert_eq!(gaps, [60, 120, 240, 480, 600, 600, 600, 600, 600]);
3326 let snapshot = core.snapshot("A").await.unwrap();
3327 assert_eq!(snapshot.failures, 10);
3328 assert_eq!(
3329 snapshot.retry_at,
3330 Some(start + *starts.last().unwrap() + Duration::from_secs(600))
3331 );
3332
3333 assert_no_orphans(&core, &radio).await;
3334 radio.assert_no_drop_teardown();
3335 })
3336 .await;
3337 }
3338
3339 #[tokio::test(start_paused = true)]
3340 async fn health_tick_gives_up_after_max_attempts_until_explicit_connect() {
3341 within(LONG_LIMIT, async {
3342 let radio = FakeRadio::new();
3343 let core = core(&radio, ManagerConfig::default());
3344 core.add_device_with_options(
3345 "A",
3346 ReconnectOptions::default()
3347 .max_attempts(3)
3348 .initial_delay(Duration::from_secs(1)),
3349 )
3350 .await
3351 .unwrap();
3352 core.connect("A").await.unwrap();
3353 radio.lose_link("A");
3354 radio.script_connects("A", [false; 3]);
3355 let mut events = core.events.subscribe();
3356
3357 for _ in 0..720 {
3359 tokio::time::sleep(Duration::from_secs(5)).await;
3360 core.health_tick().await;
3361 }
3362
3363 assert_eq!(
3364 radio.connect_count("A"),
3365 4,
3366 "the first connect and exactly 3 repairs"
3367 );
3368 let errors: Vec<String> = drain(&mut events)
3369 .into_iter()
3370 .filter(|line| line.starts_with("Error"))
3371 .collect();
3372 assert_eq!(errors, ["Error A: auto-reconnect gave up after 3 attempts"]);
3373 let snapshot = core.snapshot("A").await.unwrap();
3374 assert!(snapshot.gave_up);
3375 assert_eq!(snapshot.failures, 3);
3376
3377 core.connect("A").await.unwrap();
3379 assert!(radio.link_up("A"));
3380 let snapshot = core.snapshot("A").await.unwrap();
3381 assert!(!snapshot.gave_up);
3382 assert_eq!(snapshot.failures, 0);
3383 assert!(snapshot.wanted);
3384
3385 assert_no_orphans(&core, &radio).await;
3386 radio.assert_no_drop_teardown();
3387 })
3388 .await;
3389 }
3390
3391 #[tokio::test(start_paused = true)]
3392 async fn explicit_connect_during_a_failing_repair_starts_over() {
3393 within(LIMIT, async {
3394 let radio = FakeRadio::new();
3395 let core = core(&radio, ManagerConfig::default());
3396 core.add_device_with_options("A", ReconnectOptions::default().max_attempts(1))
3397 .await
3398 .unwrap();
3399 core.connect("A").await.unwrap();
3400 radio.lose_link("A");
3401 radio.set_connect_delay("A", Duration::from_secs(10));
3404 radio.script_connects("A", [false, false]);
3405
3406 let tick = tokio::spawn({
3407 let core = Arc::clone(&core);
3408 async move { core.health_tick().await }
3409 });
3410 tokio::time::sleep(Duration::from_secs(1)).await;
3411 let connecting = tokio::spawn({
3415 let core = Arc::clone(&core);
3416 async move { core.connect("A").await }
3417 });
3418 assert_eq!(
3419 tick.await.unwrap(),
3420 TickOutcome {
3421 healthy: 0,
3422 repaired: 0,
3423 failed: 1
3424 }
3425 );
3426 assert!(connecting.await.unwrap().is_err());
3427 assert_eq!(radio.connect_count("A"), 3);
3428
3429 assert_eq!(
3431 core.snapshot("A").await,
3432 Some(EntrySnapshot {
3433 has_link: false,
3434 failures: 0,
3435 wanted: true,
3436 gave_up: false,
3437 retry_at: None,
3438 })
3439 );
3440 assert_eq!(
3442 core.health_tick().await,
3443 TickOutcome {
3444 healthy: 0,
3445 repaired: 1,
3446 failed: 0
3447 }
3448 );
3449 assert!(radio.link_up("A"));
3450
3451 assert_no_orphans(&core, &radio).await;
3452 radio.assert_no_drop_teardown();
3453 })
3454 .await;
3455 }
3456
3457 #[tokio::test(start_paused = true)]
3458 async fn remove_device_during_connect_keeps_it_removed() {
3459 within(LIMIT, async {
3460 let radio = FakeRadio::new();
3461 let core = core(&radio, ManagerConfig::default());
3462 radio.set_connect_delay("A", Duration::from_secs(10));
3463 let connecting = tokio::spawn({
3464 let core = Arc::clone(&core);
3465 async move { core.connect("A").await }
3466 });
3467 tokio::time::sleep(Duration::from_secs(1)).await;
3468
3469 core.remove_device("A").await.unwrap();
3470
3471 let result = connecting.await.unwrap();
3472 assert!(
3473 matches!(result, Err(Error::Cancelled)),
3474 "connect returned {result:?}"
3475 );
3476 assert_eq!(core.device_count().await, 0);
3477 assert!(!radio.link_up("A"));
3478
3479 assert_no_orphans(&core, &radio).await;
3480 radio.assert_no_drop_teardown();
3481 })
3482 .await;
3483 }
3484
3485 #[tokio::test(start_paused = true)]
3486 async fn disconnect_cancels_running_and_waiting_connects() {
3487 within(LIMIT, async {
3488 let radio = FakeRadio::new();
3489 let core = core(&radio, ManagerConfig::default());
3490 radio.set_connect_delay("A", Duration::from_secs(10));
3491 let connects: Vec<_> = (0..2)
3493 .map(|_| {
3494 let core = Arc::clone(&core);
3495 tokio::spawn(async move { core.connect("A").await })
3496 })
3497 .collect();
3498 tokio::time::sleep(Duration::from_secs(1)).await;
3499
3500 core.disconnect("A").await.unwrap();
3501
3502 for connecting in connects {
3503 let result = connecting.await.unwrap();
3504 assert!(
3505 matches!(result, Err(Error::Cancelled)),
3506 "connect returned {result:?}"
3507 );
3508 }
3509 assert_eq!(radio.connect_count("A"), 1);
3511 assert!(!radio.link_up("A"));
3512 assert!(!core.snapshot("A").await.unwrap().wanted);
3513
3514 assert_no_orphans(&core, &radio).await;
3515 radio.assert_no_drop_teardown();
3516 })
3517 .await;
3518 }
3519
3520 #[tokio::test(start_paused = true)]
3526 async fn a_dropped_disconnect_finishes_after_the_check_it_waited_for() {
3527 within(LIMIT, async {
3528 let radio = FakeRadio::new();
3529 let core = core(&radio, ManagerConfig::default().with_max_connections(2));
3530 for id in ["A", "B"] {
3531 core.connect(id).await.unwrap();
3532 radio.set_probe_delay(id, Duration::from_secs(3));
3534 }
3535 let mut events = core.events.subscribe();
3536
3537 let tick = tokio::spawn({
3538 let core = Arc::clone(&core);
3539 async move { core.health_tick().await }
3540 });
3541 tokio::time::sleep(Duration::from_secs(1)).await;
3542 let (disconnect, remove) = tokio::join!(
3544 timeout(Duration::from_secs(1), core.disconnect("A")),
3545 timeout(Duration::from_secs(1), core.remove_device("B")),
3546 );
3547 assert!(disconnect.is_err(), "disconnect() returned {disconnect:?}");
3548 assert!(remove.is_err(), "remove_device() returned {remove:?}");
3549 assert_eq!(
3550 tick.await.unwrap(),
3551 TickOutcome {
3552 healthy: 2,
3553 repaired: 0,
3554 failed: 0
3555 }
3556 );
3557 tokio::time::sleep(Duration::from_secs(1)).await;
3558
3559 assert_eq!(radio.up_ids(), Vec::<String>::new());
3560 assert_eq!(core.device_ids().await, ["A"]);
3561 assert_eq!(
3562 core.snapshot("A").await,
3563 Some(EntrySnapshot {
3564 has_link: false,
3565 failures: 0,
3566 wanted: false,
3567 gave_up: false,
3568 retry_at: None,
3569 })
3570 );
3571 assert_eq!(core.available_connections().await, Some(2));
3572 let mut sent = drain(&mut events);
3573 sent.sort();
3574 assert_eq!(
3575 sent,
3576 [
3577 "Disconnected A UserRequested",
3578 "Disconnected B UserRequested"
3579 ]
3580 );
3581
3582 assert_no_orphans(&core, &radio).await;
3583 radio.assert_no_drop_teardown();
3584 })
3585 .await;
3586 }
3587
3588 #[tokio::test(start_paused = true)]
3593 async fn a_dropped_disconnect_all_disconnects_every_device_it_withdrew() {
3594 within(LIMIT, async {
3595 let radio = FakeRadio::new();
3596 let core = core(&radio, ManagerConfig::default());
3597 core.connect("A").await.unwrap();
3598 let mut events = core.events.subscribe();
3599
3600 let busy = core.devices.read().await;
3604 let disconnect_all = tokio::spawn({
3605 let core = Arc::clone(&core);
3606 async move { timeout(Duration::from_secs(1), core.disconnect_all()).await }
3607 });
3608 tokio::time::sleep(Duration::from_millis(1)).await;
3609 let reader = tokio::spawn({
3610 let core = Arc::clone(&core);
3611 async move {
3612 let _map = core.devices.read().await;
3613 tokio::time::sleep(Duration::from_secs(2)).await;
3614 }
3615 });
3616 tokio::time::sleep(Duration::from_millis(1)).await;
3617 drop(busy);
3618 let given_up = disconnect_all.await.unwrap();
3619 assert!(given_up.is_err(), "disconnect_all() returned {given_up:?}");
3620 assert!(!core.snapshot("A").await.unwrap().wanted);
3621 reader.await.unwrap();
3622 tokio::time::sleep(Duration::from_secs(1)).await;
3623
3624 assert!(
3625 !radio.link_up("A"),
3626 "A is still connected, but withdrawn: {:#?}",
3627 radio.events()
3628 );
3629 assert!(!core.snapshot("A").await.unwrap().has_link);
3630 assert_eq!(drain(&mut events), ["Disconnected A UserRequested"]);
3631
3632 assert_no_orphans(&core, &radio).await;
3633 radio.assert_no_drop_teardown();
3634 })
3635 .await;
3636 }
3637
3638 #[tokio::test(start_paused = true)]
3645 async fn a_connect_spawned_right_after_a_disconnect_keeps_the_device_connected() {
3646 within(LIMIT, async {
3647 let radio = FakeRadio::new();
3648 let core = core(&radio, ManagerConfig::default());
3649 core.connect("A").await.unwrap();
3650 let mut events = core.events.subscribe();
3651
3652 let disconnect = tokio::spawn({
3653 let core = Arc::clone(&core);
3654 async move { core.disconnect("A").await }
3655 });
3656 let connect = tokio::spawn({
3657 let core = Arc::clone(&core);
3658 async move { core.connect("A").await }
3659 });
3660 let disconnect = disconnect.await.unwrap();
3661 let connect = connect.await.unwrap();
3662
3663 assert!(disconnect.is_ok(), "disconnect() returned {disconnect:?}");
3664 assert!(connect.is_ok(), "connect() returned {connect:?}");
3665 assert!(radio.link_up("A"), "{:#?}", radio.events());
3666 assert_eq!(
3667 core.snapshot("A").await,
3668 Some(EntrySnapshot {
3669 has_link: true,
3670 failures: 0,
3671 wanted: true,
3672 gave_up: false,
3673 retry_at: None,
3674 })
3675 );
3676 assert_eq!(drain(&mut events), Vec::<String>::new());
3677 let a = || "A".to_string();
3679 assert_eq!(
3680 radio_log(&radio),
3681 [
3682 FakeEvent::ConnectStarted { id: a() },
3683 FakeEvent::Connected { id: a(), handle: 1 },
3684 ]
3685 );
3686
3687 assert_no_orphans(&core, &radio).await;
3688 radio.assert_no_drop_teardown();
3689 })
3690 .await;
3691 }
3692
3693 #[tokio::test(start_paused = true)]
3698 async fn a_connect_spawned_right_after_disconnect_all_or_an_eviction_keeps_its_device() {
3699 within(LIMIT, async {
3700 let radio = FakeRadio::new();
3701 let core = core(&radio, ManagerConfig::default());
3702 for id in ["A", "B"] {
3703 core.connect(id).await.unwrap();
3704 }
3705 let mut events = core.events.subscribe();
3706 let connect_a = || {
3707 let core = Arc::clone(&core);
3708 tokio::spawn(async move { core.connect("A").await })
3709 };
3710
3711 let disconnect_all = tokio::spawn({
3712 let core = Arc::clone(&core);
3713 async move { core.disconnect_all().await }
3714 });
3715 let connect = connect_a();
3716 let results = disconnect_all.await.unwrap();
3717 let connect = connect.await.unwrap();
3718 assert!(
3719 results.len() == 2 && results.values().all(Result::is_ok),
3720 "disconnect_all() returned {results:?}"
3721 );
3722 assert!(connect.is_ok(), "connect() returned {connect:?}");
3723 assert_eq!(radio.up_ids(), ["A"]);
3724 assert_eq!(drain(&mut events), ["Disconnected B UserRequested"]);
3725
3726 let evict = tokio::spawn({
3728 let core = Arc::clone(&core);
3729 async move { core.evict_lowest_priority().await }
3730 });
3731 let connect = connect_a();
3732 let evicted = evict.await.unwrap();
3733 let connect = connect.await.unwrap();
3734 assert!(
3735 matches!(evicted, Ok(true)),
3736 "evict_lowest_priority() returned {evicted:?}"
3737 );
3738 assert!(connect.is_ok(), "connect() returned {connect:?}");
3739 assert_eq!(radio.up_ids(), ["A"]);
3740 assert_eq!(drain(&mut events), Vec::<String>::new());
3741 assert_eq!(radio.connect_count("A"), 1);
3742 let a = core.snapshot("A").await.unwrap();
3743 assert!(a.wanted && a.has_link, "{a:?}");
3744
3745 assert_no_orphans(&core, &radio).await;
3746 radio.assert_no_drop_teardown();
3747 })
3748 .await;
3749 }
3750
3751 #[tokio::test(start_paused = true)]
3755 async fn a_connect_after_a_disconnect_polled_once_keeps_the_device_connected() {
3756 within(LIMIT, async {
3757 let radio = FakeRadio::new();
3758 let core = core(&radio, ManagerConfig::default());
3759 core.connect("A").await.unwrap();
3760 let mut events = core.events.subscribe();
3761
3762 let mut disconnect = Box::pin(core.disconnect("A"));
3763 let first = disconnect
3764 .as_mut()
3765 .poll(&mut std::task::Context::from_waker(std::task::Waker::noop()));
3766 assert!(first.is_pending(), "disconnect() returned {first:?}");
3767 drop(disconnect);
3768 assert!(!core.snapshot("A").await.unwrap().wanted, "not withdrawn");
3769
3770 core.connect("A").await.expect("connect");
3771 tokio::time::sleep(Duration::from_secs(1)).await;
3773
3774 assert!(radio.link_up("A"), "{:#?}", radio.events());
3775 assert_eq!(
3776 core.snapshot("A").await,
3777 Some(EntrySnapshot {
3778 has_link: true,
3779 failures: 0,
3780 wanted: true,
3781 gave_up: false,
3782 retry_at: None,
3783 })
3784 );
3785 assert_eq!(drain(&mut events), Vec::<String>::new());
3786
3787 assert_no_orphans(&core, &radio).await;
3788 radio.assert_no_drop_teardown();
3789 })
3790 .await;
3791 }
3792
3793 #[tokio::test(start_paused = true)]
3798 async fn a_connect_spawned_right_after_a_removal_does_not_keep_the_device() {
3799 within(LIMIT, async {
3800 let radio = FakeRadio::new();
3801 let core = core(&radio, ManagerConfig::default());
3802 core.connect("A").await.unwrap();
3803 let mut events = core.events.subscribe();
3804
3805 let remove = tokio::spawn({
3806 let core = Arc::clone(&core);
3807 async move { core.remove_device("A").await }
3808 });
3809 let connect = tokio::spawn({
3810 let core = Arc::clone(&core);
3811 async move { core.connect("A").await }
3812 });
3813 let remove = remove.await.unwrap();
3814 let connect = connect.await.unwrap();
3815
3816 assert!(remove.is_ok(), "remove_device() returned {remove:?}");
3817 assert!(
3819 matches!(connect, Ok(()) | Err(Error::Cancelled)),
3820 "connect() returned {connect:?}"
3821 );
3822 assert_eq!(core.device_count().await, 0);
3823 assert!(!radio.link_up("A"), "{:#?}", radio.events());
3824 assert_eq!(drain(&mut events), ["Disconnected A UserRequested"]);
3825
3826 assert_no_orphans(&core, &radio).await;
3827 radio.assert_no_drop_teardown();
3828 })
3829 .await;
3830 }
3831
3832 #[tokio::test(start_paused = true)]
3833 async fn repairs_skip_devices_the_user_changed_during_the_tick() {
3834 within(LIMIT, async {
3835 let radio = FakeRadio::new();
3836 let core = core(&radio, ManagerConfig::default());
3837 core.add_device_with_priority("A", DevicePriority::High)
3838 .await
3839 .unwrap();
3840 for id in ["A", "B", "C", "D"] {
3841 core.connect(id).await.unwrap();
3842 radio.lose_link(id);
3843 }
3844 radio.set_connect_delay("A", Duration::from_secs(10));
3846 let mut events = core.events.subscribe();
3847 let tick = tokio::spawn({
3848 let core = Arc::clone(&core);
3849 async move { core.health_tick().await }
3850 });
3851 tokio::time::sleep(Duration::from_secs(1)).await;
3852
3853 core.connect("B").await.unwrap();
3855 core.disconnect("C").await.unwrap();
3856 core.remove_device("D").await.unwrap();
3857
3858 assert_eq!(
3859 tick.await.unwrap(),
3860 TickOutcome {
3861 healthy: 0,
3862 repaired: 1,
3863 failed: 0
3864 }
3865 );
3866 let repairs: Vec<String> = drain(&mut events)
3867 .into_iter()
3868 .filter(|line| line.starts_with("Reconnect"))
3869 .collect();
3870 assert_eq!(repairs, ["ReconnectStarted A 1", "ReconnectSucceeded A 1"]);
3871 assert_eq!(
3872 ["A", "B", "C", "D"].map(|id| radio.connect_count(id)),
3873 [2, 2, 1, 1]
3874 );
3875 assert_eq!(radio.up_ids(), ["A", "B"]);
3876 assert_eq!(core.device_count().await, 3);
3877 assert!(!core.snapshot("C").await.unwrap().wanted);
3878
3879 assert_no_orphans(&core, &radio).await;
3880 radio.assert_no_drop_teardown();
3881 })
3882 .await;
3883 }
3884
3885 #[tokio::test(start_paused = true)]
3886 async fn add_device_with_options_rejects_invalid_options() {
3887 within(LIMIT, async {
3888 let radio = FakeRadio::new();
3889 let plain = core(&radio, ManagerConfig::default());
3890 let result = plain
3891 .add_device_with_options(
3892 "A",
3893 ReconnectOptions::default().backoff_multiplier(0.5),
3894 )
3895 .await;
3896 assert!(
3897 matches!(result, Err(Error::InvalidConfig(_))),
3898 "add_device_with_options accepted invalid options: {result:?}"
3899 );
3900 assert_eq!(plain.device_count().await, 0);
3901
3902 let strict = core(
3904 &radio,
3905 ManagerConfig {
3906 default_reconnect_options: ReconnectOptions::default()
3907 .backoff_multiplier(0.5),
3908 ..ManagerConfig::default()
3909 },
3910 );
3911 let result = strict
3912 .add_device_with_priority("B", DevicePriority::High)
3913 .await;
3914 assert!(
3915 matches!(result, Err(Error::InvalidConfig(_))),
3916 "add_device_with_priority accepted invalid options: {result:?}"
3917 );
3918 let result = strict.add_device("B").await;
3919 assert!(
3920 matches!(result, Err(Error::InvalidConfig(_))),
3921 "add_device accepted invalid options: {result:?}"
3922 );
3923 assert_eq!(strict.device_count().await, 0);
3924
3925 assert_no_orphans(&plain, &radio).await;
3926 radio.assert_no_drop_teardown();
3927 })
3928 .await;
3929 }
3930
3931 #[tokio::test(start_paused = true)]
3932 async fn repairs_wait_for_a_free_slot_without_counting_failures() {
3933 within(LIMIT, async {
3934 let radio = FakeRadio::new();
3935 let core = core(&radio, ManagerConfig::default().with_max_connections(1));
3936 core.connect("A").await.unwrap();
3937 core.add_device_with_options("B", ReconnectOptions::default())
3938 .await
3939 .unwrap();
3940 let mut events = core.events.subscribe();
3941
3942 for _ in 0..12 {
3945 tokio::time::sleep(Duration::from_secs(5)).await;
3946 assert_eq!(
3947 core.health_tick().await,
3948 TickOutcome {
3949 healthy: 1,
3950 repaired: 0,
3951 failed: 0
3952 }
3953 );
3954 }
3955 assert_eq!(radio.connect_count("B"), 0);
3956 assert_eq!(drain(&mut events), Vec::<String>::new());
3957 let waiting = core.snapshot("B").await.unwrap();
3958 assert_eq!((waiting.failures, waiting.gave_up), (0, false));
3959
3960 core.disconnect("A").await.unwrap();
3962 assert_eq!(
3963 core.health_tick().await,
3964 TickOutcome {
3965 healthy: 0,
3966 repaired: 1,
3967 failed: 0
3968 }
3969 );
3970 assert!(radio.link_up("B"));
3971 assert_eq!(
3972 drain(&mut events),
3973 [
3974 "Disconnected A UserRequested",
3975 "ReconnectStarted B 1",
3976 "Connected B",
3977 "ReconnectSucceeded B 1"
3978 ]
3979 );
3980
3981 assert_no_orphans(&core, &radio).await;
3982 radio.assert_no_drop_teardown();
3983 })
3984 .await;
3985 }
3986
3987 #[tokio::test(start_paused = true)]
3990 async fn successful_repair_starts_the_backoff_over() {
3991 within(LIMIT, async {
3992 let radio = FakeRadio::new();
3993 let core = core(&radio, ManagerConfig::default());
3994 core.add_device_with_options(
3995 "A",
3996 ReconnectOptions::default()
3997 .max_attempts(2)
3998 .initial_delay(Duration::from_secs(60)),
3999 )
4000 .await
4001 .unwrap();
4002 core.connect("A").await.unwrap();
4003 radio.lose_link("A");
4004 radio.script_connects("A", [false]);
4006 assert_eq!(
4007 core.health_tick().await,
4008 TickOutcome {
4009 healthy: 0,
4010 repaired: 0,
4011 failed: 1
4012 }
4013 );
4014 tokio::time::sleep(Duration::from_secs(60)).await;
4015 assert_eq!(
4016 core.health_tick().await,
4017 TickOutcome {
4018 healthy: 0,
4019 repaired: 1,
4020 failed: 0
4021 }
4022 );
4023 assert_eq!(
4024 core.snapshot("A").await,
4025 Some(EntrySnapshot {
4026 has_link: true,
4027 failures: 0,
4028 wanted: true,
4029 gave_up: false,
4030 retry_at: None,
4031 })
4032 );
4033
4034 radio.lose_link("A");
4037 radio.script_connects("A", [false]);
4038 let mut events = core.events.subscribe();
4039 assert_eq!(
4040 core.health_tick().await,
4041 TickOutcome {
4042 healthy: 0,
4043 repaired: 0,
4044 failed: 1
4045 }
4046 );
4047 assert_eq!(
4048 drain(&mut events),
4049 ["Disconnected A Unknown", "ReconnectStarted A 1"]
4050 );
4051 assert!(!core.snapshot("A").await.unwrap().gave_up);
4052
4053 assert_no_orphans(&core, &radio).await;
4054 radio.assert_no_drop_teardown();
4055 })
4056 .await;
4057 }
4058
4059 #[tokio::test(start_paused = true)]
4063 async fn a_repair_cancelled_during_the_info_read_starts_the_backoff_over() {
4064 within(LIMIT, async {
4065 let radio = FakeRadio::new();
4066 let core = core(&radio, ManagerConfig::default());
4067 core.add_device_with_options(
4068 "A",
4069 ReconnectOptions::default()
4070 .max_attempts(2)
4071 .initial_delay(Duration::from_secs(60)),
4072 )
4073 .await
4074 .unwrap();
4075 core.connect("A").await.unwrap();
4076 radio.lose_link("A");
4077 radio.script_connects("A", [false]);
4079 assert_eq!(
4080 core.health_tick().await,
4081 TickOutcome {
4082 healthy: 0,
4083 repaired: 0,
4084 failed: 1
4085 }
4086 );
4087 tokio::time::sleep(Duration::from_secs(60)).await;
4090 radio.set_info_delay("A", Duration::from_secs(5));
4091 let cancelled = timeout(Duration::from_secs(1), core.health_tick()).await;
4092 assert!(cancelled.is_err(), "the tick returned {cancelled:?}");
4093 assert_eq!(
4094 core.snapshot("A").await,
4095 Some(EntrySnapshot {
4096 has_link: true,
4097 failures: 0,
4098 wanted: true,
4099 gave_up: false,
4100 retry_at: None,
4101 })
4102 );
4103
4104 radio.set_info_delay("A", Duration::ZERO);
4107 radio.lose_link("A");
4108 radio.script_connects("A", [false]);
4109 let mut events = core.events.subscribe();
4110 assert_eq!(
4111 core.health_tick().await,
4112 TickOutcome {
4113 healthy: 0,
4114 repaired: 0,
4115 failed: 1
4116 }
4117 );
4118 assert_eq!(
4119 drain(&mut events),
4120 ["Disconnected A Unknown", "ReconnectStarted A 1"]
4121 );
4122 assert!(!core.snapshot("A").await.unwrap().gave_up);
4123
4124 assert_no_orphans(&core, &radio).await;
4125 radio.assert_no_drop_teardown();
4126 })
4127 .await;
4128 }
4129
4130 #[tokio::test(start_paused = true)]
4134 async fn an_overlong_reconnect_wait_is_capped_and_the_monitor_keeps_running() {
4135 within(LIMIT, async {
4136 let start = Instant::now();
4137 let radio = FakeRadio::new();
4138 let core = core(
4139 &radio,
4140 ManagerConfig::default()
4141 .health_check_interval(Duration::from_secs(5))
4142 .adaptive_interval(false),
4143 );
4144 core.add_device_with_options(
4145 "A",
4146 ReconnectOptions::fixed_delay(Duration::MAX).max_delay(Duration::MAX),
4147 )
4148 .await
4149 .unwrap();
4150 core.connect("A").await.unwrap();
4151 core.connect("B").await.unwrap();
4152 radio.lose_link("A");
4153 radio.script_connects("A", [false]);
4154
4155 let cancel = CancellationToken::new();
4156 let monitor = Arc::clone(&core).spawn_health_monitor(cancel.clone());
4157 tokio::time::sleep(Duration::from_secs(6)).await;
4159 let failed_at = connect_starts(&radio, "A")[1];
4160 assert_eq!(failed_at, Duration::from_secs(5));
4161 let snapshot = core.snapshot("A").await.unwrap();
4162 assert_eq!(snapshot.failures, 1);
4163 assert_eq!(
4164 snapshot.retry_at,
4165 Some(start + failed_at + crate::manager::MAX_RETRY_WAIT)
4166 );
4167
4168 radio.lose_link("B");
4170 tokio::time::sleep(Duration::from_secs(5)).await;
4171 assert!(
4172 radio.link_up("B"),
4173 "the monitor stopped: {:#?}",
4174 radio.events()
4175 );
4176 assert!(!monitor.is_finished());
4177 cancel.cancel();
4178 monitor.await.expect("the health monitor panicked");
4179
4180 assert_no_orphans(&core, &radio).await;
4181 radio.assert_no_drop_teardown();
4182 })
4183 .await;
4184 }
4185
4186 #[tokio::test(start_paused = true)]
4193 async fn disconnected_is_sent_even_if_the_close_is_cancelled() {
4194 within(LIMIT, async {
4195 let radio = FakeRadio::new();
4196 let core = core(
4197 &radio,
4198 ManagerConfig::default()
4199 .health_check_interval(Duration::from_secs(5))
4200 .adaptive_interval(false),
4201 );
4202 for id in ["A", "B", "C", "D"] {
4203 core.connect(id).await.unwrap();
4204 radio.set_disconnect_delay(id, Duration::from_secs(2));
4205 }
4206 radio.lose_link("A");
4207 let mut events = core.events.subscribe();
4208
4209 let cancel = CancellationToken::new();
4212 let monitor = Arc::clone(&core).spawn_health_monitor(cancel.clone());
4213 tokio::time::sleep(Duration::from_millis(5500)).await;
4214 cancel.cancel();
4215 monitor.await.unwrap();
4216 radio.lose_link("B");
4218 let given_up = timeout(Duration::from_millis(500), core.connect("B")).await;
4219 assert!(given_up.is_err(), "connect() returned {given_up:?}");
4220 let given_up = timeout(Duration::from_millis(500), core.disconnect("C")).await;
4223 assert!(given_up.is_err(), "disconnect() returned {given_up:?}");
4224 let given_up = timeout(Duration::from_millis(500), core.remove_device("D")).await;
4225 assert!(given_up.is_err(), "remove_device() returned {given_up:?}");
4226
4227 assert_eq!(
4228 drain(&mut events),
4229 [
4230 "Disconnected A Unknown",
4231 "Disconnected B Unknown",
4232 "Disconnected C UserRequested",
4233 "Disconnected D UserRequested"
4234 ]
4235 );
4236 for id in ["A", "B", "C", "D"] {
4237 assert_eq!(core.try_is_connected(id), Some(false), "{id}");
4238 }
4239
4240 tokio::time::sleep(Duration::from_secs(3)).await;
4242 let mut ids = core.device_ids().await;
4243 ids.sort();
4244 assert_eq!(ids, ["A", "B", "C"], "D wasn't removed");
4245 assert_no_orphans(&core, &radio).await;
4246 radio.assert_no_drop_teardown();
4247 })
4248 .await;
4249 }
4250
4251 #[tokio::test(start_paused = true)]
4255 async fn max_attempts_counts_from_zero() {
4256 within(LIMIT, async {
4257 let radio = FakeRadio::new();
4258 let core = core(&radio, ManagerConfig::default());
4259 for (id, max) in [("A", 0), ("B", 1)] {
4260 core.add_device_with_options(id, ReconnectOptions::default().max_attempts(max))
4261 .await
4262 .unwrap();
4263 core.connect(id).await.unwrap();
4264 radio.lose_link(id);
4265 }
4266 radio.script_connects("B", [false]);
4267 let mut events = core.events.subscribe();
4268
4269 for _ in 0..12 {
4271 tokio::time::sleep(Duration::from_secs(5)).await;
4272 core.health_tick().await;
4273 }
4274
4275 assert_eq!(radio.connect_count("A"), 1, "A was reconnected");
4276 assert_eq!(
4277 radio.connect_count("B"),
4278 2,
4279 "the first connect and one repair"
4280 );
4281 let errors: Vec<String> = drain(&mut events)
4282 .into_iter()
4283 .filter(|line| line.starts_with("Error"))
4284 .collect();
4285 assert_eq!(
4286 errors,
4287 [
4288 "Error A: auto-reconnect gave up after 0 attempts",
4289 "Error B: auto-reconnect gave up after 1 attempt"
4290 ]
4291 );
4292 let a = core.snapshot("A").await.unwrap();
4293 assert_eq!((a.failures, a.gave_up), (0, true));
4294 assert!(core.snapshot("B").await.unwrap().gave_up);
4295
4296 assert_no_orphans(&core, &radio).await;
4297 radio.assert_no_drop_teardown();
4298 })
4299 .await;
4300 }
4301
4302 #[tokio::test(start_paused = true)]
4307 async fn max_attempts_of_zero_still_connects_an_added_device_once() {
4308 within(LIMIT, async {
4309 let radio = FakeRadio::new();
4310 let core = core(&radio, ManagerConfig::default());
4311 for id in ["A", "B"] {
4312 core.add_device_with_options(id, ReconnectOptions::default().max_attempts(0))
4313 .await
4314 .unwrap();
4315 }
4316 radio.script_connects("B", [false]);
4317 let mut events = core.events.subscribe();
4318
4319 assert_eq!(
4321 core.health_tick().await,
4322 TickOutcome {
4323 healthy: 0,
4324 repaired: 1,
4325 failed: 1
4326 }
4327 );
4328 assert!(radio.link_up("A"));
4329 assert_eq!(
4330 drain(&mut events),
4331 [
4332 "ReconnectStarted A 1",
4333 "Connected A",
4334 "ReconnectSucceeded A 1",
4335 "ReconnectStarted B 1",
4336 "Error B: auto-reconnect gave up after 1 attempt"
4337 ]
4338 );
4339
4340 radio.lose_link("A");
4343 for _ in 0..12 {
4344 tokio::time::sleep(Duration::from_secs(5)).await;
4345 core.health_tick().await;
4346 }
4347 assert_eq!(radio.connect_count("A"), 1, "A was reconnected");
4348 assert_eq!(radio.connect_count("B"), 1, "B was tried again");
4349 assert_eq!(
4350 drain(&mut events),
4351 [
4352 "Disconnected A Unknown",
4353 "Error A: auto-reconnect gave up after 0 attempts"
4354 ]
4355 );
4356 for id in ["A", "B"] {
4357 let snapshot = core.snapshot(id).await.unwrap();
4358 assert!(snapshot.gave_up && !snapshot.has_link, "{id}: {snapshot:?}");
4359 }
4360
4361 assert_no_orphans(&core, &radio).await;
4362 radio.assert_no_drop_teardown();
4363 })
4364 .await;
4365 }
4366 }
4367}