use std::collections::HashMap;
use crate::adapter::net::behavior::fold::{IslandId, NodeId};
use crate::adapter::net::behavior::meshos::{DaemonIntentUpdate, DaemonLifecycleSignal, DaemonRef};
use crate::adapter::net::cortex::workflow::{ActiveClaim, TaskId, WorkflowState};
use super::claim_registry::ClaimRegistry;
use super::daemon_ref::daemon_ref;
use super::lifecycle::{
apply_lifecycle, build_daemon_task_map, signal_implies_transition, LifecycleTransition,
};
use super::migration::{ClaimHeld, MigrationEligible};
use super::projection::{project_daemon_intents, project_forced_placements};
pub fn desired_daemon_intents<F>(
workflow: &WorkflowState,
claims: &ClaimRegistry,
resolve_host: F,
) -> Vec<DaemonIntentUpdate>
where
F: Fn(IslandId) -> Option<NodeId>,
{
let mut by_daemon: HashMap<DaemonRef, DaemonIntentUpdate> = project_daemon_intents(workflow)
.into_iter()
.map(|u| (u.daemon.clone(), u))
.collect();
for (daemon, host) in project_forced_placements(claims, resolve_host) {
if let Some(intent) = by_daemon.get_mut(&daemon) {
intent.node = Some(host);
}
}
let mut out: Vec<DaemonIntentUpdate> = by_daemon.into_values().collect();
out.sort_by_key(|u| u.daemon.id);
out
}
#[derive(Debug, Default)]
pub struct SchedulerBridge {
claims: ClaimRegistry,
}
impl SchedulerBridge {
pub fn new() -> Self {
Self::default()
}
pub fn on_running(&mut self, task: TaskId, claim: ActiveClaim) {
self.claims.insert(daemon_ref(task), claim);
}
pub fn on_released(&mut self, task: TaskId) -> Option<ActiveClaim> {
self.claims.remove(&daemon_ref(task))
}
pub fn desired_intents<F>(
&self,
workflow: &WorkflowState,
resolve_host: F,
) -> Vec<DaemonIntentUpdate>
where
F: Fn(IslandId) -> Option<NodeId>,
{
desired_daemon_intents(workflow, &self.claims, resolve_host)
}
pub fn lifecycle_transition(
&self,
signal: &DaemonLifecycleSignal,
daemon: &DaemonRef,
workflow: &WorkflowState,
) -> Option<LifecycleTransition> {
if !signal_implies_transition(signal) {
return None;
}
apply_lifecycle(signal, daemon, &build_daemon_task_map(workflow))
}
pub fn check_migration(&self, daemon: DaemonRef) -> Result<MigrationEligible, ClaimHeld> {
MigrationEligible::check(daemon, &self.claims)
}
pub fn claims(&self) -> &ClaimRegistry {
&self.claims
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::net::cortex::workflow::WorkflowAdapter;
use crate::adapter::net::redex::Redex;
#[test]
fn claim_lifecycle_drives_the_migration_veto() {
let mut bridge = SchedulerBridge::new();
let d1 = daemon_ref(1);
assert!(bridge.check_migration(d1.clone()).is_ok());
bridge.on_running(1, ActiveClaim { island: 0xA0 });
assert!(bridge.check_migration(d1.clone()).is_err());
assert_eq!(bridge.claims().get(&d1).map(|c| c.island), Some(0xA0));
assert_eq!(bridge.on_released(1).map(|c| c.island), Some(0xA0));
assert!(bridge.check_migration(d1).is_ok());
}
#[tokio::test]
async fn desired_intents_pins_only_the_claim_holder() {
let redex = Redex::new();
let wf = WorkflowAdapter::open(&redex, 0x5C4E_DB06).await.unwrap();
wf.submit(1).unwrap();
wf.submit(2).unwrap();
wf.start(1).unwrap();
let seq = wf.start(2).unwrap();
wf.wait_for_seq(seq).await.unwrap();
let mut bridge = SchedulerBridge::new();
bridge.on_running(1, ActiveClaim { island: 0xA0 });
let resolve = |island| (island == 0xA0).then_some(7);
let intents = bridge.desired_intents(&wf.state().read(), resolve);
let by_name: HashMap<String, &DaemonIntentUpdate> =
intents.iter().map(|u| (u.daemon.name.clone(), u)).collect();
assert_eq!(intents.len(), 2);
assert_eq!(by_name["task/1"].node, Some(7));
assert!(by_name["task/2"].node.is_none());
}
#[tokio::test]
async fn lifecycle_transition_maps_a_crash_to_failstep() {
use std::time::Instant;
let redex = Redex::new();
let wf = WorkflowAdapter::open(&redex, 0x5C4E_DB07).await.unwrap();
let seq = wf.submit(1).unwrap();
wf.wait_for_seq(seq).await.unwrap();
let bridge = SchedulerBridge::new();
let crash = DaemonLifecycleSignal::Crashed {
at: Instant::now(),
reason: "oom".into(),
};
assert_eq!(
bridge.lifecycle_transition(&crash, &daemon_ref(1), &wf.state().read()),
Some(LifecycleTransition::FailStep(1)),
);
}
}