diff --git a/entrypoint/src/lib.rs b/entrypoint/src/lib.rs index 6d4c541..41d4022 100644 --- a/entrypoint/src/lib.rs +++ b/entrypoint/src/lib.rs @@ -2,7 +2,7 @@ use std::{path::PathBuf, str::FromStr, time::Duration}; use clap::Parser; use driver_kasa::connection::LB130USHandle; -use futures::stream::{FuturesUnordered, StreamExt}; +use futures::{TryFutureExt, stream::{FuturesUnordered, StreamExt}}; use home_assistant::{home_assistant::HomeAssistant, object_id::ObjectId}; use protocol::light::{Kelvin, TurnToTemperature}; use pyo3::{ @@ -11,7 +11,9 @@ use pyo3::{ wrap_pyfunction, Bound, PyAny, PyResult, Python, }; use shadow_rs::shadow; -use tokio::time::{interval, MissedTickBehavior}; +use tokio::{ + join, task::JoinSet, time::{MissedTickBehavior, interval} +}; use tracing::{level_filters::LevelFilter, Level}; use tracing_subscriber::{ fmt::{self, format::FmtSpan}, @@ -187,6 +189,58 @@ async fn real_main( // }, context, target, false).await; // dbg!(total_silence_result); // } + }, + async { + let plug_awp04l_1_power_object_id = "plug_awp04l_1_power"; + let plug_awp04l_2_power_object_id = "plug_awp04l_2_power"; + + let plug_awp04l_1_power_object_id = ObjectId::from_str(plug_awp04l_1_power_object_id) + .expect("statically written and known to be correct"); + let plug_awp04l_2_power_object_id = ObjectId::from_str(plug_awp04l_2_power_object_id) + .expect("statically written and known to be correct"); + + let plug_awp04l_1_power_signal_result = Python::attach(|py| { + home_assistant::sensor::device_classes::power::signal::( + py, + &home_assistant, + plug_awp04l_1_power_object_id, + ) + }); + let plug_awp04l_2_power_signal_result = Python::attach(|py| { + home_assistant::sensor::device_classes::power::signal::( + py, + &home_assistant, + plug_awp04l_2_power_object_id, + ) + }); + + let mut tasks = JoinSet::new(); + if let Ok((plug_awp04l_1_power_signal, task)) = plug_awp04l_1_power_signal_result { + tracing::error!("listening to the plug_awp04l_1_power_signal"); + tasks.spawn(task.unwrap_or_else(|e| panic!("TODO: {e}"))); + tasks.spawn( + plug_awp04l_1_power_signal + .subscribe() + .expect("TODO") + .for_each(|plug_awp04l_1_power| async move { + tracing::warn!(?plug_awp04l_1_power) + }), + ); + } + if let Ok((plug_awp04l_2_power_signal, task)) = plug_awp04l_2_power_signal_result { + tracing::error!("listening to the plug_awp04l_2_power_signal"); + tasks.spawn(task.unwrap_or_else(|e| panic!("TODO: {e}"))); + tasks.spawn( + plug_awp04l_2_power_signal + .subscribe() + .expect("TODO") + .for_each(|plug_awp04l_2_power| async move { + tracing::warn!(?plug_awp04l_2_power) + }), + ); + } + + tasks.join_all().await; } ) .0