aranet_core/
manager.rs

1//! Multi-device management.
2//!
3//! This module provides a manager for handling multiple Aranet devices
4//! simultaneously, with connection pooling and concurrent operations.
5
6use 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/// Device priority levels for connection management.
27///
28/// The manager never disconnects a device on its own to make room for
29/// another: when the connection limit is reached, `connect()` fails, and
30/// [`DeviceManager::evict_lowest_priority`] frees a slot when you call it,
31/// unless a `connect()` of the device it picks keeps that device connected.
32/// The health monitor ([`DeviceManager::start_health_monitor`]) reconnects
33/// lost devices highest priority first.
34#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default)]
35pub enum DevicePriority {
36    /// Low priority: evicted first by `evict_lowest_priority`.
37    Low,
38    /// Normal priority (default).
39    #[default]
40    Normal,
41    /// High priority: evicted only when no `Low` or `Normal` device is connected.
42    High,
43    /// Critical priority: never evicted.
44    Critical,
45}
46
47/// Adaptive interval that adjusts based on connection stability.
48///
49/// This is used by the health monitor to check connections more frequently
50/// when connections are unstable, and less frequently when stable.
51#[derive(Debug, Clone)]
52pub struct AdaptiveInterval {
53    /// Base interval when connections are stable.
54    pub base: Duration,
55    /// Current interval (may differ from base based on stability).
56    current: Duration,
57    /// Minimum interval (most frequent checking).
58    pub min: Duration,
59    /// Maximum interval (least frequent checking).
60    pub max: Duration,
61    /// Number of consecutive successes.
62    consecutive_successes: u32,
63    /// Number of consecutive failures.
64    consecutive_failures: u32,
65    /// Success threshold before increasing interval.
66    success_threshold: u32,
67    /// Failure threshold before decreasing interval.
68    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    /// Create a new adaptive interval with custom settings.
88    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    /// Get the current interval.
99    pub fn current(&self) -> Duration {
100        self.current
101    }
102
103    /// Record a successful health check.
104    ///
105    /// After enough consecutive successes, the interval will increase
106    /// (less frequent checks) up to the maximum.
107    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            // Double the interval, capped at max
113            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    /// Record a failed health check (connection lost or reconnect needed).
124    ///
125    /// After enough consecutive failures, the interval will decrease
126    /// (more frequent checks) down to the minimum.
127    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            // Halve the interval, capped at min
133            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    /// Reset to the base interval.
144    pub fn reset(&mut self) {
145        self.current = self.base;
146        self.consecutive_successes = 0;
147        self.consecutive_failures = 0;
148    }
149}
150
151/// Information about a managed device.
152///
153/// [`DeviceManager`] doesn't store `ManagedDevice` values, and none of its
154/// methods return one. The type is kept so that code which names it still
155/// compiles.
156#[derive(Debug)]
157pub struct ManagedDevice {
158    /// Device identifier.
159    pub id: String,
160    /// Device name.
161    pub name: Option<String>,
162    /// Device type.
163    pub device_type: Option<DeviceType>,
164    /// The connected device (if connected).
165    /// Wrapped in Arc to allow concurrent access without holding the manager lock.
166    device: Option<Arc<Device>>,
167    /// Whether auto-reconnect is enabled.
168    pub auto_reconnect: bool,
169    /// Last known reading.
170    pub last_reading: Option<CurrentReading>,
171    /// Device info.
172    pub info: Option<DeviceInfo>,
173    /// Reconnection options (if auto-reconnect is enabled).
174    pub reconnect_options: ReconnectOptions,
175    /// Device priority for connection management.
176    pub priority: DevicePriority,
177    /// Number of consecutive connection failures.
178    pub consecutive_failures: u32,
179    /// Last successful connection timestamp (Unix epoch millis).
180    pub last_success: Option<u64>,
181}
182
183impl ManagedDevice {
184    /// Create a new managed device entry.
185    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    /// Create a managed device with custom reconnect options.
202    pub fn with_reconnect_options(id: &str, options: ReconnectOptions) -> Self {
203        Self {
204            reconnect_options: options,
205            ..Self::new(id)
206        }
207    }
208
209    /// Create a managed device with priority.
210    pub fn with_priority(id: &str, priority: DevicePriority) -> Self {
211        Self {
212            priority,
213            ..Self::new(id)
214        }
215    }
216
217    /// Create a managed device with reconnect options and priority.
218    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    /// Record a successful operation.
227    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    /// Record a failed operation.
238    pub fn record_failure(&mut self) {
239        self.consecutive_failures += 1;
240    }
241
242    /// Check if the device is connected (sync check, doesn't query BLE).
243    pub fn has_device(&self) -> bool {
244        self.device.is_some()
245    }
246
247    /// Check if the device is connected (async, queries BLE).
248    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    /// Get a reference to the underlying device.
257    pub fn device(&self) -> Option<&Arc<Device>> {
258        self.device.as_ref()
259    }
260
261    /// Get a clone of the device Arc.
262    pub fn device_arc(&self) -> Option<Arc<Device>> {
263        self.device.clone()
264    }
265}
266
267/// Configuration for the device manager.
268#[derive(Debug, Clone)]
269pub struct ManagerConfig {
270    /// Default scan options.
271    pub scan_options: ScanOptions,
272    /// Default reconnect options for new devices.
273    ///
274    /// The health monitor ([`DeviceManager::start_health_monitor`]) waits
275    /// between automatic reconnects of a device as these options say, and
276    /// stops after `max_attempts` failures in a row, emitting one
277    /// [`DeviceEvent::Error`]; [`DeviceManager::connect`] starts over. With
278    /// `max_attempts` of 0 it never reconnects a device that has been
279    /// connected, but a device that has never been connected, such as one
280    /// just added, still gets one attempt.
281    ///
282    /// The default is [`ReconnectOptions::unlimited`] since 0.3.0. With the
283    /// five attempts of `ReconnectOptions::default()`, the monitor would give
284    /// up on a device after about a minute.
285    pub default_reconnect_options: ReconnectOptions,
286    /// Event channel capacity.
287    pub event_capacity: usize,
288    /// Health check interval for auto-reconnect (base interval).
289    pub health_check_interval: Duration,
290    /// Maximum number of concurrent BLE connections.
291    ///
292    /// Most BLE adapters support 5-7 concurrent connections.
293    /// Attempting to connect beyond this limit will return an error.
294    /// Set to 0 for no limit (not recommended).
295    pub max_concurrent_connections: usize,
296    /// Whether to use adaptive health check intervals.
297    ///
298    /// When enabled, the health check interval will automatically adjust:
299    /// - Decrease (more frequent) when connections are unstable
300    /// - Increase (less frequent) when connections are stable
301    pub use_adaptive_interval: bool,
302    /// Minimum health check interval (for adaptive mode).
303    pub min_health_check_interval: Duration,
304    /// Maximum health check interval (for adaptive mode).
305    pub max_health_check_interval: Duration,
306    /// Default priority for new devices.
307    pub default_priority: DevicePriority,
308    /// Whether to use connection validation (keepalive checks).
309    ///
310    /// When enabled, health checks use `device.validate_connection()`, which
311    /// reads the current measurements to verify the connection is alive. That
312    /// read needs no pairing that reading the sensor doesn't. This catches
313    /// "zombie connections" but uses more power.
314    ///
315    /// When disabled, health checks only ask the Bluetooth stack whether the
316    /// device is connected. A zombie connection passes that check, so the
317    /// health monitor replaces it only once the stack reports it lost.
318    pub use_connection_validation: bool,
319}
320
321impl Default for ManagerConfig {
322    fn default() -> Self {
323        // Use platform-specific defaults if available
324        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    /// Create a configuration with a specific connection limit.
343    pub fn with_max_connections(mut self, max: usize) -> Self {
344        self.max_concurrent_connections = max;
345        self
346    }
347
348    /// Create a configuration with no connection limit (not recommended).
349    pub fn unlimited_connections(mut self) -> Self {
350        self.max_concurrent_connections = 0;
351        self
352    }
353
354    /// Enable or disable adaptive health check intervals.
355    pub fn adaptive_interval(mut self, enabled: bool) -> Self {
356        self.use_adaptive_interval = enabled;
357        self
358    }
359
360    /// Set the health check interval (base interval for adaptive mode).
361    pub fn health_check_interval(mut self, interval: Duration) -> Self {
362        self.health_check_interval = interval;
363        self
364    }
365
366    /// Set the default device priority.
367    pub fn default_priority(mut self, priority: DevicePriority) -> Self {
368        self.default_priority = priority;
369        self
370    }
371
372    /// Enable or disable connection validation in health checks.
373    pub fn connection_validation(mut self, enabled: bool) -> Self {
374        self.use_connection_validation = enabled;
375        self
376    }
377}
378
379/// Longest wait the health monitor schedules before a device's next automatic
380/// reconnect. It only caps `ReconnectOptions::max_delay` values so large that
381/// adding them to an `Instant` could overflow.
382const MAX_RETRY_WAIT: Duration = Duration::from_secs(365 * 24 * 60 * 60);
383
384/// A managed device's state inside `ManagerCore`.
385struct 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    /// Consecutive failed health-monitor reconnects.
394    failures: u32,
395    /// Whether the user wants the device connected. New entries start wanted;
396    /// `connect` sets it; `disconnect`, `disconnect_all`,
397    /// `evict_lowest_priority` and `remove_device` clear it. `add_device*` on
398    /// an existing entry leave it as it is, so only `connect` re-arms a
399    /// withdrawn device. The health monitor checks and repairs only wanted
400    /// devices.
401    wanted: bool,
402    /// When the health monitor may next try to reconnect the device (`None`:
403    /// at once).
404    retry_at: Option<Instant>,
405    /// Set when automatic reconnects used up `max_attempts`; cleared by `connect`.
406    gave_up: bool,
407    /// Whether a link has ever been stored in this entry. Until then, the
408    /// health monitor's connect is the device's first connect, not a
409    /// reconnect, and `max_attempts` of 0 doesn't stop it.
410    ever_connected: bool,
411    /// The connection, if connected.
412    link: Option<Arc<L>>,
413    /// The connection-limit permit, taken before connecting and held until
414    /// the link has been disconnected.
415    slot: Option<OwnedSemaphorePermit>,
416    /// Serialises connect, disconnect and removal of this device. Taken
417    /// without holding the device map lock.
418    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    /// Whether the health monitor should reconnect this device now.
443    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    /// Gives up on automatic reconnects if they have used up `max_attempts`,
452    /// and then returns the number of failures.
453    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    /// Before an automatic reconnect: gives up without one if `max_attempts`
466    /// is already used up, and then returns the number of failures. Only
467    /// `max_attempts` of 0 can be used up here, as a failure that uses up a
468    /// larger one gives up at once. A device that has never been connected
469    /// isn't reconnected but connected for the first time, so it still gets
470    /// that attempt.
471    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    /// Records a failed automatic reconnect and schedules the next one.
480    /// Returns the number of failures when this one used up `max_attempts`.
481    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    /// Starts automatic reconnects over, for an explicit `connect` or once a
492    /// new link is stored: the next repair is due at once, with all of
493    /// `max_attempts` available again.
494    fn reset_backoff(&mut self) {
495        self.failures = 0;
496        self.retry_at = None;
497        self.gave_up = false;
498    }
499}
500
501/// What the lifecycle tests can see of an entry.
502#[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
512/// A new link and its slot, not yet stored in the link's entry, with a share
513/// of its connect's `op` guard. If it is dropped before `into_parts` (its
514/// connect was cancelled while waiting for the device map), it disconnects
515/// the link on a spawned task, which keeps the slot and the share of the
516/// guard until the link is down.
517struct NewLink<L: SensorLink> {
518    /// The device the link connects, for the log.
519    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    /// Hands over the link and its slot; dropping `self` then only gives back
539    /// its share of the guard, which the connect still holds.
540    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        // Without a runtime, dropping the link runs its own teardown.
553        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
565/// Awaits a disconnect or removal task. A task that panicked, or that was
566/// dropped when its runtime shut down, gives an `Error::Io`.
567async fn finished(task: tokio::task::JoinHandle<Result<()>>) -> Result<()> {
568    task.await.map_err(std::io::Error::from)?
569}
570
571/// What one health-monitor tick did.
572#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
573pub(crate) struct TickOutcome {
574    /// Connected devices that passed their check.
575    pub(crate) healthy: usize,
576    /// Devices this tick reconnected.
577    pub(crate) repaired: usize,
578    /// Reconnects that failed.
579    pub(crate) failed: usize,
580}
581
582/// The device manager's logic, generic over the connection type so that the
583/// lifecycle tests can run it on `test_support::FakeRadio`. `DeviceManager` is
584/// this over `Device`.
585pub(crate) struct ManagerCore<L: SensorLink> {
586    devices: RwLock<HashMap<String, Entry<L>>>,
587    events: EventDispatcher,
588    config: ManagerConfig,
589    connect: ConnectFn<L>,
590    /// One permit per allowed connection; `None` when connections are not limited.
591    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    /// The error for a connect that the connection limit rejects.
618    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    /// Takes a connection slot for `identifier`: `None` when connections are
633    /// not limited, the limit error when no slot is free.
634    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(()); // Already exists
659        }
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            // A device that the manager doesn't know yet takes its slot here,
674            // under the same lock that adds it, so a connect that the limit
675            // rejects doesn't add the device.
676            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            // An explicit connect: keep the device connected from now on.
689            entry.wanted = true;
690            (Arc::clone(&entry.op), reserved)
691        };
692        // The map lock is released before waiting. One connect, disconnect or
693        // removal of this device runs at a time. A cancelled caller releases
694        // the guard, and a reserved slot, with its future, except that a new
695        // link it abandons keeps a share of the guard until it is down.
696        let held = Arc::new(Arc::clone(&op).lock_owned().await);
697        // An explicit connect starts automatic reconnects over. This runs under
698        // the guard, so a repair that failed while this connect waited can't
699        // leave its failure count or give-up behind.
700        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    /// Connects `identifier` unless it already has a link that the Bluetooth
709    /// stack reports as up; a link that is down is closed and replaced. The
710    /// caller holds the entry's `op` guard and passes it, shared, as `held`;
711    /// a lost link being closed (see `detach`) and a new link that is
712    /// abandoned each keep a share of it until that link is down.
713    /// `reserved` is the slot that `connect` took when it added the device;
714    /// without one, a slot is taken here.
715    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                // Disconnected or removed while this connect waited for the guard.
727                _ => 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        // Stored with the link and held until the link has been disconnected.
741        // Dropped here if the connect fails or its caller is cancelled.
742        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        // Store the link as soon as the map lock is free, so that a cancel from
755        // then on leaves a link that the manager holds. A cancel while waiting
756        // for the lock drops `new_link`, which disconnects it.
757        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                    // The device is connected from here on, even if this
766                    // connect is cancelled before it returns.
767                    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                // Whoever removed the entry released its link.
797                _ => return Err(Error::Cancelled),
798            }
799        }
800
801        // Emit event
802        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    /// Marks the device as not wanted, so that a connect or repair of it
816    /// that is running closes its new link instead of installing it, and the
817    /// health monitor leaves the device alone until `connect`. Returns the
818    /// entry's `op` mutex, or `None` if the device isn't managed.
819    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    /// Disconnects a withdrawn device (see `withdraw`) on a task, once its
834    /// `op` mutex is free. The health monitor leaves a withdrawn device alone,
835    /// so nothing else would close its link: the task finishes even if the
836    /// caller is dropped, for example while it waits for a connect or a
837    /// check of the device to end. A `connect()` that re-arms the device
838    /// before the task takes the link supersedes the disconnect, and the task
839    /// then does nothing (see `disconnect_locked`).
840    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    /// Disconnects `identifier`'s link, if it has one, and returns the
854    /// close's result. It fails only after taking the link out of the entry.
855    /// With `unless_rearmed`, it does nothing if the entry is wanted again: a
856    /// `connect()` has re-armed the device since it was withdrawn, so that
857    /// connect started before this took the link, and supersedes it. The
858    /// check and the take share one lock of the device map, so no connect
859    /// can come between them.
860    /// `held` is the caller's guard of the entry's `op` mutex. The disconnect
861    /// keeps a share of that guard, and the slot, until the link is down,
862    /// even if this future is dropped: until then no other connect,
863    /// disconnect or removal of the device starts, so none can overlap the
864    /// old link's disconnect.
865    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        // The link has left the manager: say so before the close, as `detach`
886        // does. The close can fail (on macOS, closing a link that dropped on
887        // its own times out) or go on in the background if this future is
888        // dropped, and nothing would send the event later.
889        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        // On a task, as in `spawn_disconnect`, so that a dropped caller still
901        // removes the device it has withdrawn.
902        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            // Not `disconnect`: the guard is not reentrant. A removal always
907            // goes ahead, even if a `connect()` has re-armed the device.
908            let result = core.disconnect_locked(&identifier, &held, false).await;
909            // Removed even if closing the link failed: the link has left the
910            // manager, so a kept entry would have nothing left to close.
911            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        // Collect links while holding the lock briefly
959        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        // Lock is released here
967
968        // Check connection status in parallel
969        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        // Get the link while holding the lock briefly
976        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        // Lock is released here
984
985        let reading = link.read_current().await?;
986
987        // Emit reading event
988        self.events.send(DeviceEvent::Reading {
989            device: DeviceId::new(identifier),
990            reading,
991        });
992
993        // Update cached reading
994        {
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        // Collect links while holding the lock briefly
1006        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        // Lock is released here
1014
1015        // Perform all reads in parallel
1016        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        // Emit events and update cache
1024        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        // Update cached readings
1034        {
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        // Note: We can't fully parallelize connect because it modifies state,
1052        // but we can at least attempt connections concurrently
1053        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            // Withdraw every device, including those without a link (not
1065            // connected yet, or lost and waiting for a repair), so the health
1066            // monitor connects none of them afterwards.
1067            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        // Every disconnect starts before anything else is awaited, so that a
1077        // dropped caller can't leave a device it has withdrawn connected.
1078        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        // Try to acquire the lock without blocking
1090        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, // Lock was held, couldn't check
1097        }
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    /// One pass of the health monitor.
1124    ///
1125    /// Checks every connected device at once and closes the links that fail,
1126    /// then reconnects the devices that are due, one at a time and highest
1127    /// priority first. A device whose connect, disconnect or removal is
1128    /// running is left alone.
1129    pub(crate) async fn health_tick(&self) -> TickOutcome {
1130        let mut outcome = TickOutcome::default();
1131
1132        // First pass: check every connected device at once, so a slow device
1133        // doesn't hold up the others.
1134        let targets: Vec<_> = {
1135            let devices = self.devices.read().await;
1136            devices
1137                .iter()
1138                // A device being disconnected or removed belongs to that call.
1139                .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        // Second pass: reconnect one device at a time (a second search would
1157        // only wait for the scan permit), highest priority first.
1158        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            // A connect, disconnect or remove of this device is running.
1169            let Ok(guard) = Arc::clone(&op).try_lock_owned() else {
1170                continue;
1171            };
1172            let held = Arc::new(guard);
1173            // Check again now that the device is ours: it may have been
1174            // connected, disconnected or removed since the list was made.
1175            // With `max_attempts` of 0 a device that has been connected gets
1176            // no reconnect: give up.
1177            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            // With every connection slot taken, the repair would fail at once,
1194            // without a Bluetooth attempt: wait for a free slot instead, and
1195            // don't count a failure.
1196            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                // `connect_locked` started the backoff over when it stored
1211                // the new link.
1212                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                // The device was disconnected or removed while this repair ran.
1221                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    /// Logs that the health monitor stops reconnecting `id` after `failures`
1241    /// failed attempts, and emits the one `DeviceEvent::Error` that says so.
1242    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    /// Checks one connected device for `health_tick` and closes its link if
1255    /// the check fails. Returns whether the device was checked and is healthy.
1256    async fn probe(&self, id: String, link: Arc<L>, op: Arc<Mutex<()>>) -> bool {
1257        // A connect, disconnect or remove of this device is running.
1258        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    /// Closes `link` if the entry for `identifier` still holds it: takes the
1276    /// link and its slot, emits `Disconnected`, disconnects the link and frees
1277    /// the slot once the link is down. `held` is the caller's guard of the
1278    /// entry's `op` mutex. As in `disconnect_locked`, the disconnect keeps a
1279    /// share of that guard until the link is down, even if this future is
1280    /// dropped, so no connect of the device can overlap the old link's
1281    /// disconnect.
1282    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                    // The first repair is due at once.
1299                    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        // The link has left the manager: say so before the close, which can
1309        // take seconds and goes on in the background if this future is
1310        // dropped, so a dropped caller can't lose the event.
1311        self.events.send(DeviceEvent::Disconnected {
1312            device: DeviceId::new(identifier),
1313            reason,
1314        });
1315        // Close it explicitly: a handle dropped without a disconnect tears
1316        // down whatever link the sensor has by then, including a new one.
1317        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                // Cancelling mid-tick drops the tick: a connect in progress
1346                // gives back its slot and `op` guard, and releases the sensor;
1347                // a dead link being closed, or a new link not stored yet,
1348                // keeps both until it is down.
1349                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            // Update priority if device already exists
1374            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            // Create passive monitor
1418            let passive_monitor = Arc::new(PassiveMonitor::new(options));
1419            let mut passive_rx = passive_monitor.subscribe();
1420
1421            // Start passive monitoring
1422            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                                // Convert passive reading to CurrentReading and emit event
1435                                if let Some(reading) = passive_reading_to_current(&passive_reading) {
1436                                    // Update last reading in the entry if it exists
1437                                    if let Some(entry) = self.devices.write().await.get_mut(&passive_reading.device_id) {
1438                                        entry.last_reading = Some(reading);
1439                                    }
1440
1441                                    // Emit reading event
1442                                    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        // Check if we have a recent cached reading
1474        {
1475            let devices = self.devices.read().await;
1476            if let Some(entry) = devices.get(identifier)
1477                && let Some(reading) = entry.last_reading
1478            {
1479                // Check if the reading has a captured_at timestamp
1480                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        // No recent passive reading, use active connection
1493        debug!(
1494            "No recent passive reading, using active connection for {}",
1495            identifier
1496        );
1497        self.read_current(identifier).await
1498    }
1499
1500    /// A copy of the entry's lifecycle state, for the lifecycle tests.
1501    #[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
1514/// Manager for multiple Aranet devices.
1515pub struct DeviceManager {
1516    core: Arc<ManagerCore<Device>>,
1517}
1518
1519impl DeviceManager {
1520    /// Create a new device manager.
1521    pub fn new() -> Self {
1522        Self::with_config(ManagerConfig::default())
1523    }
1524
1525    /// Create a manager with custom event capacity.
1526    pub fn with_event_capacity(capacity: usize) -> Self {
1527        Self::with_config(ManagerConfig {
1528            event_capacity: capacity,
1529            ..Default::default()
1530        })
1531    }
1532
1533    /// Create a manager with full configuration.
1534    pub fn with_config(config: ManagerConfig) -> Self {
1535        Self {
1536            core: Arc::new(ManagerCore::new(config, ble_connector())),
1537        }
1538    }
1539
1540    /// Get the event dispatcher for subscribing to events.
1541    pub fn events(&self) -> &EventDispatcher {
1542        &self.core.events
1543    }
1544
1545    /// Get the manager configuration.
1546    pub fn config(&self) -> &ManagerConfig {
1547        &self.core.config
1548    }
1549
1550    /// Scan for available devices.
1551    pub async fn scan(&self) -> Result<Vec<DiscoveredDevice>> {
1552        scan_with_options(self.core.config.scan_options.clone()).await
1553    }
1554
1555    /// Scan with custom options.
1556    pub async fn scan_with_options(&self, options: ScanOptions) -> Result<Vec<DiscoveredDevice>> {
1557        let devices = scan_with_options(options).await?;
1558
1559        // Emit discovery events
1560        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    /// Add a device to the manager by identifier.
1575    ///
1576    /// `identifier` must match the device exactly, as [`Device::connect`]
1577    /// describes.
1578    ///
1579    /// # Errors
1580    ///
1581    /// Returns [`Error::InvalidConfig`] if the config's
1582    /// `default_reconnect_options` are invalid (see [`ReconnectOptions::validate`]).
1583    pub async fn add_device(&self, identifier: &str) -> Result<()> {
1584        self.core.add_device(identifier).await
1585    }
1586
1587    /// Add a device with custom reconnect options.
1588    ///
1589    /// `identifier` must match the device exactly, as [`Device::connect`]
1590    /// describes.
1591    ///
1592    /// The health monitor waits between automatic reconnects of the device as
1593    /// `reconnect_options` say. If the device is already managed, nothing
1594    /// changes.
1595    ///
1596    /// # Errors
1597    ///
1598    /// Returns [`Error::InvalidConfig`] if `reconnect_options` are invalid (see
1599    /// [`ReconnectOptions::validate`]).
1600    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    /// Connect to a device.
1611    ///
1612    /// `identifier` must match the device exactly, as [`Device::connect`]
1613    /// describes.
1614    ///
1615    /// This method performs an atomic connect-or-skip operation:
1616    /// - If the device doesn't exist, it's added and connected
1617    /// - If the device exists but is not connected, it's connected
1618    /// - If the device already has a connection that the Bluetooth stack
1619    ///   reports as up, this is a no-op; a lost one is closed and replaced
1620    ///
1621    /// A second connect of the same device waits for the first and returns its
1622    /// own real result: `Ok(())` if the first one connected the device,
1623    /// otherwise the result of its own attempt. A connect that hasn't
1624    /// connected yet when [`disconnect`](Self::disconnect),
1625    /// [`disconnect_all`](Self::disconnect_all),
1626    /// [`evict_lowest_priority`](Self::evict_lowest_priority) or
1627    /// [`remove_device`](Self::remove_device) is called for the same device
1628    /// returns [`Error::Cancelled`], closing its connection if one comes up,
1629    /// unless another `connect()` re-arms the device first.
1630    /// A connect that waits for a [`remove_device`](Self::remove_device) of
1631    /// the same device returns [`Error::Cancelled`] instead of adding the
1632    /// device back.
1633    ///
1634    /// If this future is dropped while a lost connection is being closed,
1635    /// the close still finishes in the background, and a connect of the same
1636    /// device waits for it instead of running alongside it.
1637    ///
1638    /// If this future is dropped after the connection is made but before the
1639    /// device information has been read, the device stays connected, but no
1640    /// [`DeviceEvent::Connected`] is sent and
1641    /// [`get_device_info`](Self::get_device_info) returns `None` until the
1642    /// device is reconnected.
1643    ///
1644    /// # Connection Limits
1645    ///
1646    /// If `max_concurrent_connections` is set in the config and would be exceeded,
1647    /// this method returns an error. The limit counts connects in progress; use
1648    /// `can_connect()` or `available_connections()` to check it before calling
1649    /// this method.
1650    ///
1651    /// The device map is locked only while the entry is updated, not during the
1652    /// BLE connection, so operations on other devices don't wait for it;
1653    /// operations on the same device do, as described above.
1654    pub async fn connect(&self, identifier: &str) -> Result<()> {
1655        self.core.connect(identifier).await
1656    }
1657
1658    /// Disconnect from a device.
1659    ///
1660    /// A connect of the same device that hasn't connected yet is abandoned:
1661    /// it returns [`Error::Cancelled`], closing its connection if one comes
1662    /// up, unless another `connect()` re-arms the device first. This waits
1663    /// until that connect's Bluetooth attempt has ended and such a
1664    /// connection is closed, which can take the connect's whole time budget:
1665    /// tens of seconds at default settings. The health monitor doesn't
1666    /// reconnect the device until [`connect`](Self::connect) is called.
1667    ///
1668    /// [`DeviceEvent::Disconnected`], with [`DisconnectReason::UserRequested`],
1669    /// is sent as soon as the manager lets go of the connection, before the
1670    /// connection is closed, so it is sent even if closing the connection
1671    /// fails or this future is dropped. A failed close's error is still
1672    /// returned.
1673    ///
1674    /// If this future is dropped, a disconnect that has started still
1675    /// finishes in the background, even one still waiting for a connect or a
1676    /// health check of the device to end.
1677    ///
1678    /// Whether or not this future is dropped, a connect of the same device
1679    /// that starts while the disconnect runs either waits for it and then
1680    /// reconnects, or, if it starts before the disconnect has taken the
1681    /// connection, keeps the device connected, and the disconnect then does
1682    /// nothing.
1683    pub async fn disconnect(&self, identifier: &str) -> Result<()> {
1684        self.core.disconnect(identifier).await
1685    }
1686
1687    /// Remove a device from the manager.
1688    ///
1689    /// The device is disconnected first, as [`disconnect`](Self::disconnect)
1690    /// does, and then removed, even if closing its connection fails: that
1691    /// error is still returned. A connect of the device that hasn't
1692    /// connected yet is abandoned: it returns [`Error::Cancelled`], closing
1693    /// its connection if one comes up, unless another `connect()` re-arms the
1694    /// device first, in which case it connects, and the device is removed
1695    /// after that. As with `disconnect`, this waits until that connect's
1696    /// Bluetooth attempt has ended and such a connection is closed, and a
1697    /// removal that has started still finishes in the background if this
1698    /// future is dropped.
1699    ///
1700    /// Whether or not this future is dropped, a connect of the same device
1701    /// that starts while the removal runs doesn't keep it: the device is
1702    /// removed even when that connect returns first.
1703    pub async fn remove_device(&self, identifier: &str) -> Result<()> {
1704        self.core.remove_device(identifier).await
1705    }
1706
1707    /// Get a list of all managed device IDs.
1708    pub async fn device_ids(&self) -> Vec<String> {
1709        self.core.device_ids().await
1710    }
1711
1712    /// Get the number of managed devices.
1713    pub async fn device_count(&self) -> usize {
1714        self.core.device_count().await
1715    }
1716
1717    /// Get the number of connected devices (fast, doesn't query BLE).
1718    ///
1719    /// This returns the number of devices that have an active device handle,
1720    /// without querying the BLE stack. Use `connected_count_verified` for
1721    /// an accurate count that queries each device.
1722    pub async fn connected_count(&self) -> usize {
1723        self.core.connected_count().await
1724    }
1725
1726    /// Check if a new connection can be made without exceeding the limit.
1727    ///
1728    /// Returns `true` if another connection can be made, `false` if at limit.
1729    /// Connects in progress count against the limit. Always returns `true` if
1730    /// `max_concurrent_connections` is 0 (unlimited).
1731    pub async fn can_connect(&self) -> bool {
1732        self.core.can_connect().await
1733    }
1734
1735    /// Get the connection limit status.
1736    ///
1737    /// Returns (current_connections, max_connections). If max is 0, there is no limit.
1738    ///
1739    /// `current_connections` counts devices with a connection handle; connects in
1740    /// progress are not included but already hold a slot, so use
1741    /// `available_connections` or `can_connect` to see whether `connect()` will
1742    /// pass the limit.
1743    pub async fn connection_status(&self) -> (usize, usize) {
1744        self.core.connection_status().await
1745    }
1746
1747    /// Get the number of available connection slots.
1748    ///
1749    /// Connects in progress hold a slot, and so does a disconnect until the link
1750    /// is down. Returns `None` if there is no connection limit (unlimited).
1751    pub async fn available_connections(&self) -> Option<usize> {
1752        self.core.available_connections().await
1753    }
1754
1755    /// Get the number of connected devices (verified via BLE).
1756    ///
1757    /// This method queries each device to verify its connection status.
1758    /// The lock is released before making BLE calls to avoid contention.
1759    pub async fn connected_count_verified(&self) -> usize {
1760        self.core.connected_count_verified().await
1761    }
1762
1763    /// Read current values from a specific device.
1764    pub async fn read_current(&self, identifier: &str) -> Result<CurrentReading> {
1765        self.core.read_current(identifier).await
1766    }
1767
1768    /// Read current values from all connected devices (in parallel).
1769    ///
1770    /// This method releases the lock before performing async BLE operations,
1771    /// allowing other tasks to add/remove devices while reads are in progress.
1772    /// All reads are performed in parallel for maximum performance.
1773    pub async fn read_all(&self) -> HashMap<String, Result<CurrentReading>> {
1774        self.core.read_all().await
1775    }
1776
1777    /// Connect to all known devices (in parallel).
1778    ///
1779    /// Returns a map of device IDs to connection results.
1780    pub async fn connect_all(&self) -> HashMap<String, Result<()>> {
1781        self.core.connect_all().await
1782    }
1783
1784    /// Disconnect from all devices (in parallel).
1785    ///
1786    /// Returns a map of device IDs to disconnection results, with an entry for
1787    /// each device that had a connection. Every managed device is withdrawn,
1788    /// connected or not: the health monitor reconnects none of them until
1789    /// [`connect`](Self::connect) is called. Each connected device is
1790    /// disconnected as [`disconnect`](Self::disconnect) describes.
1791    pub async fn disconnect_all(&self) -> HashMap<String, Result<()>> {
1792        self.core.disconnect_all().await
1793    }
1794
1795    /// Check if a specific device is connected (fast, doesn't query BLE).
1796    ///
1797    /// This method attempts to check if a device has an active connection handle
1798    /// without blocking. Returns `None` if the lock couldn't be acquired immediately,
1799    /// or `Some(bool)` indicating whether the device has a connection handle.
1800    ///
1801    /// Note: This only checks if we have a device handle, not whether the actual
1802    /// BLE connection is still alive. Use [`is_connected`](Self::is_connected) for
1803    /// a verified check.
1804    pub fn try_is_connected(&self, identifier: &str) -> Option<bool> {
1805        self.core.try_is_connected(identifier)
1806    }
1807
1808    /// Check if a specific device is connected (verified via BLE).
1809    ///
1810    /// The lock is released before making the BLE call.
1811    pub async fn is_connected(&self, identifier: &str) -> bool {
1812        self.core.is_connected(identifier).await
1813    }
1814
1815    /// Get device info for a specific device.
1816    pub async fn get_device_info(&self, identifier: &str) -> Option<DeviceInfo> {
1817        self.core.get_device_info(identifier).await
1818    }
1819
1820    /// Get the last cached reading for a device.
1821    pub async fn get_last_reading(&self, identifier: &str) -> Option<CurrentReading> {
1822        self.core.get_last_reading(identifier).await
1823    }
1824
1825    /// Start a background task that checks the managed devices and repairs
1826    /// lost connections.
1827    ///
1828    /// On every tick the task:
1829    ///
1830    /// 1. Checks every connected device at the same time, so a slow device
1831    ///    doesn't delay the others. A connection that fails its check is
1832    ///    disconnected explicitly, and [`DeviceEvent::Disconnected`] is emitted.
1833    /// 2. Reconnects the devices that should be connected but aren't, one at a
1834    ///    time and highest [`DevicePriority`] first, emitting
1835    ///    [`DeviceEvent::ReconnectStarted`] and
1836    ///    [`DeviceEvent::ReconnectSucceeded`].
1837    ///
1838    /// Each device waits between reconnect attempts as its [`ReconnectOptions`]
1839    /// say. After `max_attempts` failures in a row the task stops trying and
1840    /// emits one [`DeviceEvent::Error`]; [`connect`](Self::connect) starts
1841    /// over. With `max_attempts` of 0 the task never reconnects a device that
1842    /// has been connected: it gives up as soon as it finds the device's
1843    /// connection gone, without trying. A device that has never been
1844    /// connected, such as one just added, still gets one attempt.
1845    /// While `max_concurrent_connections` connections are in use, a device
1846    /// waits for a free one, and that wait isn't counted as a failed attempt.
1847    ///
1848    /// Devices added with [`add_device`](Self::add_device) and its variants
1849    /// are connected by the task. Devices that were disconnected (with
1850    /// [`disconnect`](Self::disconnect), [`disconnect_all`](Self::disconnect_all)
1851    /// or [`evict_lowest_priority`](Self::evict_lowest_priority)) or removed
1852    /// are not reconnected until `connect()` is called for them.
1853    ///
1854    /// The task runs until the provided cancellation token is cancelled, and
1855    /// then stops at once, even in the middle of a check or a reconnect. A
1856    /// connect it abandons releases the sensor, except one abandoned after the
1857    /// connection is made but before the device information has been read:
1858    /// that one keeps its connection, and no
1859    /// [`DeviceEvent::ReconnectSucceeded`] follows its
1860    /// [`DeviceEvent::ReconnectStarted`]. A lost connection it was closing is
1861    /// still closed in the background, and a [`connect`](Self::connect) of
1862    /// that device waits for it.
1863    ///
1864    /// # Adaptive Intervals
1865    ///
1866    /// If `use_adaptive_interval` is enabled in the config, the time between
1867    /// ticks adapts to connection stability:
1868    /// - after a tick with a failed reconnect, it halves (down to
1869    ///   `min_health_check_interval`);
1870    /// - after three ticks with a healthy device and no reconnect, without a
1871    ///   failed reconnect in between, it doubles (up to
1872    ///   `max_health_check_interval`).
1873    ///
1874    /// # Connection Validation
1875    ///
1876    /// If `use_connection_validation` is enabled, health checks read the current
1877    /// measurements (`device.validate_connection()`, which needs no pairing that
1878    /// a reading doesn't) to catch "zombie connections" where the BLE stack
1879    /// thinks it's connected but the device is out of range. Otherwise they only
1880    /// ask the BLE stack, which misses zombie connections.
1881    ///
1882    /// # Example
1883    ///
1884    /// ```ignore
1885    /// use tokio_util::sync::CancellationToken;
1886    ///
1887    /// let manager = Arc::new(DeviceManager::new());
1888    /// let cancel = CancellationToken::new();
1889    /// let handle = manager.start_health_monitor(cancel.clone());
1890    ///
1891    /// // Later, to stop the health monitor:
1892    /// cancel.cancel();
1893    /// handle.await.unwrap();
1894    /// ```
1895    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    /// Add a device with priority.
1903    ///
1904    /// `identifier` must match the device exactly, as [`Device::connect`]
1905    /// describes.
1906    ///
1907    /// # Errors
1908    ///
1909    /// Returns [`Error::InvalidConfig`] if the config's
1910    /// `default_reconnect_options` are invalid (see [`ReconnectOptions::validate`]).
1911    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    /// Get the lowest priority connected device that could be disconnected.
1922    ///
1923    /// Returns None if no devices can be disconnected (all are Critical priority or not connected).
1924    pub async fn lowest_priority_connected(&self) -> Option<String> {
1925        self.core.lowest_priority_connected().await
1926    }
1927
1928    /// Disconnect the lowest priority device to make room for a new connection.
1929    ///
1930    /// Returns Ok(true) if a device was chosen, Ok(false) if no eligible device found.
1931    /// The health monitor doesn't reconnect the evicted device until
1932    /// [`connect`](Self::connect) is called for it. The device is disconnected
1933    /// as [`disconnect`](Self::disconnect) does it, so this too waits for a
1934    /// connect of that device that is running to end. A connect of the device
1935    /// that starts before the eviction has taken its connection keeps it
1936    /// connected, and then no slot is freed.
1937    pub async fn evict_lowest_priority(&self) -> Result<bool> {
1938        self.core.evict_lowest_priority().await
1939    }
1940
1941    /// Start hybrid monitoring using both passive (advertisement) and active connections.
1942    ///
1943    /// This is the most efficient way to monitor multiple devices:
1944    /// - **Passive monitoring**: Uses BLE advertisements to receive real-time readings
1945    ///   without maintaining connections. Lower power consumption, unlimited devices.
1946    /// - **Active connections**: Only established when needed (history download, settings changes).
1947    ///
1948    /// # Requirements
1949    ///
1950    /// Smart Home integration must be enabled on each device for passive monitoring.
1951    ///
1952    /// # Example
1953    ///
1954    /// ```ignore
1955    /// use tokio_util::sync::CancellationToken;
1956    ///
1957    /// let manager = Arc::new(DeviceManager::new());
1958    /// let cancel = CancellationToken::new();
1959    /// let handle = manager.start_hybrid_monitor(cancel.clone(), None);
1960    ///
1961    /// // Receive readings via manager events
1962    /// let mut rx = manager.events().subscribe();
1963    /// while let Ok(event) = rx.recv().await {
1964    ///     if let DeviceEvent::Reading { device, reading } = event {
1965    ///         println!("{}: CO2 = {} ppm", device.id, reading.co2);
1966    ///     }
1967    /// }
1968    /// ```
1969    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    /// Get a reading using hybrid approach: try passive first, fall back to active.
1979    ///
1980    /// This method checks if a recent passive reading is available. If not,
1981    /// it establishes an active connection to read the value.
1982    ///
1983    /// # Arguments
1984    ///
1985    /// * `identifier` - Device identifier
1986    /// * `max_passive_age` - Maximum age of passive reading to accept (default: 60s)
1987    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    /// Check if a device supports passive monitoring (Smart Home enabled).
1996    ///
1997    /// This performs a quick scan to check if the device is broadcasting
1998    /// advertisement data with sensor readings.
1999    ///
2000    /// Scans in one process run one at a time, in the order they asked to, so
2001    /// the check's 5 s scan first waits for the scan window that is running
2002    /// and for every window queued before it. It returns `false` if no reading
2003    /// arrives within 15 s, so when those windows take more than 10 s in all,
2004    /// it can miss a device that does advertise.
2005    ///
2006    /// The check's scanning stops when it returns, or as soon as this future
2007    /// is dropped, for example by a caller's timeout.
2008    pub async fn supports_passive_monitoring(&self, identifier: &str) -> bool {
2009        // Create a short-lived passive monitor to check for advertisements
2010        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        // Wait for a reading or timeout: 5 s of scanning, after up to 10 s of
2018        // waiting for the scan windows ahead of it.
2019        receives_within(readings, Duration::from_secs(15), |cancel| {
2020            monitor.start(cancel)
2021        })
2022        .await
2023    }
2024}
2025
2026/// Starts a monitor with `start` and waits up to `limit` for a value on
2027/// `readings`, which must be subscribed to that monitor before it starts.
2028/// Returns whether one arrived. The token given to `start` is cancelled,
2029/// which stops the monitor, when this returns or is dropped.
2030async 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    // Dropping this future, as a caller's timeout does, stops the monitor
2037    // too; otherwise it would scan for the rest of the process.
2038    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
2046/// Convert a passive advertisement reading to a CurrentReading.
2047fn passive_reading_to_current(passive: &PassiveReading) -> Option<CurrentReading> {
2048    let data = &passive.data;
2049
2050    // We need at least some sensor data to create a reading
2051    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, // Not available in advertisement data
2076    })
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        // Events are only emitted for actual device operations
2130        assert_eq!(manager.events().receiver_count(), 1);
2131    }
2132
2133    /// A stand-in for `supports_passive_monitoring`'s monitor: it runs until
2134    /// its token is cancelled, sending one reading on `sender` 3 s in if
2135    /// `reading` is set, and then tells `stopped`.
2136    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    /// The passive check returns whether a reading arrives in time, and
2153    /// stops its monitor when it returns.
2154    #[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    /// The passive check's monitor also stops when the check is dropped
2178    /// before it returns, as by a caller's timeout, instead of scanning for
2179    /// the rest of the process.
2180    #[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    /// The manager events received so far, one line each (`DeviceEvent`
2218    /// has no `PartialEq`).
2219    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    /// Every sensor whose link is up must be held by the manager.
2241    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    /// BR-9: a connect dropped by its caller's timeout left the `connecting`
2262    /// flag set, and every later connect returned Ok without connecting.
2263    #[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    /// BR-9: a second connect of the same device returned Ok at once while
2291    /// the first was still running, even when the first then failed.
2292    #[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            // The second attempt starts only after the first has failed.
2308            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    /// BR-9: the limit counted installed links only, so connects that were
2328    /// still running all passed the check.
2329    #[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            // A connect that the limit rejects doesn't add the device.
2354            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    /// BR-9: `connect_all` starts every connect at once, so they all passed
2364    /// the limit check.
2365    #[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    /// Guard: a failed or cancelled connect gives its slot back.
2401    #[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    /// BR-9: the slot was free as soon as the link was taken out of the map,
2440    /// while the sensor was still connected.
2441    #[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    /// BR-9: the new link was a local variable until the device-info read
2475    /// finished, so a cancel during that read dropped a live link.
2476    #[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            // As `DeviceManager::connect` documents: the device stays connected,
2498            // but has no device information, and no `Connected` event was sent.
2499            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    /// BR-9: the limit was checked when `connect` added a new device, but the
2510    /// slot was taken later. Two connects of new devices that waited for the
2511    /// device map together both passed the check and both added their device.
2512    #[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            // Another user of the device map (a read, a snapshot, a monitor)
2521            // holds it while both connects start.
2522            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    /// BR-9: a connect cancelled while it waited for the device map to store
2550    /// its new link dropped the link, whose `Drop` tore it down, and freed
2551    /// the slot before the link was down.
2552    #[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            // Hold the device map from 1 s, while the connect is under way, so
2565            // that at 2 s the new link waits to be stored until the caller
2566            // gives up at 3 s.
2567            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    /// BR-9, in BR-4's pattern: a disconnect whose caller gave up still ran
2598    /// to the end (as `Device::disconnect` does), but a connect of the same
2599    /// device made meanwhile didn't wait for it, so the old link's disconnect
2600    /// took the new link down while the manager kept the new handle.
2601    #[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            // Until well after the first link's disconnect has finished.
2616            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            // The second connect started only once the first link was down.
2624            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    /// BR-9, in BR-4's pattern: a connect of the same device made after a
2645    /// connect was cancelled while its new link waited for the device map
2646    /// didn't wait for that link to go down, so the abandoned link's
2647    /// disconnect could take the new link down.
2648    #[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            // As in test 9: the new link waits for the map from 2 s until the
2657            // caller gives up at 3 s.
2658            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            // Until well after the abandoned link's disconnect has finished.
2670            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            // The second connect started only once the abandoned link was down.
2678            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    /// A disconnect sends `Disconnected` once it has taken the link out of
2699    /// the manager, even if closing the link then fails (as closing a link
2700    /// that dropped on its own does on macOS): the manager reports the device
2701    /// as disconnected from then on, and nothing sends the event later. The
2702    /// error is still returned.
2703    #[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            // B has the lowest priority.
2722            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            // C is the only device still connected.
2727            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    /// `remove_device` disconnects first. When closing the link fails, the
2744    /// device is still removed, as the manager no longer holds a link that a
2745    /// second call could close, and the error is returned.
2746    #[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            // Gone: the health monitor has nothing to reconnect.
2761            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    /// The health monitor: dead links, the user's intent, backoff and
2772    /// cancellation.
2773    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        /// For the tests that run an hour or more of paused time.
2788        const LONG_LIMIT: Duration = Duration::from_secs(2 * 60 * 60);
2789
2790        /// The radio log without times and without `Op` entries.
2791        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        /// When each connect to `id` started.
2801        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        /// When `id`'s link was first closed with a disconnect.
2813        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        // ---- Dead links ----
2821
2822        #[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                // A link that is up: connecting again changes nothing.
2879                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                // The new link answers.
2932                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                // Documented limitation: without validation the monitor only
2945                // asks the stack, which still reports the zombie as connected.
2946                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                // 30 s base interval, adaptive, 5 s minimum.
3005                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                // B stays healthy on every tick: a failed repair of A still
3017                // counts as a failure.
3018                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                // Each tick with a failed repair halves the interval: 30 s, 15 s,
3029                // 7.5 s, then the 5 s minimum. The first start is `connect`.
3030                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                // A's check takes 3 s and fails (a validation read that times
3056                // out); B's link is simply gone.
3057                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                // Repairs run one at a time, highest priority first.
3080                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                // Added but not connected yet: the tick at 5 s connects it,
3110                // and that connect takes 60 s.
3111                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                // The abandoned connect never completes, and its slot is free.
3128                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        /// A `connect()` given up on while it closes a dead link keeps the
3142        /// device until that link is down, as a cancelled `disconnect()` does:
3143        /// a disconnect acts on the sensor, not on the handle, so a new link
3144        /// made before it finished would be taken down by it.
3145        #[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                // Given up on 0.5 s into closing the dead link, which takes 2 s.
3155                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                // Until well after the dead link's disconnect has finished.
3162                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                // The second connect started only once the dead link was down.
3170                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        /// The same for a health tick that the monitor's cancel drops while
3182        /// the tick closes a dead link.
3183        #[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                // The tick at 5 s finds the link dead and starts closing it,
3198                // which takes 2 s; the monitor is cancelled 0.5 s later.
3199                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                // Until well after the dead link's disconnect has finished.
3207                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                // The connect started only once the dead link was down.
3215                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        // ---- The user's intent, backoff and limits ----
3227
3228        #[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                // An added device is connected by the monitor, unless
3244                // `disconnect_all` withdrew it first, link or no link.
3245                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                // A tick every 5 s. The first repair runs on the tick that
3318                // finds the link dead; after that the waits follow the
3319                // options: 60 s, doubling up to 600 s.
3320                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                // One hour of ticks, 5 s apart.
3358                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                // An explicit connect starts over.
3378                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                // The repair's connect and then the user's each take 10 s and
3402                // fail.
3403                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                // The user connects while the repair runs: the connect waits
3412                // for the repair, which fails and uses up max_attempts, and
3413                // then makes its own attempt.
3414                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                // The explicit connect started over, although it failed too.
3430                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                // So the monitor keeps repairing the device.
3441                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                // One connect runs; the other waits for it.
3492                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                // The waiting connect gave up without a Bluetooth connect.
3510                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        /// A `disconnect()` or `remove_device()` withdraws the device before
3521        /// it waits for the device's lock, and from then on the health
3522        /// monitor leaves the device alone. So one dropped during that wait
3523        /// (here, for a check that finds the link alive) must still finish
3524        /// afterwards, or the device stays connected and unchecked.
3525        #[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                    // Each check takes 3 s and finds the link alive.
3533                    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                // Both wait for the checks, and are given up on at 2 s.
3543                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        /// `disconnect_all()` withdraws every device at once and then
3589        /// disconnects each one, so one dropped in between (here, while
3590        /// another user of the device map holds it) must still disconnect
3591        /// the devices it withdrew.
3592        #[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                // The device map is busy when disconnect_all() starts, and a
3601                // reader that asks for it next holds it from the moment
3602                // disconnect_all() has withdrawn the devices until 2 s.
3603                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        /// A `disconnect()` withdraws the device at once but takes the link
3639        /// only once its task runs. A `connect()` that starts in between
3640        /// (here, spawned right after it) re-arms the device, finds the link
3641        /// up and returns `Ok`, so the disconnect must then do nothing:
3642        /// taking the link would leave the device wanted but unlinked, which
3643        /// neither order of the two calls gives.
3644        #[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                // The disconnect took nothing: the link is still the first one.
3678                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        /// As in the test above, for `disconnect_all()` and
3694        /// `evict_lowest_priority()`, which disconnect each device the same
3695        /// way: a device that a `connect()` re-arms first stays connected,
3696        /// and the others are still disconnected.
3697        #[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                // A is the only device connected, so it is the one evicted.
3727                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        /// The same race when the `disconnect()` is dropped after its first
3752        /// poll: it has withdrawn the device and left the rest to its task,
3753        /// which hasn't run yet when a `connect()` of the device starts.
3754        #[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                // Until well after the dropped disconnect's task has run.
3772                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        /// A removal always goes ahead: a `connect()` spawned right after a
3794        /// `remove_device()` of the same device can return first, and the
3795        /// device is removed after it all the same, as when the two calls run
3796        /// one after the other.
3797        #[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                // Connected, then removed; or removed while the connect waited.
3818                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                // All four are due; A goes first, and its connect takes 10 s.
3845                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                // Meanwhile the user connects B, disconnects C and removes D.
3854                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                // The same options as the manager's default.
3903                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                // A holds the only slot. B waits for it: no connect, no failed
3943                // repair and no give-up, however many ticks pass.
3944                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                // Once A's slot is free, the next tick connects B.
3961                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        /// A repair that succeeds starts the backoff over: the next lost link
3988        /// is repaired at once, as attempt 1 again, with all of `max_attempts`.
3989        #[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                // The first repair fails; the one 60 s later succeeds.
4005                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                // The next loss is repaired as attempt 1 again, and one more
4035                // failure doesn't use up the two attempts.
4036                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        /// The same holds for a repair cancelled during the device-info read
4060        /// after its link was stored: the device stays connected, so its
4061        /// backoff starts over as soon as the link is stored.
4062        #[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                // The first repair fails.
4078                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                // The one 60 s later connects, and is cancelled 1 s into the
4088                // device-info read, which takes 5 s.
4089                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                // The next loss is repaired as attempt 1 again, and one more
4105                // failure doesn't use up the two attempts.
4106                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        /// A wait too long to add to an `Instant` is capped at
4131        /// `MAX_RETRY_WAIT`. These options pass `validate()`; without the cap
4132        /// the failed repair would panic and end the monitor for every device.
4133        #[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                // The tick at 5 s fails A's repair.
4158                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                // The monitor still runs: the tick at 10 s repairs B.
4169                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        /// A link leaves the manager when a check, a `connect()`, a
4187        /// `disconnect()` or a `remove_device()` starts closing it, and
4188        /// `Disconnected` is sent then: a caller dropped during the close,
4189        /// such as the monitor when it is cancelled or a call under a
4190        /// timeout, must not lose the event. A removal dropped during the
4191        /// close still removes the device once the link is down.
4192        #[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                // The tick at 5 s starts closing A's dead link, which takes
4210                // 2 s; the monitor is cancelled 0.5 s later.
4211                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                // A connect() given up on 0.5 s into closing B's dead link.
4217                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                // A disconnect() and a remove_device() given up on 0.5 s into
4221                // closing C's and D's links.
4222                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                // Until well after the closes have finished.
4241                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        /// `max_attempts` counts reconnect attempts: 0 allows none, so the
4252        /// monitor gives up on the first loss it finds without trying, and 1
4253        /// allows one.
4254        #[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                // A minute of ticks, 5 s apart.
4270                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        /// `max_attempts` limits reconnects, not the monitor's first connect
4303        /// of a device that has never been connected: with 0, an added device
4304        /// is still connected, or tried once, and the monitor gives up on it
4305        /// without trying only once it has been connected and is lost.
4306        #[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                // The first tick connects A and tries B once.
4320                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                // A minute of ticks, 5 s apart, after A is lost: the monitor
4341                // gives up on A without trying, and never tries B again.
4342                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}