Move lights and info tasks to module

This commit is contained in:
2022-09-30 21:17:02 +02:00
parent d2054caf69
commit 62887371c9
5 changed files with 151 additions and 134 deletions

69
backend/src/tasks/info.rs Normal file
View File

@ -0,0 +1,69 @@
use std::time::Duration;
use common::ServerMessage;
use tokio::time::sleep;
use crate::{
collector::{Collector, MarkdownWeb, WeatherApi},
State,
};
pub async fn info_task(state: &State) {
let mut collectors: Vec<Box<dyn Collector + Send>> = vec![];
for url in &state.config.collectors.markdown_web_links {
collectors.push(Box::new(MarkdownWeb {
url: url.to_string(),
}));
}
if !state.config.collectors.weatherapi_locations.is_empty() {
let api_key = state
.config
.collectors
.weatherapi_key
.as_deref()
.expect("Missing weatherapi_key");
for location in state.config.collectors.weatherapi_locations.iter().cloned() {
collectors.push(Box::new(WeatherApi {
api_key: api_key.to_string(),
location,
}));
}
}
let mut collectors = collectors.into_boxed_slice();
let server_message = &state.server_message;
let collectors_len = collectors.len();
let next = move |i: usize| (i + 1) % collectors_len;
let mut i = 0;
loop {
sleep(Duration::from_secs(30)).await;
// don't bother collecting if no clients are connected
// there is always 1 receiver held by main process
if server_message.receiver_count() <= 1 {
continue;
}
i = next(i);
let collector = &mut collectors[i];
let msg = match collector.collect().await {
Ok(html) => ServerMessage::InfoPage { html },
Err(e) => {
warn!("collector error: {e}");
continue;
}
};
if let Err(e) = server_message.send(msg) {
error!("broadcast channel error: {e}");
return;
}
}
}

View File

@ -0,0 +1,63 @@
use common::{ClientMessage, ServerMessage};
use lighter_manager::manager::{BulbCommand, BulbManager, BulbSelector};
use tokio::select;
use tokio::time::{sleep, Duration};
use crate::State;
pub async fn lights_task(state: &State) {
let config = &state.config;
let server_message = &state.server_message;
let mut client_message = state.client_message.subscribe();
let (cmd, bulb_states) = BulbManager::launch(config.bulbs.clone(), config.mqtt.clone())
.await
.expect("Failed to launch bulb manager");
loop {
let notify = bulb_states.notify_on_change();
sleep(Duration::from_millis(1000 / 10)).await; // limit to 10 updates/second
select! {
_ = notify => {
for (id, mode) in bulb_states.bulbs().await.clone().into_iter() {
if let Err(e) = server_message.send(ServerMessage::BulbMode { id, mode }) {
error!("broadcast channel error: {e}");
return;
}
}
}
request = client_message.recv() => {
let request = match request {
Ok(r) => r,
Err(_) => continue,
};
match request.message {
ClientMessage::SetBulbColor { id, color } => {
if let Err(e) = cmd.send(BulbCommand::SetColor(BulbSelector::Id(id), color)).await {
error!("bulb manager error: {e}");
}
}
ClientMessage::SetBulbPower { id, power } => {
if let Err(e) = cmd.send(BulbCommand::SetPower(BulbSelector::Id(id), power)).await {
error!("bulb manager error: {e}");
}
}
ClientMessage::GetBulbs => {
if let Err(e) = request.response.send(ServerMessage::BulbMap(config.bulb_map.clone())).await {
error!("GetBulbs response channel error: {e}");
return;
}
for (id, mode) in bulb_states.bulbs().await.clone().into_iter() {
if let Err(e) = request.response.send(ServerMessage::BulbMode { id, mode }).await {
error!("GetBulbs response channel error: {e}");
return;
}
}
}
_ => {}
}
}
}
}
}

5
backend/src/tasks/mod.rs Normal file
View File

@ -0,0 +1,5 @@
pub mod info;
pub mod lights;
pub use info::info_task;
pub use lights::lights_task;