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
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
use crate::db::Db;
use crate::reading::Reading;
use crate::services::Service;
use crate::{threading, Result};
use chrono::Local;
use crossbeam_channel::{Receiver, Sender};
use log::{debug, info};
use serde::Deserialize;
use std::sync::{Arc, Mutex};
#[derive(Deserialize, Debug, Clone)]
pub struct Settings {
scenarios: Vec<Scenario>,
}
#[derive(Deserialize, Debug, Clone)]
pub struct Scenario {
#[serde(default = "String::new")]
description: String,
conditions: Vec<Condition>,
actions: Vec<Action>,
}
#[derive(Deserialize, Debug, Clone)]
pub enum Condition {
Sensor(String),
}
#[derive(Deserialize, Debug, Clone)]
pub enum Action {
Reading(),
}
pub struct Automator {
service_id: String,
settings: Settings,
}
impl Automator {
pub fn new(service_id: &str, settings: &Settings) -> Automator {
Automator {
service_id: service_id.into(),
settings: settings.clone(),
}
}
}
impl Service for Automator {
fn spawn(self: Box<Self>, _db: Arc<Mutex<Db>>, tx: Sender<Reading>, rx: Receiver<Reading>) -> Result<()> {
threading::spawn(self.service_id.clone(), move || loop {
for reading in rx.iter() {
for scenario in self.settings.scenarios.iter() {
if scenario.conditions.iter().all(|s| s.is_met(&reading)) {
info!(r#"Running scenario: "{}"."#, scenario.description);
for action in scenario.actions.iter() {
action.execute(&self.service_id, &reading, &tx).unwrap();
}
} else {
debug!(r#"Conditions are not met for scenario: "{}"."#, scenario.description)
}
}
}
})?;
Ok(())
}
}
impl Condition {
pub fn is_met(&self, reading: &Reading) -> bool {
match self {
Condition::Sensor(sensor) => &reading.sensor == sensor,
}
}
}
impl Action {
pub fn execute(&self, service_id: &str, reading: &Reading, tx: &Sender<Reading>) -> Result<()> {
match self {
Action::Reading() => tx
.send(Reading {
sensor: format!("{}::{}", &service_id, &reading.sensor),
timestamp: Local::now(),
value: reading.value.clone(),
is_persisted: false,
})
.map_err(|e| e.into()),
}
}
}