Struct DeviceManager

Source
pub struct DeviceManager { /* private fields */ }
Expand description

Manager for multiple Aranet devices.

Implementations§

Source§

impl DeviceManager

Source

pub fn new() -> Self

Create a new device manager.

Source

pub fn with_event_capacity(capacity: usize) -> Self

Create a manager with custom event capacity.

Source

pub fn with_config(config: ManagerConfig) -> Self

Create a manager with full configuration.

Source

pub fn events(&self) -> &EventDispatcher

Get the event dispatcher for subscribing to events.

Source

pub fn config(&self) -> &ManagerConfig

Get the manager configuration.

Source

pub async fn scan(&self) -> Result<Vec<DiscoveredDevice>>

Scan for available devices.

Source

pub async fn scan_with_options( &self, options: ScanOptions, ) -> Result<Vec<DiscoveredDevice>>

Scan with custom options.

Source

pub async fn add_device(&self, identifier: &str) -> Result<()>

Add a device to the manager by identifier.

identifier must match the device exactly, as Device::connect describes.

§Errors

Returns Error::InvalidConfig if the config’s default_reconnect_options are invalid (see ReconnectOptions::validate).

Source

pub async fn add_device_with_options( &self, identifier: &str, reconnect_options: ReconnectOptions, ) -> Result<()>

Add a device with custom reconnect options.

identifier must match the device exactly, as Device::connect describes.

The health monitor waits between automatic reconnects of the device as reconnect_options say. If the device is already managed, nothing changes.

§Errors

Returns Error::InvalidConfig if reconnect_options are invalid (see ReconnectOptions::validate).

Source

pub async fn connect(&self, identifier: &str) -> Result<()>

Connect to a device.

identifier must match the device exactly, as Device::connect describes.

This method performs an atomic connect-or-skip operation:

  • If the device doesn’t exist, it’s added and connected
  • If the device exists but is not connected, it’s connected
  • If the device already has a connection that the Bluetooth stack reports as up, this is a no-op; a lost one is closed and replaced

A second connect of the same device waits for the first and returns its own real result: Ok(()) if the first one connected the device, otherwise the result of its own attempt. A connect that hasn’t connected yet when disconnect, disconnect_all, evict_lowest_priority or remove_device is called for the same device returns Error::Cancelled, closing its connection if one comes up, unless another connect() re-arms the device first. A connect that waits for a remove_device of the same device returns Error::Cancelled instead of adding the device back.

If this future is dropped while a lost connection is being closed, the close still finishes in the background, and a connect of the same device waits for it instead of running alongside it.

If this future is dropped after the connection is made but before the device information has been read, the device stays connected, but no DeviceEvent::Connected is sent and get_device_info returns None until the device is reconnected.

§Connection Limits

If max_concurrent_connections is set in the config and would be exceeded, this method returns an error. The limit counts connects in progress; use can_connect() or available_connections() to check it before calling this method.

The device map is locked only while the entry is updated, not during the BLE connection, so operations on other devices don’t wait for it; operations on the same device do, as described above.

Source

pub async fn disconnect(&self, identifier: &str) -> Result<()>

Disconnect from a device.

A connect of the same device that hasn’t connected yet is abandoned: it returns Error::Cancelled, closing its connection if one comes up, unless another connect() re-arms the device first. This waits until that connect’s Bluetooth attempt has ended and such a connection is closed, which can take the connect’s whole time budget: tens of seconds at default settings. The health monitor doesn’t reconnect the device until connect is called.

DeviceEvent::Disconnected, with DisconnectReason::UserRequested, is sent as soon as the manager lets go of the connection, before the connection is closed, so it is sent even if closing the connection fails or this future is dropped. A failed close’s error is still returned.

If this future is dropped, a disconnect that has started still finishes in the background, even one still waiting for a connect or a health check of the device to end.

Whether or not this future is dropped, a connect of the same device that starts while the disconnect runs either waits for it and then reconnects, or, if it starts before the disconnect has taken the connection, keeps the device connected, and the disconnect then does nothing.

Source

pub async fn remove_device(&self, identifier: &str) -> Result<()>

Remove a device from the manager.

The device is disconnected first, as disconnect does, and then removed, even if closing its connection fails: that error is still returned. A connect of the device that hasn’t connected yet is abandoned: it returns Error::Cancelled, closing its connection if one comes up, unless another connect() re-arms the device first, in which case it connects, and the device is removed after that. As with disconnect, this waits until that connect’s Bluetooth attempt has ended and such a connection is closed, and a removal that has started still finishes in the background if this future is dropped.

