Compare commits

...

7 Commits

14 changed files with 125 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 {
initial: Option<f64>,
pub struct InputNumberAttributes<Number> {
initial: Option<Number>,
editable: bool,
min: f64,
max: f64,
step: f64,
min: Number,
max: Number,
step: Number,
mode: InputNumberMode,
// todo: CustomUnitOfMeasurement type? probably not?
unit_of_measurement: Option<String>,

View File

@@ -1,7 +1,7 @@
use std::{future::Future, sync::Arc};
use std::{future::Future, str::FromStr, sync::Arc};
use emitter_and_signal::{Signal, SignalExt};
use pyo3::{Py, PyAny, PyErr, Python};
use pyo3::{FromPyObject, Py, PyAny, PyErr, Python};
use snafu::{ResultExt, Snafu};
use crate::{
@@ -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,30 +28,32 @@ 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>(
pub fn signal<
'py,
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<f64>, StateObjectSignalError<Arc<PyErr>>>>>>,
Signal<Option<Arc<Result<HomeAssistantState<Number>, StateObjectSignalError<Arc<PyErr>>>>>>,
impl Future<Output = Result<(), emitter_and_signal::signal::JoinError>>,
),
CreateSignalError,
> {
let entity_id = EntityId(Domain::InputNumber, object_id);
let (signal, task1) =
StateObject::<HomeAssistantState<f64>, InputNumberAttributes, Py<PyAny>>::signal(
py,
home_assistant,
entity_id,
)
.context(StateObjectSignalSnafu)?;
let (signal, task1) = StateObject::<
HomeAssistantState<Number>,
InputNumberAttributes<Number>,
Py<PyAny>,
>::signal(py, home_assistant, entity_id)
.context(StateObjectSignalSnafu)?;
let (signal, task2) = signal
.map(|state_object_arc_result_option| {

View File

@@ -1,10 +1,19 @@
use std::{future::Future, sync::Arc};
use std::{future::Future, str::FromStr, sync::Arc};
use emitter_and_signal::{Signal, SignalExt};
use pyo3::{FromPyObject, Py, PyAny, PyErr, Python};
use python_utils::{FromPyFromStr, ToStrToPy};
use snafu::{ResultExt, Snafu};
use string_literal::StringLiteral;
use uom::{
si::{
energy::btu,
power::{gigawatt, kilowatt, megawatt, milliwatt, terawatt, watt},
time::hour,
Units, SI,
},
Conversion,
};
use super::super::state_classes::measurement::Measurement;
use crate::{
@@ -38,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>(
pub fn signal<'py, Number>(
py: Python<'py>,
home_assistant: &'py HomeAssistant,
object_id: ObjectId,
@@ -50,17 +59,33 @@ pub fn signal<'py>(
(
Signal<
Option<
Result<HomeAssistantState<uom::si::f64::Power>, StateObjectSignalError<Arc<PyErr>>>,
Result<
HomeAssistantState<uom::si::quantities::Power<Number>>,
StateObjectSignalError<Arc<PyErr>>,
>,
>,
>,
impl Future<Output = Result<(), emitter_and_signal::signal::JoinError>>,
),
CreateSignalError,
> {
>
where
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<f64>, PowerSensorAttributes, Py<PyAny>>::signal(
StateObject::<HomeAssistantState<Number>, PowerSensorAttributes, Py<PyAny>>::signal(
py,
home_assistant,
entity_id,
@@ -75,9 +100,9 @@ pub fn signal<'py>(
|StateObject {
state, attributes, ..
}| {
state
.as_ref()
.map(|&amount| attributes.unit_of_measurement.into_uom(amount))
state.as_ref().map(|amount| {
attributes.unit_of_measurement.into_uom(amount.clone())
})
},
)
.map_err(Clone::clone)

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 },

View File

@@ -34,7 +34,7 @@ pub enum UnitOfMeasurement {
impl UnitOfMeasurement {
pub fn into_uom<V>(&self, amount: V) -> Power<V>
where
V: uom::num::Num + uom::Conversion<V, T = V>,
V: uom::num::Num + Conversion<V, T = V>,
milliwatt: Conversion<V, T = V>,
watt: Conversion<V, T = V>,
kilowatt: Conversion<V, T = V>,