1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
use crate::core::persistence::{select_last_reading, upsert_reading};
use crate::prelude::*;
pub fn spawn(db: Arc<Mutex<Connection>>, bus: &mut Bus) -> Result<()> {
info!("Spawning readings persistence…");
let tx = bus.add_tx();
let rx = bus.add_rx();
crate::core::supervisor::spawn("my-iot::persistence", tx.clone(), move || {
for message in &rx {
if let Err(error) = process_message(&message, &db, &tx) {
error!("{}: {:?}", error, &message);
}
}
unreachable!();
})?;
Ok(())
}
fn process_message(message: &Message, db: &Arc<Mutex<Connection>>, tx: &Sender<Message>) -> Result<()> {
info!(
"{}: {:?} {:?}",
&message.sensor.sensor_id, &message.type_, &message.reading.value
);
debug!("{:?}", &message);
if message.type_ == MessageType::ReadLogged {
let db = db.lock().unwrap();
let previous_reading = select_last_reading(&db, &message.sensor.sensor_id)?;
upsert_reading(&db, &message.sensor, &message.reading)?;
send_messages(&previous_reading, &message, &tx)?;
}
Ok(())
}
fn send_messages(previous_reading: &Option<Reading>, message: &Message, tx: &Sender<Message>) -> Result<()> {
if let Some(existing) = previous_reading {
if message.reading.timestamp > existing.timestamp {
tx.send(
Composer::new(format!("{}::update", &message.sensor.sensor_id))
.type_(MessageType::ReadNonLogged)
.value(message.reading.value.clone())
.into(),
)?;
if message.reading.value != existing.value {
tx.send(
Composer::new(format!("{}::change", &message.sensor.sensor_id))
.type_(MessageType::ReadNonLogged)
.value(message.reading.value.clone())
.into(),
)?;
}
}
}
Ok(())
}