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
//! Readings receiver that actually processes all readings coming from services.

use crate::db::*;
use crate::message::*;
use crate::threading;
use crate::value::Value;
use crate::Result;
use bus::Bus;
use crossbeam_channel::Sender;
use log::{debug, info};
use std::sync::{Arc, Mutex};

/// Start readings receiver thread.
pub fn spawn(bus: &mut Bus<Message>, db: Arc<Mutex<Db>>, tx: &Sender<Message>) -> Result<()> {
    info!("Spawning message receiver…");
    let rx = bus.add_rx();
    let tx = tx.clone();

    threading::spawn("my-iot::receiver", move || {
        for message in rx {
            process_message(message, &db, &tx).unwrap();
        }
        unreachable!();
    })?;
    Ok(())
}

/// Process broadcasted message.
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(())
}

/// Check if sensor value has been updated or changed and send corresponding messages.
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(())
}