use crate::adapter::net::behavior::fold::IslandId;
use crate::adapter::net::behavior::meshos::{DaemonIntent, DaemonIntentUpdate, DaemonRef, NodeId};
use crate::adapter::net::cortex::workflow::{TaskStatus, WorkflowState};
use super::claim_registry::ClaimRegistry;
use super::daemon_ref::daemon_ref;
pub fn project_daemon_intents(workflow: &WorkflowState) -> Vec<DaemonIntentUpdate> {
let mut out: Vec<DaemonIntentUpdate> = workflow
.all()
.map(|(id, state)| DaemonIntentUpdate {
daemon: daemon_ref(id),
intent: if state.status == TaskStatus::Running {
DaemonIntent::Run
} else {
DaemonIntent::Stop
},
node: None,
})
.collect();
out.sort_by_key(|update| update.daemon.id);
out
}
pub fn project_forced_placements<F>(
claims: &ClaimRegistry,
resolve_host: F,
) -> Vec<(DaemonRef, NodeId)>
where
F: Fn(IslandId) -> Option<NodeId>,
{
let mut out: Vec<(DaemonRef, NodeId)> = claims
.iter()
.filter_map(|(daemon, claim)| resolve_host(claim.island).map(|host| (daemon.clone(), host)))
.collect();
out.sort_by_key(|(daemon, _)| daemon.id);
out
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::net::cortex::workflow::WorkflowAdapter;
use crate::adapter::net::redex::Redex;
#[tokio::test]
async fn running_projects_run_everything_else_projects_stop() {
let redex = Redex::new();
let wf = WorkflowAdapter::open(&redex, 0x5C4E_DB01).await.unwrap();
for id in 1..=4u64 {
wf.submit(id).unwrap();
}
wf.start(1).unwrap(); wf.start(2).unwrap();
wf.complete(2).unwrap(); let seq = wf.wait(3).unwrap(); wf.wait_for_seq(seq).await.unwrap();
let state = wf.state();
let guard = state.read();
let intents = project_daemon_intents(&guard);
let by_name: std::collections::HashMap<String, &DaemonIntentUpdate> =
intents.iter().map(|u| (u.daemon.name.clone(), u)).collect();
assert_eq!(by_name["task/1"].intent, DaemonIntent::Run);
assert_eq!(
by_name["task/2"].intent,
DaemonIntent::Stop,
"terminal → Stop"
);
assert_eq!(
by_name["task/3"].intent,
DaemonIntent::Stop,
"Waiting → Stop"
);
assert_eq!(
by_name["task/4"].intent,
DaemonIntent::Stop,
"Submitted (never ran) → Stop (a harmless no-op at reconcile)",
);
assert!(
intents.iter().all(|u| u.node.is_none()),
"Projection 1 is claim-agnostic: every intent runs anywhere",
);
}
#[tokio::test]
async fn output_is_one_intent_per_task_in_stable_order() {
let redex = Redex::new();
let wf = WorkflowAdapter::open(&redex, 0x5C4E_DB02).await.unwrap();
for id in [42u64, 7, 99, 1] {
wf.submit(id).unwrap();
}
let seq = wf.start(7).unwrap();
wf.wait_for_seq(seq).await.unwrap();
let state = wf.state();
let guard = state.read();
let intents = project_daemon_intents(&guard);
assert_eq!(intents.len(), 4, "one intent per live task");
let ids: Vec<u64> = intents.iter().map(|u| u.daemon.id).collect();
let mut sorted = ids.clone();
sorted.sort_unstable();
assert_eq!(ids, sorted, "intents are emitted in stable daemon-id order");
}
#[test]
fn forced_placements_pin_claim_bearing_daemons_to_their_resolved_host() {
use crate::adapter::net::cortex::workflow::ActiveClaim;
let mut claims = ClaimRegistry::new();
claims.insert(daemon_ref(1), ActiveClaim { island: 0xA0 });
claims.insert(daemon_ref(2), ActiveClaim { island: 0xB0 });
let resolve = |island| match island {
0xA0 => Some(7),
0xB0 => Some(9),
_ => None,
};
let by_name: std::collections::HashMap<String, NodeId> =
project_forced_placements(&claims, resolve)
.into_iter()
.map(|(daemon, host)| (daemon.name, host))
.collect();
assert_eq!(by_name["task/1"], 7, "pinned to its island's host");
assert_eq!(by_name["task/2"], 9);
}
#[test]
fn forced_placement_skips_a_claim_whose_island_vanished() {
use crate::adapter::net::cortex::workflow::ActiveClaim;
let mut claims = ClaimRegistry::new();
claims.insert(daemon_ref(1), ActiveClaim { island: 0xDEAD });
let placements = project_forced_placements(&claims, |_| None);
assert!(
placements.is_empty(),
"a stale claim with no resolvable host produces no intent",
);
}
#[tokio::test]
async fn projection1_then_projection2_overlay_pins_only_the_claim_holder() {
use super::super::desired_daemon_intents;
use crate::adapter::net::behavior::meshos::state::DesiredState;
use crate::adapter::net::cortex::workflow::ActiveClaim;
let redex = Redex::new();
let wf = WorkflowAdapter::open(&redex, 0x5C4E_DB03).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 claims = ClaimRegistry::new();
claims.insert(daemon_ref(1), ActiveClaim { island: 0xA0 });
let resolve = |island| if island == 0xA0 { Some(7) } else { None };
let mut desired = DesiredState::default();
for intent in desired_daemon_intents(&wf.state().read(), &claims, resolve) {
desired.apply_daemon_intent(&intent);
}
assert_eq!(desired.desired_daemon_nodes.get(&daemon_ref(1)), Some(&7));
assert!(
!desired.desired_daemon_nodes.contains_key(&daemon_ref(2)),
"the non-claim task is never pinned",
);
}
#[tokio::test]
async fn running_claim_starts_its_daemon_on_the_claim_node_only() {
use super::super::desired_daemon_intents;
use crate::adapter::net::behavior::meshos::state::{DesiredState, MeshOsState};
use crate::adapter::net::behavior::meshos::{
reconcile, LocalityConfig, MaintenanceConfig, MeshOsAction, SchedulerConfig,
};
use crate::adapter::net::cortex::workflow::{ActiveClaim, WorkflowAdapter};
const CLAIM_HOST: NodeId = 7;
const OTHER_NODE: NodeId = 99;
let redex = Redex::new();
let wf = WorkflowAdapter::open(&redex, 0x5C4E_DB04).await.unwrap();
wf.submit(1).unwrap();
let seq = wf.start(1).unwrap(); wf.wait_for_seq(seq).await.unwrap();
let mut claims = ClaimRegistry::new();
claims.insert(daemon_ref(1), ActiveClaim { island: 0xA0 });
let mut desired = DesiredState::default();
{
let state = wf.state();
let guard = state.read();
for intent in desired_daemon_intents(&guard, &claims, |island| {
(island == 0xA0).then_some(CLAIM_HOST)
}) {
desired.apply_daemon_intent(&intent);
}
}
let actual = MeshOsState::default();
let (loc, maint, sched) = (
LocalityConfig::default(),
MaintenanceConfig::default(),
SchedulerConfig::default(),
);
let on_claim_host = reconcile(&actual, &desired, CLAIM_HOST, &loc, &maint, &sched, None);
assert!(
on_claim_host.iter().any(
|a| matches!(a, MeshOsAction::StartDaemon { daemon } if *daemon == daemon_ref(1)),
),
"the claim's node starts the pinned daemon; got {on_claim_host:?}",
);
let elsewhere = reconcile(&actual, &desired, OTHER_NODE, &loc, &maint, &sched, None);
assert!(
!elsewhere
.iter()
.any(|a| matches!(a, MeshOsAction::StartDaemon { .. })),
"a non-claim node never starts the pinned daemon; got {elsewhere:?}",
);
}
#[tokio::test]
async fn held_claim_on_a_non_running_task_merges_to_stop_pinned() {
use super::super::desired_daemon_intents;
use crate::adapter::net::cortex::workflow::ActiveClaim;
let redex = Redex::new();
let wf = WorkflowAdapter::open(&redex, 0x5C4E_DB08).await.unwrap();
wf.submit(1).unwrap();
wf.start(1).unwrap(); let seq = wf.complete(1).unwrap(); wf.wait_for_seq(seq).await.unwrap();
let mut claims = ClaimRegistry::new();
claims.insert(daemon_ref(1), ActiveClaim { island: 0xA0 });
let resolve = |island| (island == 0xA0).then_some(7);
let intents = desired_daemon_intents(&wf.state().read(), &claims, resolve);
let one = intents
.iter()
.find(|u| u.daemon == daemon_ref(1))
.expect("task 1 has an intent");
assert_eq!(
one.intent,
DaemonIntent::Stop,
"a non-Running held claim still projects Stop",
);
assert_eq!(one.node, Some(7), "but is still pinned to the claim host");
}
}