1use 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
51fn 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#[derive(Debug, Clone)]
62pub struct PassiveReading {
63 pub device_id: String,
65 pub device_name: Option<String>,
67 pub rssi: Option<i16>,
69 pub data: AdvertisementData,
71 pub received_at: std::time::Instant,
73}
74
75#[derive(Debug, Clone)]
77pub struct PassiveMonitorOptions {
78 pub scan_duration: Duration,
80 pub scan_interval: Duration,
82 pub channel_capacity: usize,
84 pub deduplicate: bool,
86 pub max_reading_age: Duration,
88 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 pub fn new() -> Self {
108 Self::default()
109 }
110
111 pub fn scan_duration(mut self, duration: Duration) -> Self {
113 self.scan_duration = duration;
114 self
115 }
116
117 pub fn scan_interval(mut self, interval: Duration) -> Self {
119 self.scan_interval = interval;
120 self
121 }
122
123 pub fn deduplicate(mut self, enable: bool) -> Self {
125 self.deduplicate = enable;
126 self
127 }
128
129 pub fn filter_devices(mut self, device_ids: Vec<String>) -> Self {
131 self.device_filter = device_ids;
132 self
133 }
134}
135
136struct CachedReading {
138 data: AdvertisementData,
139 received_at: std::time::Instant,
140}
141
142pub struct PassiveMonitor {
150 options: PassiveMonitorOptions,
151 sender: broadcast::Sender<PassiveReading>,
153 cache: Arc<RwLock<HashMap<String, CachedReading>>>,
155}
156
157impl PassiveMonitor {
158 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 pub fn subscribe(&self) -> broadcast::Receiver<PassiveReading> {
172 self.sender.subscribe()
173 }
174
175 pub fn subscriber_count(&self) -> usize {
177 self.sender.receiver_count()
178 }
179
180 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 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 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 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 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 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 async fn scan_cycle_with_adapter(&self, adapter: &btleplug::platform::Adapter) -> Result<()> {
268 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 let peripherals = adapter.peripherals().await?;
281
282 for peripheral in peripherals {
283 if let Ok(Some(props)) = peripheral.properties().await {
284 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 if !self.options.device_filter.is_empty()
293 && !self.options.device_filter.contains(&device_id)
294 {
295 continue;
296 }
297
298 match parse_advertisement_with_name(data, props.local_name.as_deref()) {
300 Ok(adv_data) => {
301 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 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 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 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 if cached.received_at.elapsed() > self.options.max_reading_age {
348 return true;
349 }
350
351 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 if cached.data.counter != data.counter {
365 return true;
366 }
367
368 false
369 } else {
370 true
372 }
373 }
374
375 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 pub async fn known_devices(&self) -> Vec<String> {
383 let cache = self.cache.read().await;
384 cache.keys().cloned().collect()
385 }
386
387 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
399async 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 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 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 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 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 let mut changed = data.clone();
516 changed.co2 = Some(900);
517 assert!(monitor.should_emit("device-1", &changed).await);
518
519 let mut changed = data.clone();
521 changed.battery = 50;
522 assert!(monitor.should_emit("device-1", &changed).await);
523
524 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 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 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 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 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 assert!(monitor.should_emit("device-2", &data).await);
587 }
588
589 const TEST_LIMIT: Duration = Duration::from_secs(600);
593
594 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 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 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 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}