1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
use crate::db::*;
use crate::reading::*;
use crate::threading::spawn;
use crate::Result;
use log::info;
use multiqueue::BroadcastReceiver;
use std::sync::{Arc, Mutex};
pub fn start(rx: &BroadcastReceiver<Message>, db: Arc<Mutex<Db>>) -> Result<()> {
let rx = rx.add_stream().into_single().unwrap();
spawn(module_path!(), move || {
for message in rx {
info!("{}: {:?}", &message.reading.sensor, &message.reading.value);
if message.type_ == Type::Actual {
db.lock().unwrap().insert_reading(&message.reading).unwrap();
}
}
unreachable!();
})?;
Ok(())
}