use std::collections::HashSet;
use std::sync::{Arc, Mutex};
use std::time::Instant;
use crate::command::{Command, CommandError};
use crate::source::StateSource;
use crate::state::{EventRecord, Group, Link, LinkKind, Mode, Org, Reading, Sensor, State, Status};
const HISTORY: usize = 32;
struct Inner {
state: State,
commands: Vec<Command>,
started: Instant,
allowed: Option<HashSet<String>>,
}
#[derive(Clone)]
pub struct Fleet {
inner: Arc<Mutex<Inner>>,
}
impl Fleet {
pub fn builder() -> FleetBuilder {
FleetBuilder { orgs: Vec::new() }
}
pub fn from_state(state: State) -> Self {
Self {
inner: Arc::new(Mutex::new(Inner {
state,
commands: Vec::new(),
started: Instant::now(),
allowed: None,
})),
}
}
pub fn allow_sensors(&self, keys: impl IntoIterator<Item = impl Into<String>>) {
let mut inner = self.inner.lock().expect("fleet lock");
inner.allowed = Some(keys.into_iter().map(Into::into).collect());
}
pub fn report_reading(&self, group: &str, sensor: &str, reading: Reading) {
let mut inner = self.inner.lock().expect("fleet lock");
if let Some(target) = sensor_mut(&mut inner.state, group, sensor) {
let value = reading.value;
target.reading = reading;
target.history.push(value);
let len = target.history.len();
if len > HISTORY {
target.history.drain(0..len - HISTORY);
}
}
recompute(&mut inner.state);
}
pub fn report_event(&self, group: &str, sensor: &str, event: EventRecord) {
let mut inner = self.inner.lock().expect("fleet lock");
if let Some(target) = sensor_mut(&mut inner.state, group, sensor) {
target.events.insert(0, event);
target.events.truncate(8);
}
recompute(&mut inner.state);
}
pub fn report_link(&self, group: &str, link: Link) {
let mut inner = self.inner.lock().expect("fleet lock");
if let Some(target) = group_mut(&mut inner.state, group) {
target.link = link;
}
recompute(&mut inner.state);
}
pub fn report_power(&self, group: &str, sensor: &str, mode: Mode, battery: Option<f32>) {
let mut inner = self.inner.lock().expect("fleet lock");
if let Some(target) = sensor_mut(&mut inner.state, group, sensor) {
target.mode = mode;
target.battery = battery;
}
}
pub fn take_commands(&self) -> Vec<Command> {
let mut inner = self.inner.lock().expect("fleet lock");
std::mem::take(&mut inner.commands)
}
pub fn add_group(&self, org: &str, group: Group) {
self.mutate(Command::AddGroup {
org: org.to_owned(),
group,
});
}
pub fn add_sensor(&self, group: &str, sensor: Sensor) {
self.mutate(Command::AddSensor {
group: group.to_owned(),
sensor,
binding: None,
});
}
pub fn remove_group(&self, id: &str) {
self.mutate(Command::RemoveGroup { id: id.to_owned() });
}
pub fn remove_sensor(&self, target: &str) {
self.mutate(Command::RemoveSensor {
target: target.to_owned(),
});
}
fn mutate(&self, command: Command) {
let mut inner = self.inner.lock().expect("fleet lock");
let _ = apply(&mut inner.state, &command);
recompute(&mut inner.state);
}
}
impl StateSource for Fleet {
fn snapshot(&mut self) -> State {
let mut inner = self.inner.lock().expect("fleet lock");
let uptime = inner.started.elapsed().as_secs();
inner.state.uptime_secs = Some(uptime);
inner.state.clone()
}
fn command(&mut self, command: &Command) -> Result<(), CommandError> {
let mut inner = self.inner.lock().expect("fleet lock");
if let Command::AddSensor { sensor, .. } = command {
if let Some(allowed) = &inner.allowed {
if !allowed.contains(&sensor.reading.key) {
return Err(CommandError::UnknownSensor);
}
}
}
let outcome = apply(&mut inner.state, command);
if outcome.is_ok() {
inner.commands.push(command.clone());
recompute(&mut inner.state);
}
outcome
}
}
fn apply(state: &mut State, command: &Command) -> Result<(), CommandError> {
match command {
Command::Actuate { target, action } => {
let (group, sensor) = target.split_once('/').ok_or(CommandError::UnknownTarget)?;
let reading = sensor_mut(state, group, sensor)
.map(|s| &mut s.reading)
.ok_or(CommandError::UnknownTarget)?;
match &reading.actions {
Some(actions) if actions.iter().any(|a| a == action) => {
reading.state = Some(format!("state.{action}"));
Ok(())
}
Some(_) => Err(CommandError::InvalidAction),
None => Err(CommandError::Unsupported),
}
}
Command::AddGroup { org, group } => match org_mut(state, org) {
Some(target) => {
target.groups.push(group.clone());
Ok(())
}
None => Err(CommandError::UnknownTarget),
},
Command::RemoveGroup { id } => {
for org in &mut state.orgs {
org.groups.retain(|g| g.id != *id);
}
Ok(())
}
Command::AddSensor { group, sensor, .. } => match group_mut(state, group) {
Some(target) => {
target.sensors.push(sensor.clone());
Ok(())
}
None => Err(CommandError::UnknownTarget),
},
Command::RemoveSensor { target } => {
let (group_id, sensor_id) = target.split_once('/').unwrap_or(("", target));
for org in &mut state.orgs {
for group in &mut org.groups {
if group.id == group_id {
group.sensors.retain(|s| s.id != sensor_id);
}
}
}
Ok(())
}
}
}
fn recompute(state: &mut State) {
for org in &mut state.orgs {
for group in &mut org.groups {
group.recompute_status();
}
}
state.recompute_status();
}
fn org_mut<'a>(state: &'a mut State, org: &str) -> Option<&'a mut Org> {
state.orgs.iter_mut().find(|o| o.id == org)
}
fn group_mut<'a>(state: &'a mut State, group: &str) -> Option<&'a mut Group> {
state
.orgs
.iter_mut()
.flat_map(|o| &mut o.groups)
.find(|g| g.id == group)
}
fn sensor_mut<'a>(state: &'a mut State, group: &str, sensor: &str) -> Option<&'a mut Sensor> {
state
.orgs
.iter_mut()
.flat_map(|o| &mut o.groups)
.filter(|g| g.id == group)
.flat_map(|g| &mut g.sensors)
.find(|s| s.id == sensor)
}
pub struct FleetBuilder {
orgs: Vec<Org>,
}
impl FleetBuilder {
pub fn org(mut self, id: impl Into<String>, name: impl Into<String>) -> Self {
self.orgs.push(Org {
id: id.into(),
name: name.into(),
groups: Vec::new(),
});
self
}
pub fn group(
mut self,
org: &str,
id: impl Into<String>,
name: impl Into<String>,
kind: LinkKind,
) -> Self {
if let Some(target) = self.orgs.iter_mut().find(|o| o.id == org) {
target.groups.push(Group {
id: id.into(),
name: name.into(),
link: Link {
kind,
strength: 4,
online: true,
},
status: Status::Ok,
sensors: Vec::new(),
lat: None,
lon: None,
});
}
self
}
pub fn sensor(mut self, group: &str, sensor: Sensor) -> Self {
for org in &mut self.orgs {
if let Some(target) = org.groups.iter_mut().find(|g| g.id == group) {
target.sensors.push(sensor);
break;
}
}
self
}
pub fn build(mut self) -> Fleet {
let mut state = State {
orgs: std::mem::take(&mut self.orgs),
status: Status::Ok,
uptime_secs: None,
demo: false,
};
recompute(&mut state);
Fleet::from_state(state)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn fleet() -> Fleet {
Fleet::builder()
.org("clinic", "Kano clinic")
.group("clinic", "fridges", "Cold chain", LinkKind::Cellular)
.sensor(
"fridges",
Sensor::new("fridge-1", Reading::new("fridge_temp", 4.5, "celsius")),
)
.sensor(
"fridges",
Sensor::new(
"valve",
Reading::new("drip_valve", 0.0, "state")
.with_state("state.closed")
.with_actions(["open", "closed"]),
),
)
.build()
}
#[test]
fn a_reported_reading_shows_in_the_snapshot_with_history() {
let fleet = fleet();
fleet.report_reading(
"fridges",
"fridge-1",
Reading::new("fridge_temp", 9.0, "celsius").with_status(Status::Alarm),
);
let mut handle = fleet.clone();
let state = handle.snapshot();
let sensor = &state.orgs[0].groups[0].sensors[0];
assert_eq!(sensor.reading.value, 9.0);
assert_eq!(sensor.history, vec![9.0]);
assert_eq!(
state.status,
Status::Alarm,
"the alarm reading lifts fleet status"
);
}
#[test]
fn an_actuate_command_updates_state_and_queues_for_the_project() {
let mut fleet = fleet();
fleet
.command(&Command::Actuate {
target: "fridges/valve".to_owned(),
action: "open".to_owned(),
})
.expect("valve accepts open");
let queued = fleet.take_commands();
assert_eq!(
queued.len(),
1,
"the command is queued for the project to apply"
);
let valve = sensor_after(&fleet, "fridges", "valve");
assert_eq!(valve.reading.state.as_deref(), Some("state.open"));
assert!(fleet.take_commands().is_empty());
}
#[test]
fn an_invalid_actuate_is_refused_and_not_queued() {
let mut fleet = fleet();
assert_eq!(
fleet.command(&Command::Actuate {
target: "fridges/fridge-1".to_owned(),
action: "open".to_owned(),
}),
Err(CommandError::Unsupported)
);
assert!(fleet.take_commands().is_empty());
}
#[test]
fn provisioning_commands_change_the_structure() {
let mut fleet = fleet();
fleet
.command(&Command::AddSensor {
group: "fridges".to_owned(),
sensor: Sensor::new("fridge-2", Reading::new("fridge_temp", 5.0, "celsius")),
binding: None,
})
.expect("add sensor to a known group");
assert!(sensor_present(&fleet, "fridges", "fridge-2"));
fleet
.command(&Command::RemoveSensor {
target: "fridges/fridge-2".to_owned(),
})
.expect("remove the sensor");
assert!(!sensor_present(&fleet, "fridges", "fridge-2"));
}
#[test]
fn runtime_mutators_add_and_remove_for_discovery() {
let fleet = fleet();
fleet.add_group(
"clinic",
Group {
id: "ward".to_owned(),
name: "Ward".to_owned(),
link: Link {
kind: LinkKind::Wifi,
strength: 4,
online: true,
},
status: Status::Ok,
sensors: Vec::new(),
lat: None,
lon: None,
},
);
fleet.add_sensor(
"ward",
Sensor::new("o2", Reading::new("oxygen_stock", 80.0, "percent")),
);
assert!(
sensor_present(&fleet, "ward", "o2"),
"discovered sensor shows"
);
fleet.remove_group("ward");
let mut handle = fleet.clone();
assert!(
!handle
.snapshot()
.orgs
.iter()
.flat_map(|o| &o.groups)
.any(|g| g.id == "ward"),
"removed group is gone"
);
}
#[test]
fn an_allow_list_rejects_unsupported_client_adds_but_not_discovery() {
let mut fleet = fleet();
fleet.allow_sensors(["fridge_temp"]);
fleet
.command(&Command::AddSensor {
group: "fridges".to_owned(),
sensor: Sensor::new("f2", Reading::new("fridge_temp", 5.0, "celsius")),
binding: None,
})
.expect("a supported sensor is added");
assert_eq!(
fleet.command(&Command::AddSensor {
group: "fridges".to_owned(),
sensor: Sensor::new("w", Reading::new("wind_speed", 3.0, "meter_per_second")),
binding: None,
}),
Err(CommandError::UnknownSensor)
);
assert!(!sensor_present(&fleet, "fridges", "w"));
fleet.add_sensor(
"fridges",
Sensor::new("disc", Reading::new("wind_speed", 3.0, "meter_per_second")),
);
assert!(sensor_present(&fleet, "fridges", "disc"));
}
#[test]
fn from_state_restores_a_saved_fleet() {
let mut original = fleet();
let saved = original.snapshot();
let mut restored = Fleet::from_state(saved.clone());
assert_eq!(restored.snapshot().orgs.len(), saved.orgs.len());
}
fn sensor_after(fleet: &Fleet, group: &str, sensor: &str) -> Sensor {
let mut handle = fleet.clone();
let state = handle.snapshot();
state
.orgs
.iter()
.flat_map(|o| &o.groups)
.filter(|g| g.id == group)
.flat_map(|g| &g.sensors)
.find(|s| s.id == sensor)
.expect("sensor")
.clone()
}
fn sensor_present(fleet: &Fleet, group: &str, sensor: &str) -> bool {
let mut handle = fleet.clone();
handle
.snapshot()
.orgs
.iter()
.flat_map(|o| &o.groups)
.filter(|g| g.id == group)
.flat_map(|g| &g.sensors)
.any(|s| s.id == sensor)
}
}