Whether or not this future is dropped, a connect of the same device that starts while the removal runs doesn’t keep it: the device is removed even when that connect returns first.

Source

pub async fn device_ids(&self) -> Vec<String>

Get a list of all managed device IDs.

Source

pub async fn device_count(&self) -> usize

Get the number of managed devices.

Source

pub async fn connected_count(&self) -> usize

Get the number of connected devices (fast, doesn’t query BLE).

This returns the number of devices that have an active device handle, without querying the BLE stack. Use connected_count_verified for an accurate count that queries each device.

Source

pub async fn can_connect(&self) -> bool

Check if a new connection can be made without exceeding the limit.

Returns true if another connection can be made, false if at limit. Connects in progress count against the limit. Always returns true if max_concurrent_connections is 0 (unlimited).

Source

pub async fn connection_status(&self) -> (usize, usize)

Get the connection limit status.

Returns (current_connections, max_connections). If max is 0, there is no limit.

current_connections counts devices with a connection handle; connects in progress are not included but already hold a slot, so use available_connections or can_connect to see whether connect() will pass the limit.

Source

pub async fn available_connections(&self) -> Option<usize>

Get the number of available connection slots.

Connects in progress hold a slot, and so does a disconnect until the link is down. Returns None if there is no connection limit (unlimited).

Source

pub async fn connected_count_verified(&self) -> usize

Get the number of connected devices (verified via BLE).

This method queries each device to verify its connection status. The lock is released before making BLE calls to avoid contention.

Source

pub async fn read_current(&self, identifier: &str) -> Result<CurrentReading>

Read current values from a specific device.

Source

pub async fn read_all(&self) -> HashMap<String, Result<CurrentReading>>

Read current values from all connected devices (in parallel).

This method releases the lock before performing async BLE operations, allowing other tasks to add/remove devices while reads are in progress. All reads are performed in parallel for maximum performance.

Source

pub async fn connect_all(&self) -> HashMap<String, Result<()>>

Connect to all known devices (in parallel).

Returns a map of device IDs to connection results.

Source

pub async fn disconnect_all(&self) -> HashMap<String, Result<()>>

Disconnect from all devices (in parallel).

Returns a map of device IDs to disconnection results, with an entry for each device that had a connection. Every managed device is withdrawn, connected or not: the health monitor reconnects none of them until connect is called. Each connected device is disconnected as disconnect describes.

Source

pub fn try_is_connected(&self, identifier: &str) -> Option<bool>

Check if a specific device is connected (fast, doesn’t query BLE).

This method attempts to check if a device has an active connection handle without blocking. Returns None if the lock couldn’t be acquired immediately, or Some(bool) indicating whether the device has a connection handle.

Note: This only checks if we have a device handle, not whether the actual BLE connection is still alive. Use is_connected for a verified check.

Source

pub async fn is_connected(&self, identifier: &str) -> bool

Check if a specific device is connected (verified via BLE).

The lock is released before making the BLE call.

Source

pub async fn get_device_info(&self, identifier: &str) -> Option<DeviceInfo>

Get device info for a specific device.

Source

pub async fn get_last_reading(&self, identifier: &str) -> Option<CurrentReading>

Get the last cached reading for a device.

Source

pub fn start_health_monitor( self: &Arc<Self>, cancel_token: CancellationToken, ) -> JoinHandle<()>

Start a background task that checks the managed devices and repairs lost connections.

On every tick the task:

  1. Checks every connected device at the same time, so a slow device doesn’t delay the others. A connection that fails its check is disconnected explicitly, and DeviceEvent::Disconnected is emitted.
  2. Reconnects the devices that should be connected but aren’t, one at a time and highest DevicePriority first, emitting DeviceEvent::ReconnectStarted and DeviceEvent::ReconnectSucceeded.

Each device waits between reconnect attempts as its ReconnectOptions say. After max_attempts failures in a row the task stops trying and emits one DeviceEvent::Error; connect starts over. With max_attempts of 0 the task never reconnects a device that has been connected: it gives up as soon as it finds the device’s connection gone, without trying. A device that has never been connected, such as one just added, still gets one attempt. While max_concurrent_connections connections are in use, a device waits for a free one, and that wait isn’t counted as a failed attempt.

Devices added with add_device and its variants are connected by the task. Devices that were disconnected (with disconnect, disconnect_all or evict_lowest_priority) or removed are not reconnected until connect() is called for them.

