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
61
62
63
64
65
66
use crate::db::*;
use crate::message::*;
use crate::threading;
use crate::value::Value;
use crate::Result;
use crossbeam_channel::Sender;
use log::{debug, info};
use std::sync::{Arc, Mutex};
pub fn spawn(db: Arc<Mutex<Db>>, tx: &Sender<Message>) -> Result<Sender<Message>> {
info!("Spawning readings persistence…");
let tx = tx.clone();
let (out_tx, rx) = crossbeam_channel::unbounded::<Message>();
threading::spawn("my-iot::persistence", move || {
for message in rx {
process_message(message, &db, &tx).unwrap();
}
unreachable!();
})?;
Ok(out_tx)
}
fn process_message(message: Message, db: &Arc<Mutex<Db>>, tx: &Sender<Message>) -> Result<()> {
info!(
"{}: {:?} {:?}",
&message.reading.sensor, &message.type_, &message.reading.value,
);
debug!("{:?}", &message);
if message.type_ == Type::Actual {
let db = db.lock().unwrap();
let previous_reading = db.select_last_reading(&message.reading.sensor)?;
db.insert_reading(&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(Message::now(
Type::OneOff,
format!("{}::update", &message.reading.sensor),
Value::Update(
Box::new(existing.value.clone()),
Box::new(message.reading.value.clone()),
),
))?;
if message.reading.value != existing.value {
tx.send(Message::now(
Type::OneOff,
format!("{}::change", &message.reading.sensor),
message.reading.value.clone(),
))?;
}
}
}
Ok(())
}