aranet_core/
passive.rs

1//! Passive monitoring via BLE advertisements.
2//!
3//! This module provides functionality to monitor Aranet devices without
4//! establishing a connection, using BLE advertisement data instead.
5//!
6//! # Benefits
7//!
8//! - **Lower power consumption**: No connection overhead
9//! - **More devices**: Can monitor more than the BLE connection limit
10//! - **Simpler**: No connection management needed
11//!
12//! # Requirements
13//!
14//! Smart Home integration must be enabled on each device:
15//! - Go to device Settings > Smart Home > Enable
16//!
17//! # Example
18//!
19//! ```ignore
20//! use aranet_core::passive::{PassiveMonitor, PassiveMonitorOptions};
21//! use tokio_util::sync::CancellationToken;
22//!
23//! let monitor = PassiveMonitor::new(PassiveMonitorOptions::default());
24//! let cancel = CancellationToken::new();
25//!
26//! // Start monitoring in background
27//! let handle = monitor.start(cancel.clone());
28//!
29//! // Receive readings
30//! let mut rx = monitor.subscribe();
31//! while let Ok(reading) = rx.recv().await {
32//!     println!("Device: {} CO2: {:?}", reading.device_name, reading.data.co2);
33//! }
34//! ```
35
36use std::collections::HashMap;
37use std::sync::Arc;
38use std::time::Duration;
39
40use btleplug::api::{Central, Peripheral as _, ScanFilter};
41use tokio::sync::{RwLock, broadcast};
42use tokio::time::sleep;
43use tokio_util::sync::CancellationToken;
44use tracing::{debug, info, warn};
45
46use crate::advertisement::{AdvertisementData, parse_advertisement_with_name};
47use crate::error::Result;
48use crate::scan::get_adapter;
49use crate::uuid::MANUFACTURER_ID;
50
51/// Bitwise-exact comparison of two `Option<f32>` values (handles NaN correctly).
52fn opt_f32_eq(a: Option<f32>, b: Option<f32>) -> bool {
53    match (a, b) {
54        (Some(x), Some(y)) => x.to_bits() == y.to_bits(),
55        (None, None) => true,
56        _ => false,
57    }
58}
59
60/// A reading from passive advertisement monitoring.
61#[derive(Debug, Clone)]
62pub struct PassiveReading {
63    /// Device identifier (MAC address or UUID).
64    pub device_id: String,
65    /// Device name if available.
66    pub device_name: Option<String>,
67    /// RSSI signal strength.
68    pub rssi: Option<i16>,
69    /// Parsed advertisement data.
70    pub data: AdvertisementData,
71    /// When this reading was received.
72    pub received_at: std::time::Instant,
73}
74
75/// Options for passive monitoring.
76#[derive(Debug, Clone)]
77pub struct PassiveMonitorOptions {
78    /// How long to scan between processing cycles.
79    pub scan_duration: Duration,
80    /// Delay between scan cycles.
81    pub scan_interval: Duration,
82    /// Channel capacity for readings.
83    pub channel_capacity: usize,
84    /// Only emit readings when values change (deduplicate).
85    pub deduplicate: bool,
86    /// Maximum age of cached readings before re-emitting (if deduplicate is true).
87    pub max_reading_age: Duration,
88    /// Filter to only these device IDs (empty = all Aranet devices).
89    pub device_filter: Vec<String>,
90}
91
92impl Default for PassiveMonitorOptions {
93    fn default() -> Self {
94        Self {
95            scan_duration: Duration::from_secs(5),
96            scan_interval: Duration::from_secs(1),
97            channel_capacity: 100,
98            deduplicate: true,
99            max_reading_age: Duration::from_secs(60),
100            device_filter: Vec::new(),
101        }
102    }
103}
104
105impl PassiveMonitorOptions {
106    /// Create new options with default settings.
107    pub fn new() -> Self {
108        Self::default()
109    }
110
111    /// Set the scan duration.
112    pub fn scan_duration(mut self, duration: Duration) -> Self {
113        self.scan_duration = duration;
114        self
115    }
116
117    /// Set the interval between scan cycles.
118    pub fn scan_interval(mut self, interval: Duration) -> Self {
119        self.scan_interval = interval;
120        self
121    }
122
123    /// Enable or disable deduplication.
124    pub fn deduplicate(mut self, enable: bool) -> Self {
125        self.deduplicate = enable;
126        self
127    }
128
129    /// Filter to specific device IDs.
130    pub fn filter_devices(mut self, device_ids: Vec<String>) -> Self {
131        self.device_filter = device_ids;
132        self
133    }
134}
135
136/// Cached reading for deduplication.
137struct CachedReading {
138    data: AdvertisementData,
139    received_at: std::time::Instant,
140}
141
142/// Passive monitor for Aranet devices using BLE advertisements.
143///
144/// This allows monitoring multiple devices without establishing connections,
145/// which is useful for scenarios where:
146/// - You need to monitor more devices than the BLE connection limit
147/// - Low power consumption is important
148/// - Real-time data isn't critical (advertisement interval is typically 4+ seconds)
149pub struct PassiveMonitor {
150    options: PassiveMonitorOptions,
151    /// Broadcast sender for readings.
152    sender: broadcast::Sender<PassiveReading>,
153    /// Cache of last readings for deduplication.
154    cache: Arc<RwLock<HashMap<String, CachedReading>>>,
155}
156
157impl PassiveMonitor {
158    /// Create a new passive monitor with the given options.
159    pub fn new(options: PassiveMonitorOptions) -> Self {
160        let (sender, _) = broadcast::channel(options.channel_capacity);
161        Self {
162            options,
163            sender,
164            cache: Arc::new(RwLock::new(HashMap::new())),
165        }
166    }
167
168    /// Subscribe to passive readings.
169    ///
170    /// Returns a receiver that will receive readings as they are detected.
171    pub fn subscribe(&self) -> broadcast::Receiver<PassiveReading> {
172        self.sender.subscribe()
173    }
174
175    /// Get the number of active subscribers.
176    pub fn subscriber_count(&self) -> usize {
177        self.sender.receiver_count()
178    }
179
180    /// Start the passive monitor.
181    ///
182    /// This spawns a background task that continuously scans for BLE
183    /// advertisements and parses Aranet device data.
184    ///
185    /// The task stops as soon as the cancellation token is triggered, even
186    /// while it gets the Bluetooth adapter, scans or waits.
187    pub fn start(self: &Arc<Self>, cancel_token: CancellationToken) -> tokio::task::JoinHandle<()> {
188        let monitor = Arc::clone(self);
189
190        tokio::spawn(async move {
191            info!("Starting passive monitor");
192
193            // Acquire the adapter once and reuse across scan cycles.
194            // On persistent errors we re-acquire it in case the adapter
195            // was reset or the D-Bus connection was lost.
196            let Some(mut adapter) = first_adapter(&cancel_token, get_adapter).await else {
197                info!("Passive monitor cancelled while waiting for adapter");
198                return;
199            };
200            let mut consecutive_errors: u32 = 0;
201
202            loop {
203                // One cycle: a scan, its error handling and the wait before the
204                // next scan. The whole cycle, waits included, ends as soon as the
205                // monitor is cancelled.
206                let cycle = async {
207                    match monitor.scan_cycle_with_adapter(&adapter).await {
208                        Ok(()) => {
209                            consecutive_errors = 0;
210                        }
211                        Err(e) => {
212                            consecutive_errors += 1;
213                            warn!(
214                                "Passive monitor scan error ({consecutive_errors} consecutive): {e}"
215                            );
216                            // After several consecutive failures, try to
217                            // re-acquire the adapter — it may have been
218                            // reset or the D-Bus connection may have died.
219                            if consecutive_errors >= 5 {
220                                warn!(
221                                    "Passive monitor: re-acquiring adapter after {} consecutive errors",
222                                    consecutive_errors
223                                );
224                                match get_adapter().await {
225                                    Ok(a) => {
226                                        adapter = a;
227                                        info!("Passive monitor: adapter re-acquired");
228                                        consecutive_errors = 0;
229                                    }
230                                    Err(e2) => {
231                                        // Adapter re-acquisition failed — back off
232                                        // longer to avoid thrashing when the adapter
233                                        // is permanently unavailable.
234                                        warn!(
235                                            "Passive monitor: failed to re-acquire adapter: {}. Backing off.",
236                                            e2
237                                        );
238                                        let backoff = std::cmp::min(
239                                            monitor
240                                                .options
241                                                .scan_interval
242                                                .saturating_mul(consecutive_errors),
243                                            std::time::Duration::from_secs(300),
244                                        );
245                                        sleep(backoff).await;
246                                    }
247                                }
248                            }
249                        }
250                    }
251                    // Wait before next scan cycle
252                    sleep(monitor.options.scan_interval).await;
253                };
254
255                tokio::select! {
256                    _ = cancel_token.cancelled() => {
257                        info!("Passive monitor cancelled");
258                        break;
259                    }
260                    () = cycle => {}
261                }
262            }
263        })
264    }
265
266    /// Perform a single scan cycle using a pre-existing adapter.
267    async fn scan_cycle_with_adapter(&self, adapter: &btleplug::platform::Adapter) -> Result<()> {
268        // Scan under the process-wide permit. The scan stops even if this cycle
269        // is cancelled part-way through.
270        let permit = crate::scan::scan_lock().acquire().await;
271        crate::scan::run_scan(
272            adapter,
273            permit,
274            ScanFilter::default(),
275            self.options.scan_duration,
276        )
277        .await?;
278
279        // Process discovered peripherals
280        let peripherals = adapter.peripherals().await?;
281
282        for peripheral in peripherals {
283            if let Ok(Some(props)) = peripheral.properties().await {
284                // Check if this is an Aranet device by manufacturer data
285                if let Some(data) = props.manufacturer_data.get(&MANUFACTURER_ID) {
286                    let device_id = crate::util::create_identifier(
287                        &props.address.to_string(),
288                        &peripheral.id(),
289                    );
290
291                    // Check device filter
292                    if !self.options.device_filter.is_empty()
293                        && !self.options.device_filter.contains(&device_id)
294                    {
295                        continue;
296                    }
297
298                    // Try to parse the advertisement
299                    match parse_advertisement_with_name(data, props.local_name.as_deref()) {
300                        Ok(adv_data) => {
301                            // Check for deduplication
302                            let should_emit = if self.options.deduplicate {
303                                self.should_emit(&device_id, &adv_data).await
304                            } else {
305                                true
306                            };
307
308                            if should_emit {
309                                let reading = PassiveReading {
310                                    device_id: device_id.clone(),
311                                    device_name: props.local_name.clone(),
312                                    rssi: props.rssi,
313                                    data: adv_data.clone(),
314                                    received_at: std::time::Instant::now(),
315                                };
316
317                                // Update cache
318                                self.cache.write().await.insert(
319                                    device_id,
320                                    CachedReading {
321                                        data: adv_data,
322                                        received_at: std::time::Instant::now(),
323                                    },
324                                );
325
326                                // Send to subscribers (ignore if no receivers)
327                                let _ = self.sender.send(reading);
328                            }
329                        }
330                        Err(e) => {
331                            debug!("Failed to parse advertisement from {}: {}", device_id, e);
332                        }
333                    }
334                }
335            }
336        }
337
338        Ok(())
339    }
340
341    /// Check if a reading should be emitted (for deduplication).
342    async fn should_emit(&self, device_id: &str, data: &AdvertisementData) -> bool {
343        let cache = self.cache.read().await;
344
345        if let Some(cached) = cache.get(device_id) {
346            // Check if reading is too old
347            if cached.received_at.elapsed() > self.options.max_reading_age {
348                return true;
349            }
350
351            // Check if values have changed (use total_cmp for floats to handle NaN correctly)
352            if cached.data.co2 != data.co2
353                || !opt_f32_eq(cached.data.temperature, data.temperature)
354                || cached.data.humidity != data.humidity
355                || !opt_f32_eq(cached.data.pressure, data.pressure)
356                || cached.data.radon != data.radon
357                || !opt_f32_eq(cached.data.radiation_dose_rate, data.radiation_dose_rate)
358                || cached.data.battery != data.battery
359            {
360                return true;
361            }
362
363            // Check if counter changed (new measurement)
364            if cached.data.counter != data.counter {
365                return true;
366            }
367
368            false
369        } else {
370            // Not in cache, emit
371            true
372        }
373    }
374
375    /// Get the last known reading for a device.
376    pub async fn get_last_reading(&self, device_id: &str) -> Option<AdvertisementData> {
377        let cache = self.cache.read().await;
378        cache.get(device_id).map(|c| c.data.clone())
379    }
380
381    /// Get all known device IDs.
382    pub async fn known_devices(&self) -> Vec<String> {
383        let cache = self.cache.read().await;
384        cache.keys().cloned().collect()
385    }
386
387    /// Clear the reading cache.
388    pub async fn clear_cache(&self) {
389        self.cache.write().await.clear();
390    }
391}
392
393impl Default for PassiveMonitor {
394    fn default() -> Self {
395        Self::new(PassiveMonitorOptions::default())
396    }
397}
398
399/// The adapter that a passive monitor starts with, from `get` (`get_adapter` in
400/// the library), which is called again 10 s after each failure. `None` as soon
401/// as `cancel` is triggered, even while `get` is still waiting for an answer,
402/// which on Linux can take up to the 30 s D-Bus timeout when bluetoothd is slow.
403async fn first_adapter<A, F>(cancel: &CancellationToken, mut get: impl FnMut() -> F) -> Option<A>
404where
405    F: Future<Output = Result<A>>,
406{
407    loop {
408        match cancel.run_until_cancelled(get()).await? {
409            Ok(adapter) => return Some(adapter),
410            Err(e) => {
411                warn!("Passive monitor failed to get adapter: {e} — retrying in 10s");
412                cancel
413                    .run_until_cancelled(sleep(Duration::from_secs(10)))
414                    .await?;
415            }
416        }
417    }
418}
419
420#[cfg(test)]
421mod tests {
422    use super::*;
423
424    use crate::error::{DeviceNotFoundReason, Error};
425    use crate::test_support::within;
426
427    #[test]
428    fn test_passive_monitor_options_default() {
429        let opts = PassiveMonitorOptions::default();
430        assert_eq!(opts.scan_duration, Duration::from_secs(5));
431        assert!(opts.deduplicate);
432        assert!(opts.device_filter.is_empty());
433    }
434
435    #[test]
436    fn test_passive_monitor_options_builder() {
437        let opts = PassiveMonitorOptions::new()
438            .scan_duration(Duration::from_secs(10))
439            .deduplicate(false)
440            .filter_devices(vec!["device1".to_string()]);
441
442        assert_eq!(opts.scan_duration, Duration::from_secs(10));
443        assert!(!opts.deduplicate);
444        assert_eq!(opts.device_filter, vec!["device1"]);
445    }
446
447    #[test]
448    fn test_passive_monitor_subscribe() {
449        let monitor = Arc::new(PassiveMonitor::default());
450        let _rx1 = monitor.subscribe();
451        let _rx2 = monitor.subscribe();
452        assert_eq!(monitor.subscriber_count(), 2);
453    }
454
455    /// Helper to create test advertisement data with sensible defaults.
456    fn make_adv_data() -> AdvertisementData {
457        AdvertisementData {
458            device_type: aranet_types::DeviceType::Aranet4,
459            co2: Some(800),
460            temperature: Some(22.5),
461            pressure: Some(1013.2),
462            humidity: Some(45),
463            battery: 85,
464            status: aranet_types::Status::Green,
465            interval: 300,
466            age: 120,
467            radon: None,
468            radiation_dose_rate: None,
469            counter: Some(5),
470            flags: 0x22,
471        }
472    }
473
474    #[tokio::test]
475    async fn test_should_emit_first_reading() {
476        let monitor = PassiveMonitor::default();
477        let data = make_adv_data();
478
479        // First reading for a device should always be emitted.
480        assert!(monitor.should_emit("device-1", &data).await);
481    }
482
483    #[tokio::test]
484    async fn test_should_emit_duplicate_suppressed() {
485        let monitor = PassiveMonitor::default();
486        let data = make_adv_data();
487
488        // Populate the cache.
489        monitor.cache.write().await.insert(
490            "device-1".to_string(),
491            CachedReading {
492                data: data.clone(),
493                received_at: std::time::Instant::now(),
494            },
495        );
496
497        // Identical reading should be suppressed.
498        assert!(!monitor.should_emit("device-1", &data).await);
499    }
500
501    #[tokio::test]
502    async fn test_should_emit_on_value_change() {
503        let monitor = PassiveMonitor::default();
504        let data = make_adv_data();
505
506        monitor.cache.write().await.insert(
507            "device-1".to_string(),
508            CachedReading {
509                data: data.clone(),
510                received_at: std::time::Instant::now(),
511            },
512        );
513
514        // Changed CO2 should trigger emission.
515        let mut changed = data.clone();
516        changed.co2 = Some(900);
517        assert!(monitor.should_emit("device-1", &changed).await);
518
519        // Changed battery should trigger emission.
520        let mut changed = data.clone();
521        changed.battery = 50;
522        assert!(monitor.should_emit("device-1", &changed).await);
523
524        // Changed temperature should trigger emission.
525        let mut changed = data;
526        changed.temperature = Some(23.0);
527        assert!(monitor.should_emit("device-1", &changed).await);
528    }
529
530    #[tokio::test]
531    async fn test_should_emit_on_counter_change() {
532        let monitor = PassiveMonitor::default();
533        let data = make_adv_data();
534
535        monitor.cache.write().await.insert(
536            "device-1".to_string(),
537            CachedReading {
538                data: data.clone(),
539                received_at: std::time::Instant::now(),
540            },
541        );
542
543        // Counter increment means a new measurement was taken.
544        let mut changed = data;
545        changed.counter = Some(6);
546        assert!(monitor.should_emit("device-1", &changed).await);
547    }
548
549    #[tokio::test]
550    async fn test_should_emit_on_stale_cache() {
551        let opts = PassiveMonitorOptions {
552            max_reading_age: Duration::from_millis(10),
553            ..Default::default()
554        };
555        let monitor = PassiveMonitor::new(opts);
556        let data = make_adv_data();
557
558        // Insert a reading that is already "old".
559        monitor.cache.write().await.insert(
560            "device-1".to_string(),
561            CachedReading {
562                data: data.clone(),
563                received_at: std::time::Instant::now() - Duration::from_millis(50),
564            },
565        );
566
567        // Identical data should still be emitted because the cache entry expired.
568        assert!(monitor.should_emit("device-1", &data).await);
569    }
570
571    #[tokio::test]
572    async fn test_should_emit_different_device() {
573        let monitor = PassiveMonitor::default();
574        let data = make_adv_data();
575
576        // Cache a reading for device-1.
577        monitor.cache.write().await.insert(
578            "device-1".to_string(),
579            CachedReading {
580                data: data.clone(),
581                received_at: std::time::Instant::now(),
582            },
583        );
584
585        // device-2 has no cache entry, so it should emit even with identical data.
586        assert!(monitor.should_emit("device-2", &data).await);
587    }
588
589    // ==================== First Adapter Tests ====================
590
591    /// Longest any first-adapter test may take on the paused clock.
592    const TEST_LIMIT: Duration = Duration::from_secs(600);
593
594    /// What `get_adapter` returns when the system has no Bluetooth adapter.
595    fn no_adapter() -> Error {
596        Error::DeviceNotFound(DeviceNotFoundReason::NoAdapter)
597    }
598
599    #[tokio::test(start_paused = true)]
600    async fn a_cancel_stops_the_first_adapter_fetch() {
601        within(TEST_LIMIT, async {
602            // A fetch that never answers, like a D-Bus call to a stalled
603            // bluetoothd, cancelled 1 s in.
604            let cancel = CancellationToken::new();
605            let started = tokio::time::Instant::now();
606            let (adapter, ()) = tokio::join!(
607                within(
608                    Duration::from_secs(10),
609                    first_adapter(&cancel, std::future::pending::<Result<()>>)
610                ),
611                async {
612                    sleep(Duration::from_secs(1)).await;
613                    cancel.cancel();
614                }
615            );
616            assert_eq!(adapter, None);
617            assert_eq!(started.elapsed(), Duration::from_secs(1));
618        })
619        .await;
620    }
621
622    #[tokio::test(start_paused = true)]
623    async fn the_first_adapter_fetch_is_retried_every_10_s_until_cancelled() {
624        within(TEST_LIMIT, async {
625            // Two failures, then an adapter from the third fetch, 20 s in.
626            let calls = std::cell::Cell::new(0);
627            let started = tokio::time::Instant::now();
628            let adapter = first_adapter(&CancellationToken::new(), || {
629                calls.set(calls.get() + 1);
630                let call = calls.get();
631                async move {
632                    if call < 3 {
633                        Err(no_adapter())
634                    } else {
635                        Ok(call)
636                    }
637                }
638            })
639            .await;
640            assert_eq!(adapter, Some(3));
641            assert_eq!(started.elapsed(), Duration::from_secs(20));
642
643            // Fetches that keep failing, cancelled 15 s in, during the wait
644            // before the third.
645            let cancel = CancellationToken::new();
646            let started = tokio::time::Instant::now();
647            let (adapter, ()) = tokio::join!(
648                first_adapter(&cancel, || async { Err::<(), _>(no_adapter()) }),
649                async {
650                    sleep(Duration::from_secs(15)).await;
651                    cancel.cancel();
652                }
653            );
654            assert_eq!(adapter, None);
655            assert_eq!(started.elapsed(), Duration::from_secs(15));
656        })
657        .await;
658    }
659}