The task runs until the provided cancellation token is cancelled, and then stops at once, even in the middle of a check or a reconnect. A connect it abandons releases the sensor, except one abandoned after the connection is made but before the device information has been read: that one keeps its connection, and no DeviceEvent::ReconnectSucceeded follows its DeviceEvent::ReconnectStarted. A lost connection it was closing is still closed in the background, and a connect of that device waits for it.

§Adaptive Intervals

If use_adaptive_interval is enabled in the config, the time between ticks adapts to connection stability:

  • after a tick with a failed reconnect, it halves (down to min_health_check_interval);
  • after three ticks with a healthy device and no reconnect, without a failed reconnect in between, it doubles (up to max_health_check_interval).
§Connection Validation

If use_connection_validation is enabled, health checks read the current measurements (device.validate_connection(), which needs no pairing that a reading doesn’t) to catch “zombie connections” where the BLE stack thinks it’s connected but the device is out of range. Otherwise they only ask the BLE stack, which misses zombie connections.

§Example
ⓘ
use tokio_util::sync::CancellationToken;

let manager = Arc::new(DeviceManager::new());
let cancel = CancellationToken::new();
let handle = manager.start_health_monitor(cancel.clone());

// Later, to stop the health monitor:
cancel.cancel();
handle.await.unwrap();
Source

pub async fn add_device_with_priority( &self, identifier: &str, priority: DevicePriority, ) -> Result<()>

Add a device with priority.

identifier must match the device exactly, as Device::connect describes.

§Errors

Returns Error::InvalidConfig if the config’s default_reconnect_options are invalid (see ReconnectOptions::validate).

Source

pub async fn lowest_priority_connected(&self) -> Option<String>

Get the lowest priority connected device that could be disconnected.

Returns None if no devices can be disconnected (all are Critical priority or not connected).

Source

pub async fn evict_lowest_priority(&self) -> Result<bool>

Disconnect the lowest priority device to make room for a new connection.

Returns Ok(true) if a device was chosen, Ok(false) if no eligible device found. The health monitor doesn’t reconnect the evicted device until connect is called for it. The device is disconnected as disconnect does it, so this too waits for a connect of that device that is running to end. A connect of the device that starts before the eviction has taken its connection keeps it connected, and then no slot is freed.

Source

pub fn start_hybrid_monitor( self: &Arc<Self>, cancel_token: CancellationToken, passive_options: Option<PassiveMonitorOptions>, ) -> JoinHandle<()>

Start hybrid monitoring using both passive (advertisement) and active connections.

This is the most efficient way to monitor multiple devices:

  • Passive monitoring: Uses BLE advertisements to receive real-time readings without maintaining connections. Lower power consumption, unlimited devices.
  • Active connections: Only established when needed (history download, settings changes).
§Requirements

Smart Home integration must be enabled on each device for passive monitoring.

§Example
ⓘ
use tokio_util::sync::CancellationToken;

let manager = Arc::new(DeviceManager::new());
let cancel = CancellationToken::new();
let handle = manager.start_hybrid_monitor(cancel.clone(), None);

// Receive readings via manager events
let mut rx = manager.events().subscribe();
while let Ok(event) = rx.recv().await {
    if let DeviceEvent::Reading { device, reading } = event {
        println!("{}: CO2 = {} ppm", device.id, reading.co2);
    }
}
Source

pub async fn read_hybrid( &self, identifier: &str, max_passive_age: Option<Duration>, ) -> Result<CurrentReading>

Get a reading using hybrid approach: try passive first, fall back to active.

This method checks if a recent passive reading is available. If not, it establishes an active connection to read the value.

§Arguments
  • identifier - Device identifier
  • max_passive_age - Maximum age of passive reading to accept (default: 60s)
Source

pub async fn supports_passive_monitoring(&self, identifier: &str) -> bool

Check if a device supports passive monitoring (Smart Home enabled).

This performs a quick scan to check if the device is broadcasting advertisement data with sensor readings.

Scans in one process run one at a time, in the order they asked to, so the check’s 5 s scan first waits for the scan window that is running and for every window queued before it. It returns false if no reading arrives within 15 s, so when those windows take more than 10 s in all, it can miss a device that does advertise.

The check’s scanning stops when it returns, or as soon as this future is dropped, for example by a caller’s timeout.

Trait Implementations§

Source§

impl Default for DeviceManager

Source§

fn default() -> Self

Returns the “default value” for a type. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

§

impl<T> PolicyExt for T
where T: ?Sized,

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] only if self and other return Action::Follow. Read more
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] if either self or other returns Action::Follow. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<A, B, T> HttpServerConnExec<A, B> for T
where B: Body,

§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,