1use std::collections::HashSet;
42use std::sync::{Arc, Mutex};
43use std::time::Instant;
44
45use crate::command::{Command, CommandError};
46use crate::source::StateSource;
47use crate::state::{EventRecord, Group, Link, LinkKind, Mode, Org, Reading, Sensor, State, Status};
48
49const HISTORY: usize = 32;
51
52struct Inner {
53 state: State,
54 commands: Vec<Command>,
55 started: Instant,
56 allowed: Option<HashSet<String>>,
60}
61
62#[derive(Clone)]
65pub struct Fleet {
66 inner: Arc<Mutex<Inner>>,
67}
68
69impl Fleet {
70 pub fn builder() -> FleetBuilder {
76 FleetBuilder { orgs: Vec::new() }
77 }
78
79 pub fn from_state(state: State) -> Self {
91 Self {
92 inner: Arc::new(Mutex::new(Inner {
93 state,
94 commands: Vec::new(),
95 started: Instant::now(),
96 allowed: None,
97 })),
98 }
99 }
100
101 pub fn allow_sensors(&self, keys: impl IntoIterator<Item = impl Into<String>>) {
116 let mut inner = self.inner.lock().expect("fleet lock");
117 inner.allowed = Some(keys.into_iter().map(Into::into).collect());
118 }
119
120 pub fn report_reading(&self, group: &str, sensor: &str, reading: Reading) {
128 let mut inner = self.inner.lock().expect("fleet lock");
129 if let Some(target) = sensor_mut(&mut inner.state, group, sensor) {
130 let value = reading.value;
131 target.reading = reading;
132 target.history.push(value);
133 let len = target.history.len();
134 if len > HISTORY {
135 target.history.drain(0..len - HISTORY);
136 }
137 }
138 recompute(&mut inner.state);
139 }
140
141 pub fn report_event(&self, group: &str, sensor: &str, event: EventRecord) {
149 let mut inner = self.inner.lock().expect("fleet lock");
150 if let Some(target) = sensor_mut(&mut inner.state, group, sensor) {
151 target.events.insert(0, event);
152 target.events.truncate(8);
153 }
154 recompute(&mut inner.state);
155 }
156
157 pub fn report_link(&self, group: &str, link: Link) {
164 let mut inner = self.inner.lock().expect("fleet lock");
165 if let Some(target) = group_mut(&mut inner.state, group) {
166 target.link = link;
167 }
168 recompute(&mut inner.state);
169 }
170
171 pub fn report_power(&self, group: &str, sensor: &str, mode: Mode, battery: Option<f32>) {
180 let mut inner = self.inner.lock().expect("fleet lock");
181 if let Some(target) = sensor_mut(&mut inner.state, group, sensor) {
182 target.mode = mode;
183 target.battery = battery;
184 }
185 }
186
187 pub fn take_commands(&self) -> Vec<Command> {
194 let mut inner = self.inner.lock().expect("fleet lock");
195 std::mem::take(&mut inner.commands)
196 }
197
198 pub fn add_group(&self, org: &str, group: Group) {
206 self.mutate(Command::AddGroup {
207 org: org.to_owned(),
208 group,
209 });
210 }
211
212 pub fn add_sensor(&self, group: &str, sensor: Sensor) {
220 self.mutate(Command::AddSensor {
221 group: group.to_owned(),
222 sensor,
223 binding: None,
224 });
225 }
226
227 pub fn remove_group(&self, id: &str) {
233 self.mutate(Command::RemoveGroup { id: id.to_owned() });
234 }
235
236 pub fn remove_sensor(&self, target: &str) {
242 self.mutate(Command::RemoveSensor {
243 target: target.to_owned(),
244 });
245 }
246
247 fn mutate(&self, command: Command) {
250 let mut inner = self.inner.lock().expect("fleet lock");
251 let _ = apply(&mut inner.state, &command);
252 recompute(&mut inner.state);
253 }
254}
255
256impl StateSource for Fleet {
257 fn snapshot(&mut self) -> State {
258 let mut inner = self.inner.lock().expect("fleet lock");
259 let uptime = inner.started.elapsed().as_secs();
260 inner.state.uptime_secs = Some(uptime);
261 inner.state.clone()
262 }
263
264 fn command(&mut self, command: &Command) -> Result<(), CommandError> {
265 let mut inner = self.inner.lock().expect("fleet lock");
266 if let Command::AddSensor { sensor, .. } = command {
269 if let Some(allowed) = &inner.allowed {
270 if !allowed.contains(&sensor.reading.key) {
271 return Err(CommandError::UnknownSensor);
272 }
273 }
274 }
275 let outcome = apply(&mut inner.state, command);
276 if outcome.is_ok() {
277 inner.commands.push(command.clone());
278 recompute(&mut inner.state);
279 }
280 outcome
281 }
282}
283
284fn apply(state: &mut State, command: &Command) -> Result<(), CommandError> {
287 match command {
288 Command::Actuate { target, action } => {
289 let (group, sensor) = target.split_once('/').ok_or(CommandError::UnknownTarget)?;
290 let reading = sensor_mut(state, group, sensor)
291 .map(|s| &mut s.reading)
292 .ok_or(CommandError::UnknownTarget)?;
293 match &reading.actions {
294 Some(actions) if actions.iter().any(|a| a == action) => {
295 reading.state = Some(format!("state.{action}"));
296 Ok(())
297 }
298 Some(_) => Err(CommandError::InvalidAction),
299 None => Err(CommandError::Unsupported),
300 }
301 }
302 Command::AddGroup { org, group } => match org_mut(state, org) {
303 Some(target) => {
304 target.groups.push(group.clone());
305 Ok(())
306 }
307 None => Err(CommandError::UnknownTarget),
308 },
309 Command::RemoveGroup { id } => {
310 for org in &mut state.orgs {
311 org.groups.retain(|g| g.id != *id);
312 }
313 Ok(())
314 }
315 Command::AddSensor { group, sensor, .. } => match group_mut(state, group) {
316 Some(target) => {
317 target.sensors.push(sensor.clone());
318 Ok(())
319 }
320 None => Err(CommandError::UnknownTarget),
321 },
322 Command::RemoveSensor { target } => {
323 let (group_id, sensor_id) = target.split_once('/').unwrap_or(("", target));
324 for org in &mut state.orgs {
325 for group in &mut org.groups {
326 if group.id == group_id {
327 group.sensors.retain(|s| s.id != sensor_id);
328 }
329 }
330 }
331 Ok(())
332 }
333 }
334}
335
336fn recompute(state: &mut State) {
337 for org in &mut state.orgs {
338 for group in &mut org.groups {
339 group.recompute_status();
340 }
341 }
342 state.recompute_status();
343}
344
345fn org_mut<'a>(state: &'a mut State, org: &str) -> Option<&'a mut Org> {
346 state.orgs.iter_mut().find(|o| o.id == org)
347}
348
349fn group_mut<'a>(state: &'a mut State, group: &str) -> Option<&'a mut Group> {
350 state
351 .orgs
352 .iter_mut()
353 .flat_map(|o| &mut o.groups)
354 .find(|g| g.id == group)
355}
356
357fn sensor_mut<'a>(state: &'a mut State, group: &str, sensor: &str) -> Option<&'a mut Sensor> {
358 state
359 .orgs
360 .iter_mut()
361 .flat_map(|o| &mut o.groups)
362 .filter(|g| g.id == group)
363 .flat_map(|g| &mut g.sensors)
364 .find(|s| s.id == sensor)
365}
366
367pub struct FleetBuilder {
371 orgs: Vec<Org>,
372}
373
374impl FleetBuilder {
375 pub fn org(mut self, id: impl Into<String>, name: impl Into<String>) -> Self {
386 self.orgs.push(Org {
387 id: id.into(),
388 name: name.into(),
389 groups: Vec::new(),
390 });
391 self
392 }
393
394 pub fn group(
407 mut self,
408 org: &str,
409 id: impl Into<String>,
410 name: impl Into<String>,
411 kind: LinkKind,
412 ) -> Self {
413 if let Some(target) = self.orgs.iter_mut().find(|o| o.id == org) {
414 target.groups.push(Group {
415 id: id.into(),
416 name: name.into(),
417 link: Link {
418 kind,
419 strength: 4,
420 online: true,
421 },
422 status: Status::Ok,
423 sensors: Vec::new(),
424 lat: None,
425 lon: None,
426 });
427 }
428 self
429 }
430
431 pub fn sensor(mut self, group: &str, sensor: Sensor) -> Self {
442 for org in &mut self.orgs {
443 if let Some(target) = org.groups.iter_mut().find(|g| g.id == group) {
444 target.sensors.push(sensor);
445 break;
446 }
447 }
448 self
449 }
450
451 pub fn build(mut self) -> Fleet {
457 let mut state = State {
458 orgs: std::mem::take(&mut self.orgs),
459 status: Status::Ok,
460 uptime_secs: None,
461 demo: false,
462 };
463 recompute(&mut state);
464 Fleet::from_state(state)
465 }
466}
467
468#[cfg(test)]
469mod tests {
470 use super::*;
471
472 fn fleet() -> Fleet {
473 Fleet::builder()
474 .org("clinic", "Kano clinic")
475 .group("clinic", "fridges", "Cold chain", LinkKind::Cellular)
476 .sensor(
477 "fridges",
478 Sensor::new("fridge-1", Reading::new("fridge_temp", 4.5, "celsius")),
479 )
480 .sensor(
481 "fridges",
482 Sensor::new(
483 "valve",
484 Reading::new("drip_valve", 0.0, "state")
485 .with_state("state.closed")
486 .with_actions(["open", "closed"]),
487 ),
488 )
489 .build()
490 }
491
492 #[test]
493 fn a_reported_reading_shows_in_the_snapshot_with_history() {
494 let fleet = fleet();
495 fleet.report_reading(
496 "fridges",
497 "fridge-1",
498 Reading::new("fridge_temp", 9.0, "celsius").with_status(Status::Alarm),
499 );
500 let mut handle = fleet.clone();
501 let state = handle.snapshot();
502 let sensor = &state.orgs[0].groups[0].sensors[0];
503 assert_eq!(sensor.reading.value, 9.0);
504 assert_eq!(sensor.history, vec![9.0]);
505 assert_eq!(
506 state.status,
507 Status::Alarm,
508 "the alarm reading lifts fleet status"
509 );
510 }
511
512 #[test]
513 fn an_actuate_command_updates_state_and_queues_for_the_project() {
514 let mut fleet = fleet();
515 fleet
516 .command(&Command::Actuate {
517 target: "fridges/valve".to_owned(),
518 action: "open".to_owned(),
519 })
520 .expect("valve accepts open");
521 let queued = fleet.take_commands();
522 assert_eq!(
523 queued.len(),
524 1,
525 "the command is queued for the project to apply"
526 );
527 let valve = sensor_after(&fleet, "fridges", "valve");
528 assert_eq!(valve.reading.state.as_deref(), Some("state.open"));
529 assert!(fleet.take_commands().is_empty());
531 }
532
533 #[test]
534 fn an_invalid_actuate_is_refused_and_not_queued() {
535 let mut fleet = fleet();
536 assert_eq!(
537 fleet.command(&Command::Actuate {
538 target: "fridges/fridge-1".to_owned(),
539 action: "open".to_owned(),
540 }),
541 Err(CommandError::Unsupported)
542 );
543 assert!(fleet.take_commands().is_empty());
544 }
545
546 #[test]
547 fn provisioning_commands_change_the_structure() {
548 let mut fleet = fleet();
549 fleet
550 .command(&Command::AddSensor {
551 group: "fridges".to_owned(),
552 sensor: Sensor::new("fridge-2", Reading::new("fridge_temp", 5.0, "celsius")),
553 binding: None,
554 })
555 .expect("add sensor to a known group");
556 assert!(sensor_present(&fleet, "fridges", "fridge-2"));
557
558 fleet
559 .command(&Command::RemoveSensor {
560 target: "fridges/fridge-2".to_owned(),
561 })
562 .expect("remove the sensor");
563 assert!(!sensor_present(&fleet, "fridges", "fridge-2"));
564 }
565
566 #[test]
567 fn runtime_mutators_add_and_remove_for_discovery() {
568 let fleet = fleet();
569 fleet.add_group(
570 "clinic",
571 Group {
572 id: "ward".to_owned(),
573 name: "Ward".to_owned(),
574 link: Link {
575 kind: LinkKind::Wifi,
576 strength: 4,
577 online: true,
578 },
579 status: Status::Ok,
580 sensors: Vec::new(),
581 lat: None,
582 lon: None,
583 },
584 );
585 fleet.add_sensor(
586 "ward",
587 Sensor::new("o2", Reading::new("oxygen_stock", 80.0, "percent")),
588 );
589 assert!(
590 sensor_present(&fleet, "ward", "o2"),
591 "discovered sensor shows"
592 );
593 fleet.remove_group("ward");
594 let mut handle = fleet.clone();
595 assert!(
596 !handle
597 .snapshot()
598 .orgs
599 .iter()
600 .flat_map(|o| &o.groups)
601 .any(|g| g.id == "ward"),
602 "removed group is gone"
603 );
604 }
605
606 #[test]
607 fn an_allow_list_rejects_unsupported_client_adds_but_not_discovery() {
608 let mut fleet = fleet();
609 fleet.allow_sensors(["fridge_temp"]);
610
611 fleet
613 .command(&Command::AddSensor {
614 group: "fridges".to_owned(),
615 sensor: Sensor::new("f2", Reading::new("fridge_temp", 5.0, "celsius")),
616 binding: None,
617 })
618 .expect("a supported sensor is added");
619
620 assert_eq!(
622 fleet.command(&Command::AddSensor {
623 group: "fridges".to_owned(),
624 sensor: Sensor::new("w", Reading::new("wind_speed", 3.0, "meter_per_second")),
625 binding: None,
626 }),
627 Err(CommandError::UnknownSensor)
628 );
629 assert!(!sensor_present(&fleet, "fridges", "w"));
630
631 fleet.add_sensor(
633 "fridges",
634 Sensor::new("disc", Reading::new("wind_speed", 3.0, "meter_per_second")),
635 );
636 assert!(sensor_present(&fleet, "fridges", "disc"));
637 }
638
639 #[test]
640 fn from_state_restores_a_saved_fleet() {
641 let mut original = fleet();
642 let saved = original.snapshot();
643 let mut restored = Fleet::from_state(saved.clone());
644 assert_eq!(restored.snapshot().orgs.len(), saved.orgs.len());
645 }
646
647 fn sensor_after(fleet: &Fleet, group: &str, sensor: &str) -> Sensor {
648 let mut handle = fleet.clone();
649 let state = handle.snapshot();
650 state
651 .orgs
652 .iter()
653 .flat_map(|o| &o.groups)
654 .filter(|g| g.id == group)
655 .flat_map(|g| &g.sensors)
656 .find(|s| s.id == sensor)
657 .expect("sensor")
658 .clone()
659 }
660
661 fn sensor_present(fleet: &Fleet, group: &str, sensor: &str) -> bool {
662 let mut handle = fleet.clone();
663 handle
664 .snapshot()
665 .orgs
666 .iter()
667 .flat_map(|o| &o.groups)
668 .filter(|g| g.id == group)
669 .flat_map(|g| &g.sensors)
670 .any(|s| s.id == sensor)
671 }
672}