Compare commits

...

5 Commits

13 changed files with 96 additions and 49 deletions

View File

@@ -41,9 +41,10 @@ impl From<MapKey> for Arbitrary {
#[derive(Debug, Snafu)]
pub enum MapKeyFromArbitraryError {
#[snafu(display("floats aren't supported as map keys yet. got {value:?}"))]
/// floats aren't supported as map keys yet. got {value:?}
FloatNotSupported { value: FiniteF64 },
#[snafu(display("a map cannot be a map key. got {value:?}"))]
/// a map cannot be a map key. got {value:?}
MapCannotBeAMapKey { value: Map },
}

View File

@@ -4,8 +4,8 @@ use snafu::Snafu;
#[derive(Debug, Clone, derive_more::Into)]
pub struct FiniteF64(f64);
/// {value:?} is not finite
#[derive(Debug, Snafu)]
#[snafu(display("{value:?} is not finite"))]
pub struct NotFinite {
value: f64,
}

View File

@@ -1,3 +1,3 @@
pub mod connection;
mod impl_protocol;
pub mod messages;
mod protocol;

View File

@@ -1,9 +1,9 @@
use snafu::Snafu;
pub mod emitter;
pub mod emitter_ext;
mod emitter_ext;
pub mod signal;
pub mod signal_ext;
mod signal_ext;
pub use emitter::Emitter;
pub use emitter_ext::EmitterExt;

View File

@@ -1,30 +1,24 @@
use std::future::Future;
use ext_trait::extension;
use snafu::{ResultExt, Snafu};
use tokio::select;
use crate::ProducerExited;
use super::signal::{JoinError, Signal};
#[derive(Debug, Snafu)]
pub struct ProducerAlreadyExited {
source: ProducerExited,
}
#[extension(pub trait SignalExt)]
impl<T> Signal<T> {
fn map<M, F>(
self,
mut func: F,
) -> Result<(Signal<M>, impl Future<Output = Result<(), JoinError>>), ProducerAlreadyExited>
) -> Result<(Signal<M>, impl Future<Output = Result<(), JoinError>>), ProducerExited>
where
T: 'static + Sync + Send + Clone,
M: 'static + Sync + Send + Clone,
F: 'static + Send + FnMut(T) -> M,
{
let initial = func(self.subscribe().context(ProducerAlreadyExitedSnafu)?.get());
let initial = func(self.subscribe()?.get());
Ok(Signal::new(initial, |mut publisher_stream| async move {
while let Some(publisher) = publisher_stream.wait().await {

View File

@@ -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::<f64>(
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::<f64>(
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

View File

@@ -12,13 +12,13 @@ pub struct EntityId(pub Domain, pub ObjectId);
#[derive(Debug, Clone, Snafu)]
pub enum EntityIdParsingError {
#[snafu(display("entity IDs have a dot / period in them, e.g. light.kitchen_lamp"))]
/// entity IDs have a dot / period in them, e.g. light.kitchen_lamp
MissingDot,
#[snafu(display("could not parse the domain part of the entity ID"))]
/// could not parse the domain part of the entity ID
ParsingDomain { source: <Domain as FromStr>::Err },
#[snafu(display("could not parse the object ID part of the entity ID"))]
/// could not parse the object ID part of the entity ID
ParsingObjectId { source: ObjectIdParsingError },
}

View File

@@ -4,12 +4,12 @@ use super::InputNumberMode;
#[derive(Debug, FromPyObject)]
#[pyo3(from_item_all)]
pub struct InputNumberAttributes<Numeric> {
initial: Option<Numeric>,
pub struct InputNumberAttributes<Number> {
initial: Option<Number>,
editable: bool,
min: Numeric,
max: Numeric,
step: Numeric,
min: Number,
max: Number,
step: Number,
mode: InputNumberMode,
// todo: CustomUnitOfMeasurement type? probably not?
unit_of_measurement: Option<String>,

View File

@@ -19,7 +19,7 @@ mod mode;
pub use attributes::InputNumberAttributes;
pub use mode::InputNumberMode;
#[derive(Debug, Snafu)]
#[derive(Debug, Clone, Snafu)]
pub enum CreateSignalError {
/// couldn't get the underlying state object signal
StateObjectSignalError {
@@ -28,22 +28,20 @@ pub enum CreateSignalError {
/// couldn't map the state object to a power value
MappedSignalError {
source: emitter_and_signal::signal_ext::ProducerAlreadyExited,
source: emitter_and_signal::ProducerExited,
},
}
pub fn signal<
'py,
Numeric: 'static + uom::num::Num + Clone + Send + Sync + FromStr + for<'a, 'py2> FromPyObject<'a, 'py2>,
Number: 'static + uom::num::Num + Clone + Send + Sync + FromStr + for<'a, 'py2> FromPyObject<'a, 'py2>,
>(
py: Python<'py>,
home_assistant: &'py HomeAssistant,
object_id: ObjectId,
) -> Result<
(
Signal<
Option<Arc<Result<HomeAssistantState<Numeric>, StateObjectSignalError<Arc<PyErr>>>>>,
>,
Signal<Option<Arc<Result<HomeAssistantState<Number>, StateObjectSignalError<Arc<PyErr>>>>>>,
impl Future<Output = Result<(), emitter_and_signal::signal::JoinError>>,
),
CreateSignalError,
@@ -51,8 +49,8 @@ pub fn signal<
let entity_id = EntityId(Domain::InputNumber, object_id);
let (signal, task1) = StateObject::<
HomeAssistantState<Numeric>,
InputNumberAttributes<Numeric>,
HomeAssistantState<Number>,
InputNumberAttributes<Number>,
Py<PyAny>,
>::signal(py, home_assistant, entity_id)
.context(StateObjectSignalSnafu)?;

View File

@@ -47,11 +47,11 @@ pub enum CreateSignalError {
/// couldn't map the state object to a power value
MappedSignalError {
source: emitter_and_signal::signal_ext::ProducerAlreadyExited,
source: emitter_and_signal::ProducerExited,
},
}
pub fn signal<'py, V>(
pub fn signal<'py, Number>(
py: Python<'py>,
home_assistant: &'py HomeAssistant,
object_id: ObjectId,
@@ -60,7 +60,7 @@ pub fn signal<'py, V>(
Signal<
Option<
Result<
HomeAssistantState<uom::si::quantities::Power<V>>,
HomeAssistantState<uom::si::quantities::Power<Number>>,
StateObjectSignalError<Arc<PyErr>>,
>,
>,
@@ -70,22 +70,22 @@ pub fn signal<'py, V>(
CreateSignalError,
>
where
V: 'static + uom::num::Num + Clone + Send + Sync + FromStr,
V: uom::num::Num + uom::Conversion<V, T = V>,
milliwatt: Conversion<V, T = V>,
watt: Conversion<V, T = V>,
kilowatt: Conversion<V, T = V>,
megawatt: Conversion<V, T = V>,
gigawatt: Conversion<V, T = V>,
terawatt: Conversion<V, T = V>,
btu: Conversion<V, T = V>,
hour: Conversion<V, T = V>,
SI<V>: Units<V>,
Number: 'static + uom::num::Num + Clone + Send + Sync + FromStr,
Number: uom::num::Num + uom::Conversion<Number, T = Number>,
milliwatt: Conversion<Number, T = Number>,
watt: Conversion<Number, T = Number>,
kilowatt: Conversion<Number, T = Number>,
megawatt: Conversion<Number, T = Number>,
gigawatt: Conversion<Number, T = Number>,
terawatt: Conversion<Number, T = Number>,
btu: Conversion<Number, T = Number>,
hour: Conversion<Number, T = Number>,
SI<Number>: Units<Number>,
{
let entity_id = EntityId(Domain::Sensor, object_id);
let (signal, task1) =
StateObject::<HomeAssistantState<V>, PowerSensorAttributes, Py<PyAny>>::signal(
StateObject::<HomeAssistantState<Number>, PowerSensorAttributes, Py<PyAny>>::signal(
py,
home_assistant,
entity_id,

View File

@@ -7,8 +7,8 @@ use snafu::Snafu;
#[derive(Debug, Clone, derive_more::Display, FromPyFromStr, ToStrToPy)]
pub struct Slug(Arc<str>);
/// expected a lowercase ASCII alphabetical character (i.e. a through z) or a digit (i.e. 0 through 9) or an underscore (i.e. _) but encountered {encountered}
#[derive(Debug, Clone, Snafu)]
#[snafu(display("expected a lowercase ASCII alphabetical character (i.e. a through z) or a digit (i.e. 0 through 9) or an underscore (i.e. _) but encountered {encountered}"))]
pub struct SlugParsingError {
encountered: char,
}

View File

@@ -28,7 +28,7 @@ pub struct StateObject<State, Attributes, ContextEvent> {
pub type ExtractStateObjectError<'a, 'py, State, Attributes, ContextEvent> =
<StateObject<State, Attributes, ContextEvent> as FromPyObject<'a, 'py>>::Error;
#[derive(Debug, Snafu)]
#[derive(Debug, Clone, Snafu)]
pub enum CreateSignalError {
/// couldn't get the state machine from the Home Assistant object
GetStatesError { source: GetStatesError },