use std::collections::HashMap;
use crate::adapter::net::behavior::meshos::{DaemonLifecycleSignal, DaemonRef};
use crate::adapter::net::cortex::workflow::{TaskId, WorkflowState};
use super::daemon_ref::daemon_ref;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LifecycleTransition {
ConfirmRunning(TaskId),
FailStep(TaskId),
}
pub fn build_daemon_task_map(workflow: &WorkflowState) -> HashMap<DaemonRef, TaskId> {
workflow.all().map(|(id, _)| (daemon_ref(id), id)).collect()
}
pub fn apply_lifecycle(
signal: &DaemonLifecycleSignal,
daemon: &DaemonRef,
daemon_task: &HashMap<DaemonRef, TaskId>,
) -> Option<LifecycleTransition> {
let task = *daemon_task.get(daemon)?;
match signal {
DaemonLifecycleSignal::Started { .. } => Some(LifecycleTransition::ConfirmRunning(task)),
DaemonLifecycleSignal::Crashed { .. } => Some(LifecycleTransition::FailStep(task)),
DaemonLifecycleSignal::ExitedCleanly { .. }
| DaemonLifecycleSignal::HealthChanged { .. }
| DaemonLifecycleSignal::SaturationChanged { .. } => None,
}
}
pub fn signal_implies_transition(signal: &DaemonLifecycleSignal) -> bool {
match signal {
DaemonLifecycleSignal::Started { .. } | DaemonLifecycleSignal::Crashed { .. } => true,
DaemonLifecycleSignal::ExitedCleanly { .. }
| DaemonLifecycleSignal::HealthChanged { .. }
| DaemonLifecycleSignal::SaturationChanged { .. } => false,
}
}
#[cfg(test)]
mod tests {
use std::time::Instant;
use super::*;
fn map_of(ids: &[TaskId]) -> HashMap<DaemonRef, TaskId> {
ids.iter().map(|&id| (daemon_ref(id), id)).collect()
}
#[test]
fn started_confirms_running_crashed_fails_clean_exit_is_noop() {
let map = map_of(&[1]);
let d = daemon_ref(1);
let at = Instant::now();
assert_eq!(
apply_lifecycle(&DaemonLifecycleSignal::Started { at }, &d, &map),
Some(LifecycleTransition::ConfirmRunning(1)),
);
assert_eq!(
apply_lifecycle(
&DaemonLifecycleSignal::Crashed {
at,
reason: "segfault".into(),
},
&d,
&map,
),
Some(LifecycleTransition::FailStep(1)),
);
assert_eq!(
apply_lifecycle(&DaemonLifecycleSignal::ExitedCleanly { at }, &d, &map),
None,
"graceful exit is the expected terminal/cancel shutdown",
);
}
#[test]
fn only_started_and_crashed_imply_a_transition() {
use crate::adapter::net::behavior::meshos::DaemonHealth;
let at = Instant::now();
assert!(signal_implies_transition(&DaemonLifecycleSignal::Started {
at
}));
assert!(signal_implies_transition(&DaemonLifecycleSignal::Crashed {
at,
reason: "x".into(),
}));
assert!(!signal_implies_transition(
&DaemonLifecycleSignal::ExitedCleanly { at }
));
assert!(!signal_implies_transition(
&DaemonLifecycleSignal::HealthChanged {
at,
health: DaemonHealth::Healthy,
}
));
assert!(!signal_implies_transition(
&DaemonLifecycleSignal::SaturationChanged {
at,
saturation: 0.5
}
));
}
#[test]
fn a_system_daemon_signal_yields_no_transition() {
let map = map_of(&[1, 2]);
let other = DaemonRef {
id: 7,
name: "telemetry".into(),
};
assert_eq!(
apply_lifecycle(
&DaemonLifecycleSignal::Crashed {
at: Instant::now(),
reason: "x".into(),
},
&other,
&map,
),
None,
);
}
#[test]
fn resolves_each_daemon_back_to_its_own_task() {
let map = map_of(&[10, 20, 30]);
let at = Instant::now();
for id in [10u64, 20, 30] {
assert_eq!(
apply_lifecycle(
&DaemonLifecycleSignal::Started { at },
&daemon_ref(id),
&map
),
Some(LifecycleTransition::ConfirmRunning(id)),
);
}
}
#[tokio::test]
async fn build_map_recovers_tasks_from_live_workflow_state() {
use crate::adapter::net::cortex::workflow::WorkflowAdapter;
use crate::adapter::net::redex::Redex;
let redex = Redex::new();
let wf = WorkflowAdapter::open(&redex, 0x5C4E_DB05).await.unwrap();
wf.submit(1).unwrap();
let seq = wf.submit(2).unwrap();
wf.wait_for_seq(seq).await.unwrap();
let map = build_daemon_task_map(&wf.state().read());
assert_eq!(map.get(&daemon_ref(1)), Some(&1));
assert_eq!(map.get(&daemon_ref(2)), Some(&2));
assert_eq!(map.len(), 2);
}
}