use core::time::Duration;
use pamoja_codec::Codec;
use pamoja_core::{Actuator, Result, Sensor, Transport};
use pamoja_power::PowerMode;
use crate::{Controller, Profile, Reaction};
#[derive(Clone, Copy, Debug, Default)]
pub struct NoActuator;
impl Actuator for NoActuator {
type Command = bool;
async fn apply(&mut self, _command: bool) -> Result<()> {
Ok(())
}
}
pub struct Node<S, A, T, C> {
profile: Profile,
controller: Controller,
sensor: S,
actuator: A,
transport: T,
codec: C,
}
impl<S, A, T, C> Node<S, A, T, C> {
pub fn new(profile: Profile, sensor: S, actuator: A, transport: T, codec: C) -> Self {
let controller = profile.controller();
Self {
profile,
controller,
sensor,
actuator,
transport,
codec,
}
}
pub fn profile(&self) -> &Profile {
&self.profile
}
pub fn schedule(&self, soc: f32, charging: bool) -> (PowerMode, Duration) {
let plan = self.profile.power.plan();
let mode = plan.mode_while_charging(soc, charging);
(mode, plan.interval_for(mode))
}
}
impl<S, T, C> Node<S, NoActuator, T, C> {
pub fn monitor(profile: Profile, sensor: S, transport: T, codec: C) -> Self {
Node::new(profile, sensor, NoActuator, transport, codec)
}
}
impl<S, A, T, C> Node<S, A, T, C>
where
S: Sensor<Reading = f32>,
A: Actuator<Command = bool>,
T: Transport,
C: Codec<f32>,
{
pub async fn tick(&mut self) -> Result<Reaction> {
let reading = self.sensor.read().await?;
let reaction = self.controller.evaluate(reading);
if let Some(on) = reaction.actuator {
self.actuator.apply(on).await?;
}
let payload = self.codec.encode(&reading)?;
self.transport.send(&self.profile.topic, &payload).await?;
Ok(reaction)
}
}
#[cfg(test)]
mod tests {
use std::sync::{Arc, Mutex};
use pamoja_codec::CborCodec;
use pamoja_core::Error;
use pamoja_loopback::{LoopbackBroker, LoopbackTransport};
use pamoja_power::PowerMode;
use super::*;
struct ScriptedSensor {
readings: std::vec::IntoIter<f32>,
}
impl ScriptedSensor {
fn new(readings: Vec<f32>) -> Self {
Self {
readings: readings.into_iter(),
}
}
}
impl Sensor for ScriptedSensor {
type Reading = f32;
async fn read(&mut self) -> Result<f32> {
self.readings.next().ok_or(Error::Closed)
}
}
#[derive(Clone)]
struct Recording(Arc<Mutex<Vec<bool>>>);
impl Actuator for Recording {
type Command = bool;
async fn apply(&mut self, on: bool) -> Result<()> {
self.0.lock().expect("commands lock").push(on);
Ok(())
}
}
async fn connected_pair(filter: &str) -> (LoopbackTransport, LoopbackTransport) {
let broker = LoopbackBroker::new();
let mut gateway = LoopbackTransport::new(broker.clone());
let mut link = LoopbackTransport::new(broker);
gateway.connect().await.expect("gateway connect");
link.connect().await.expect("link connect");
gateway.subscribe(filter).await.expect("subscribe");
(gateway, link)
}
#[tokio::test]
async fn tick_reads_actuates_and_publishes() {
let (mut gateway, link) = connected_pair("cold-chain/#").await;
let commands = Arc::new(Mutex::new(Vec::new()));
let mut node = Node::new(
Profile::vaccine_fridge_monitor(),
ScriptedSensor::new(vec![9.0]),
Recording(commands.clone()),
link,
CborCodec,
);
let reaction = node.tick().await.expect("tick");
assert_eq!(reaction.actuator, Some(true));
assert!(reaction.alert.is_some());
assert_eq!(*commands.lock().expect("commands"), vec![true]);
let message = gateway.recv().await.expect("recv").expect("a reading");
assert_eq!(message.topic, "cold-chain/fridge/temperature");
let reading: f32 = CborCodec.decode(&message.payload).expect("decode");
assert_eq!(reading, 9.0);
}
#[tokio::test]
async fn monitor_publishes_without_an_actuator() {
let (mut gateway, link) = connected_pair("water/#").await;
let mut node = Node::monitor(
Profile::well_level(),
ScriptedSensor::new(vec![3.2]),
link,
CborCodec,
);
let reaction = node.tick().await.expect("tick");
assert_eq!(reaction.actuator, None);
let message = gateway.recv().await.expect("recv").expect("a reading");
assert_eq!(message.topic, "water/well/level");
let reading: f32 = CborCodec.decode(&message.payload).expect("decode");
assert_eq!(reading, 3.2);
}
#[test]
fn schedule_follows_state_of_charge() {
let node = Node::monitor(Profile::vaccine_fridge_monitor(), (), (), ());
assert_eq!(node.schedule(0.9, false).0, PowerMode::Active);
assert_eq!(node.schedule(0.1, false).0, PowerMode::Critical);
assert_eq!(node.schedule(0.1, true).0, PowerMode::Saver);
assert_eq!(node.schedule(0.9, false).1, Duration::from_secs(60));
}
}