Stuff
This commit is contained in:
442
Cargo.lock
generated
442
Cargo.lock
generated
File diff suppressed because it is too large
Load Diff
@ -25,8 +25,7 @@ path = "../common"
|
||||
|
||||
[dependencies.lighter_lib]
|
||||
git = "https://git.nubo.sh/hulthe/lighter.git"
|
||||
#path = "../../lighter/lib"
|
||||
|
||||
[dependencies.lighter_manager]
|
||||
git = "https://git.nubo.sh/hulthe/lighter.git"
|
||||
#path = "../../lighter/manager"
|
||||
|
||||
|
||||
@ -1,3 +1,5 @@
|
||||
persistence_dir = "/tmp/"
|
||||
|
||||
[mqtt]
|
||||
#address = "hostname"
|
||||
#port = 1883
|
||||
|
||||
@ -1,4 +1,5 @@
|
||||
mod collector;
|
||||
mod persistence;
|
||||
mod tasks;
|
||||
|
||||
use clap::Parser;
|
||||
@ -7,6 +8,7 @@ use common::{BulbMap, ClientMessage, ServerMessage};
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
use lighter_manager::{manager::BulbsConfig, mqtt_conf::MqttConfig};
|
||||
use log::LevelFilter;
|
||||
use persistence::Persistence;
|
||||
use serde::Deserialize;
|
||||
use std::convert::Infallible;
|
||||
use std::net::SocketAddr;
|
||||
@ -46,6 +48,8 @@ pub struct Config {
|
||||
|
||||
collectors: CollectorConfig,
|
||||
|
||||
persistence_dir: Option<PathBuf>,
|
||||
|
||||
#[serde(flatten)]
|
||||
bulbs: BulbsConfig,
|
||||
|
||||
@ -55,6 +59,7 @@ pub struct Config {
|
||||
|
||||
pub struct State {
|
||||
config: Config,
|
||||
persistence: Persistence,
|
||||
client_message: broadcast::Sender<ClientRequest>,
|
||||
server_message: broadcast::Sender<ServerMessage>,
|
||||
}
|
||||
@ -85,9 +90,15 @@ async fn main() {
|
||||
let (client_message, _) = broadcast::channel(100);
|
||||
|
||||
let state = State {
|
||||
config,
|
||||
client_message,
|
||||
server_message,
|
||||
persistence: match &config.persistence_dir {
|
||||
Some(path) => Persistence::new_persistence(path.to_owned())
|
||||
.await
|
||||
.expect("Failed to open persistence dir"),
|
||||
None => Persistence::new().await,
|
||||
},
|
||||
config,
|
||||
};
|
||||
let state = Box::leak(Box::new(state));
|
||||
|
||||
|
||||
111
backend/src/persistence.rs
Normal file
111
backend/src/persistence.rs
Normal file
@ -0,0 +1,111 @@
|
||||
use std::io;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use serde::{de::DeserializeOwned, Serialize};
|
||||
use tokio::{
|
||||
fs::{self, File},
|
||||
io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt},
|
||||
};
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Persistence {
|
||||
state: Arc<PersistenceState>,
|
||||
}
|
||||
|
||||
pub struct PersistenceFile<S> {
|
||||
content: S,
|
||||
file: Option<File>,
|
||||
}
|
||||
|
||||
enum PersistenceState {
|
||||
NoPersistence,
|
||||
Persistence { directory: PathBuf },
|
||||
}
|
||||
|
||||
impl Persistence {
|
||||
pub async fn new() -> Persistence {
|
||||
Persistence {
|
||||
state: Arc::new(PersistenceState::NoPersistence),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn new_persistence(directory: PathBuf) -> io::Result<Persistence> {
|
||||
let _ = fs::read_dir(&directory).await?;
|
||||
|
||||
Ok(Persistence {
|
||||
state: Arc::new(PersistenceState::Persistence { directory }),
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn open<S: Default + Serialize + DeserializeOwned>(
|
||||
&self,
|
||||
name: String,
|
||||
) -> Result<PersistenceFile<S>, io::Error> {
|
||||
Ok(match &*self.state {
|
||||
PersistenceState::NoPersistence => PersistenceFile {
|
||||
content: S::default(),
|
||||
file: None,
|
||||
},
|
||||
PersistenceState::Persistence { directory } => {
|
||||
let file_path = directory.join(&name);
|
||||
let mut file = fs::OpenOptions::new()
|
||||
.create(true)
|
||||
.read(true)
|
||||
.write(true)
|
||||
.open(file_path)
|
||||
.await?;
|
||||
|
||||
let mut content = String::new();
|
||||
file.read_to_string(&mut content).await?;
|
||||
|
||||
let content = ron::from_str(&content).unwrap_or_else(|_| {
|
||||
info!("Failed to load persistence {name}, using default");
|
||||
S::default()
|
||||
});
|
||||
|
||||
PersistenceFile {
|
||||
content,
|
||||
file: Some(file),
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl<S> PersistenceFile<S>
|
||||
where
|
||||
S: Serialize,
|
||||
{
|
||||
pub fn get(&self) -> &S {
|
||||
&self.content
|
||||
}
|
||||
|
||||
pub async fn set(&mut self, s: S) -> anyhow::Result<()> {
|
||||
self.content = s;
|
||||
|
||||
if let Some(file) = &mut self.file {
|
||||
let serialized = ron::ser::to_string_pretty(&self.content, Default::default())?;
|
||||
file.rewind().await?;
|
||||
file.set_len(0).await?;
|
||||
file.write_all(serialized.as_bytes()).await?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn update(&mut self, f: impl FnOnce(&mut S)) -> anyhow::Result<()>
|
||||
where
|
||||
S: Clone + PartialEq,
|
||||
{
|
||||
let mut new = self.content.clone();
|
||||
|
||||
f(&mut new);
|
||||
|
||||
if new != self.content {
|
||||
self.set(new).await?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@ -1,14 +1,49 @@
|
||||
use common::{ClientMessage, ServerMessage};
|
||||
use lighter_manager::manager::{BulbCommand, BulbManager, BulbSelector};
|
||||
use tokio::select;
|
||||
use tokio::time::{sleep, Duration};
|
||||
use std::collections::HashMap;
|
||||
|
||||
use crate::State;
|
||||
use chrono::{Datelike, Local, NaiveTime, Weekday};
|
||||
use common::{ClientMessage, ServerMessage};
|
||||
use lighter_lib::{BulbColor, BulbId};
|
||||
use lighter_manager::manager::{BulbCommand, BulbManager, BulbSelector};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::time::Duration;
|
||||
use tokio::select;
|
||||
use tokio::sync::{broadcast, mpsc};
|
||||
use tokio::task::{spawn, JoinHandle};
|
||||
use tokio::time::sleep;
|
||||
|
||||
use crate::persistence::PersistenceFile;
|
||||
use crate::{ClientRequest, State};
|
||||
|
||||
#[derive(Default, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
struct LightsState {
|
||||
wake_schedule: HashMap<(BulbId, Weekday), NaiveTime>,
|
||||
}
|
||||
|
||||
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 mut lights_state: PersistenceFile<LightsState> = state
|
||||
.persistence
|
||||
.open("lights".into())
|
||||
.await
|
||||
.expect("Failed to open lights config");
|
||||
|
||||
let mut wake_tasks: HashMap<(BulbId, Weekday), JoinHandle<()>> = lights_state
|
||||
.get()
|
||||
.wake_schedule
|
||||
.iter()
|
||||
.map(|((bulb, day), time)| {
|
||||
let handle = spawn(wake_task(
|
||||
state.client_message.clone(),
|
||||
bulb.clone(),
|
||||
*day,
|
||||
*time,
|
||||
));
|
||||
|
||||
((bulb.clone(), *day), handle)
|
||||
})
|
||||
.collect();
|
||||
|
||||
let (cmd, bulb_states) = BulbManager::launch(config.bulbs.clone(), config.mqtt.clone())
|
||||
.await
|
||||
@ -16,7 +51,7 @@ pub async fn lights_task(state: &State) {
|
||||
|
||||
loop {
|
||||
let notify = bulb_states.notify_on_change();
|
||||
sleep(Duration::from_millis(1000 / 10)).await; // limit to 10 updates/second
|
||||
sleep(tokio::time::Duration::from_millis(1000 / 10)).await; // limit to 10 updates/second
|
||||
select! {
|
||||
_ = notify => {
|
||||
for (id, mode) in bulb_states.bulbs().await.clone().into_iter() {
|
||||
@ -55,9 +90,70 @@ pub async fn lights_task(state: &State) {
|
||||
}
|
||||
}
|
||||
}
|
||||
ClientMessage::SetBulbWakeTime { id, day, time } => {
|
||||
if let Err(e) = lights_state.update(|lights_state| {
|
||||
lights_state.wake_schedule.insert((id.clone(), day), time);
|
||||
}).await {
|
||||
error!("Failed to save wake schedule: {e}");
|
||||
};
|
||||
|
||||
let handle = spawn(wake_task(state.client_message.clone(), id.clone(), day, time));
|
||||
if let Some(old_handle) = wake_tasks.insert((id, day), handle) {
|
||||
old_handle.abort();
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn wake_task(
|
||||
channel: broadcast::Sender<ClientRequest>,
|
||||
id: BulbId,
|
||||
day: Weekday,
|
||||
time: NaiveTime,
|
||||
) {
|
||||
let now = Local::now();
|
||||
let day_num = day.num_days_from_monday();
|
||||
let now_day = now.weekday();
|
||||
let now_day_num = now_day.num_days_from_monday();
|
||||
|
||||
let mut alarm = now;
|
||||
if day_num >= now_day_num {
|
||||
// next alarm is this week
|
||||
alarm += chrono::Duration::days((day_num - now_day_num).into());
|
||||
alarm = alarm.date().and_time(time).unwrap();
|
||||
} else {
|
||||
// next alarm is next week
|
||||
alarm += chrono::Duration::weeks(1);
|
||||
alarm -= chrono::Duration::days((now_day_num - day_num).into());
|
||||
alarm = alarm.date().and_time(time).unwrap();
|
||||
}
|
||||
|
||||
loop {
|
||||
info!("sleeping until {alarm}");
|
||||
sleep((alarm - Local::now()).to_std().unwrap()).await;
|
||||
alarm += chrono::Duration::weeks(1);
|
||||
|
||||
for brightness in (1..=50).map(|i| (i as f32) * 0.01) {
|
||||
sleep(Duration::from_secs(12)).await;
|
||||
|
||||
let message = ClientMessage::SetBulbColor {
|
||||
id: id.clone(),
|
||||
color: BulbColor::Kelvin {
|
||||
t: 0.0,
|
||||
b: brightness,
|
||||
},
|
||||
};
|
||||
|
||||
let (response, _) = mpsc::channel(1);
|
||||
let request = ClientRequest { message, response };
|
||||
|
||||
if channel.send(request).is_err() {
|
||||
return;
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user