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
121
122
123
124
use crate::core::supervisor;
use crate::prelude::*;
use crate::settings::Settings;
use crossbeam_channel::{Receiver, Sender};
use log::Level;
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use structopt::StructOpt;
mod consts;
mod core;
mod format;
mod prelude;
mod services;
mod settings;
mod templates;
mod web;
#[derive(StructOpt, Debug)]
#[structopt(name = "my-iot", author, about)]
struct Opt {
#[structopt(short = "s", long = "silent")]
silent: bool,
#[structopt(short = "v", long = "verbose", conflicts_with = "silent")]
verbose: bool,
#[structopt(long, parse(from_os_str), env = "MYIOT_SETTINGS", default_value = "my-iot.toml")]
settings: PathBuf,
#[structopt(long, parse(from_os_str), env = "MYIOT_DB", default_value = "my-iot.sqlite3")]
db: PathBuf,
}
fn main() -> Result<()> {
let opt: Opt = Opt::from_args();
simple_logger::init_with_level(if opt.silent {
Level::Warn
} else if opt.verbose {
Level::Debug
} else {
Level::Info
})?;
info!("Reading settings…");
let settings = settings::read(opt.settings)?;
debug!("Settings: {:?}", &settings);
info!("Opening database…");
let db = Arc::new(Mutex::new(Db::new(opt.db)?));
info!("Starting services…");
let (dispatcher_tx, dispatcher_rx) = crossbeam_channel::unbounded();
dispatcher_tx.send(Composer::new("my-iot::start").type_(MessageType::ReadNonLogged).into())?;
let mut all_txs = vec![core::persistence::spawn(db.clone(), &dispatcher_tx)?];
all_txs.extend(spawn_services(&settings, &db, &dispatcher_tx)?);
spawn_dispatcher(dispatcher_rx, dispatcher_tx, all_txs)?;
info!("Starting web server on port {}…", settings.http_port);
web::start_server(settings, db.clone())
}
fn spawn_services(settings: &Settings, db: &Arc<Mutex<Db>>, tx: &Sender<Message>) -> Result<Vec<Sender<Message>>> {
let mut service_txs = Vec::new();
for (service_id, service_settings) in settings.services.iter() {
if !settings.disabled_services.contains(service_id.as_str()) {
info!("Spawning service `{}`…", service_id);
debug!("Settings `{}`: {:?}", service_id, service_settings);
let txs = services::spawn(service_id, service_settings, &db, tx)?;
debug!("Got {} txs from `{}`", txs.len(), service_id);
service_txs.extend(txs);
} else {
warn!("Service `{}` is disabled.", &service_id);
}
}
Ok(service_txs)
}
fn spawn_dispatcher(rx: Receiver<Message>, tx: Sender<Message>, txs: Vec<Sender<Message>>) -> Result<()> {
info!("Spawning message dispatcher…");
supervisor::spawn("my-iot::dispatcher", tx, move || -> Result<()> {
for message in &rx {
debug!("Dispatching {}", &message.sensor);
for tx in txs.iter() {
if let Err(error) = tx.send(message.clone()) {
error!("Could not send message to {:?}: {:?}", tx, error);
}
}
debug!("Dispatched {}", &message.sensor);
}
Err(format_err!("Receiver channel is unexpectedly exhausted"))
})?;
Ok(())
}