From 60c1f7622bb0b2eda5ea3673e517950efa8ed4e7 Mon Sep 17 00:00:00 2001 From: Willow C Reed Date: Thu, 27 Aug 2026 17:01:46 -0600 Subject: [PATCH 1/7] split scheduler_actor.rs into multipel files --- odorobo/src/actors/scheduler_actor.rs | 1728 +---------------- odorobo/src/actors/scheduler_actor/cache.rs | 310 +++ .../src/actors/scheduler_actor/discovery.rs | 387 ++++ .../src/actors/scheduler_actor/handlers.rs | 248 +++ .../src/actors/scheduler_actor/scheduling.rs | 371 ++++ odorobo/src/actors/scheduler_actor/tests.rs | 383 ++++ 6 files changed, 1758 insertions(+), 1669 deletions(-) create mode 100644 odorobo/src/actors/scheduler_actor/cache.rs create mode 100644 odorobo/src/actors/scheduler_actor/discovery.rs create mode 100644 odorobo/src/actors/scheduler_actor/handlers.rs create mode 100644 odorobo/src/actors/scheduler_actor/scheduling.rs create mode 100644 odorobo/src/actors/scheduler_actor/tests.rs diff --git a/odorobo/src/actors/scheduler_actor.rs b/odorobo/src/actors/scheduler_actor.rs index 1de4b01..505cbda 100644 --- a/odorobo/src/actors/scheduler_actor.rs +++ b/odorobo/src/actors/scheduler_actor.rs @@ -1,34 +1,41 @@ -use std::cmp::Ordering; +//! The scheduler actor and its in-memory cluster view. +//! +//! `SchedulerActor` makes placement decisions from periodically refreshed agent +//! status rather than issuing network requests while scoring. Its state is +//! eventually consistent and is not a durable source of truth: +//! +//! - [`SchedulerActor::vm_manifests`] records VM intent. +//! - [`SchedulerActor::vm_placements`] records desired and observed placements. +//! A VM can have more than one entry while it is migrating. +//! - [`SchedulerActor::vm_data_cache`] tracks discovered VM actors. A `None` +//! actor reference is a deliberate unresolved-placement placeholder. +//! - [`SchedulerActor::agent_vm_index`] is an observation index, not placement +//! authority. +//! +//! Implementation is split by responsibility into cache reconciliation, +//! discovery/polling, scheduling policy, and actor message handlers. + +mod cache; +mod discovery; +mod handlers; +mod scheduling; -use std::ops::ControlFlow; +#[cfg(test)] +mod tests; -use std::time::{Duration, Instant}; +use std::time::Instant; -use crate::actors::agent_actor::AgentActor; -use crate::ch_driver::actor::VMActor; -use crate::manifest::{ - AffinityRequirement, AffinityStrictness, AffinityType, MetadataTable, Operator, VmManifest, -}; -use crate::messages::agent::{AgentStatus, AgentStatusUpdate, GetAgentStatus, apply_status_update}; -use crate::messages::vm::{ - AgentListVMs, AgentListVMsReply, CreateVM, CreateVMReply, DeleteVM, DeleteVMReply, - GetConsoleHistory, GetConsoleHistoryReply, GetVMHeartbeat, GetVMInfo, GetVMInfoReply, - SendConsoleInput, SendConsoleInputReply, ShutdownVM, ShutdownVMReply, -}; -use crate::messages::{Ping, Pong}; -use crate::utils::actor_names::AGENT; -use crate::utils::actor_names::VM; -use crate::utils::actor_names::vm_actor_id; use ahash::{AHashMap, AHashSet}; use kameo::prelude::*; -use libp2p::futures::TryStreamExt; -use stable_eyre::eyre::OptionExt; -use stable_eyre::{Report, Result, eyre::eyre}; use tokio::task::JoinHandle; -use tracing::trace; -use tracing::{info, warn}; use ulid::Ulid; +use crate::actors::agent_actor::AgentActor; +use crate::ch_driver::actor::VMActor; +use crate::manifest::VmManifest; +use crate::messages::agent::{AgentStatus, AgentStatusUpdate}; +use crate::messages::vm::GetVMInfoReply; + #[derive(Debug)] struct VmActorDiscovered { actor_ref: RemoteActorRef, @@ -71,6 +78,10 @@ struct AgentUpdaterStopped { #[derive(Debug)] struct ReconcileVmPlacements; +/// The latest scheduler-approved status snapshot for an agent. +/// +/// The first accepted update must be a full snapshot. Later updates are +/// revision-ordered; stale or duplicate revisions are ignored. #[derive(Debug, Clone)] pub struct CachedAgentActor { pub actor_ref: RemoteActorRef, @@ -78,12 +89,21 @@ pub struct CachedAgentActor { pub status_revision: u64, } +/// Scheduler-side lifecycle for a VM placement. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum VmLifecycle { + /// Resources are reserved after create dispatch but before the destination + /// agent has reported the VM. This does not prove the VM process exists. Pending, + /// The destination agent has reported the VM in its status. Running, } +/// A desired or observed placement of a VM on an agent. +/// +/// Multiple entries for one VM are valid during migration. `last_confirmed_at` +/// is updated only from agent status; `created_at` is used to expire unresolved +/// pending placeholders. #[derive(Debug, Clone)] pub struct VmPlacement { pub agent_id: ActorId, @@ -92,1669 +112,39 @@ pub struct VmPlacement { pub last_confirmed_at: Option, } +/// A discovered VM actor associated with one concurrent placement. +/// +/// `actor_ref: None` intentionally represents a placement reserved before its +/// VM actor is discovered. #[derive(Debug, Clone)] pub struct CachedVMActor { pub actor_ref: Option>, } -type MetadataTables<'a> = [( - &'a std::collections::BTreeMap, - &'a std::collections::BTreeMap, -)]; - +/// An eventually consistent, in-memory VM scheduler. +/// +/// Public caches are exposed for inspection, but correlated maps must be +/// updated together through the actor's message handlers to preserve placement +/// and pending-resource accounting invariants. #[derive(RemoteActor)] pub struct SchedulerActor { /// Latest status snapshot for each known agent, used for scheduling decisions. pub agent_data_cache: AHashMap, - /// Heartbeat tasks that refresh the corresponding entry in `agent_data_cache`. + /// Polling tasks that refresh corresponding agent cache entries. pub agent_keepalive_tasks: AHashMap>, - - /// Maps discovered VM actor IDs to their canonical VM IDs. + /// Maps discovered VM actor IDs to canonical VM IDs. pub vm_actorid_ulid_map: AHashMap, - /// Canonical VM intent. A manifest remains available while the VM is being - /// reconciled or migrated. + /// Canonical VM intent retained while a VM is reconciled or migrated. pub vm_manifests: AHashMap, - /// Desired VM placements. Entries are retained when an agent observation says - /// a VM is absent so reconciliation can restore the desired placement. + /// Desired and observed VM placements; multiple entries allow migration. pub vm_placements: AHashMap>, - /// Cached VM actor references, with one entry per concurrent placement during - /// migration. `None` represents a placement whose actor is not discovered yet. + /// VM actor references. `None` marks a placement awaiting discovery. pub vm_data_cache: AHashMap>, - /// Heartbeat tasks that refresh the corresponding VM actor cache entry. + /// Polling tasks that refresh corresponding VM actor cache entries. pub vm_keepalive_tasks: AHashMap>, - /// Cached resource totals for pending placements; invalidated on state changes. pending_resources_cache: Option>, - /// VM membership observed from each agent, used to discover placements without - /// scanning every placement during scheduling. agent_vm_index: AHashMap>, - /// Whether a discovered actor is an agent or VM actor. actor_kinds: AHashMap, - - /// Background task that discovers actors and periodically triggers cleanup. + /// Background discovery and reconciliation task. pub cache_actor_finder: Option>, } - -// todo: this might need to be a runtime thing but this makes it easy to write for now and could easily be switched out later. -static VCPU_OVERPROVISIONMENT_NUMERATOR: u32 = 2; -static VCPU_OVERPROVISIONMENT_DENOMINATOR: u32 = 1; -const UNRESOLVED_VM_CACHE_TIMEOUT: Duration = Duration::from_secs(30); - -impl SchedulerActor { - fn shrink_non_migrating_entries(entries: &mut Vec) { - if entries.len() <= 1 && entries.capacity() >= 4 { - entries.shrink_to_fit(); - } - } - - fn update_cached_vm_entry( - entries: &mut Vec, - actor_id: ActorId, - cached_vm: CachedVMActor, - ) { - if let Some(entry) = entries.iter_mut().find(|entry| { - entry - .actor_ref - .as_ref() - .is_some_and(|actor| actor.id() == actor_id) - }) { - *entry = cached_vm; - } else if let Some(entry) = entries.iter_mut().find(|entry| entry.actor_ref.is_none()) { - *entry = cached_vm; - } else { - entries.push(cached_vm); - } - Self::shrink_non_migrating_entries(entries); - } - - fn cleanup_unresolved_vm_cache( - manifests: &mut AHashMap, - placements: &mut AHashMap>, - data_cache: &mut AHashMap>, - ) { - let now = Instant::now(); - let empty_vmids: Vec<_> = placements - .iter_mut() - .filter_map(|(vmid, entries)| { - entries.retain(|entry| { - entry.lifecycle != VmLifecycle::Pending - || now.duration_since(entry.created_at) < UNRESOLVED_VM_CACHE_TIMEOUT - }); - Self::shrink_non_migrating_entries(entries); - entries.is_empty().then_some(*vmid) - }) - .collect(); - - for vmid in empty_vmids { - Self::remove_vm_state(vmid, manifests, placements, data_cache); - } - } - - /// Returns every VM that could be on an agent, deduplicating each source in - /// constant expected time. - fn placement_vm_ids( - placements: &AHashMap>, - indexed: Option<&AHashSet>, - agent_id: ActorId, - observed: &[Ulid], - ) -> Vec { - let indexed_len = indexed.map_or(0, |index| index.len()); - let mut vmids = Vec::with_capacity(observed.len().max(indexed_len)); - let mut seen = AHashSet::with_capacity(observed.len().saturating_add(indexed_len)); - - for vmid in observed { - if seen.insert(*vmid) { - vmids.push(*vmid); - } - } - if let Some(indexed) = indexed { - for vmid in indexed { - if seen.insert(*vmid) { - vmids.push(*vmid); - } - } - } - for vmid in placements.iter().filter_map(|(vmid, entries)| { - entries - .iter() - .any(|entry| entry.agent_id == agent_id) - .then_some(vmid) - }) { - if seen.insert(*vmid) { - vmids.push(*vmid); - } - } - vmids - } - - fn remove_vm_state( - vmid: Ulid, - manifests: &mut AHashMap, - placements: &mut AHashMap>, - data_cache: &mut AHashMap>, - ) { - manifests.remove(&vmid); - placements.remove(&vmid); - data_cache.remove(&vmid); - } - - fn remove_vm_actor(actor_id: ActorId, data_cache: &mut AHashMap>) { - let empty_vmids: Vec<_> = data_cache - .iter_mut() - .filter_map(|(vmid, entries)| { - entries.retain(|entry| { - entry - .actor_ref - .as_ref() - .is_none_or(|actor| actor.id() != actor_id) - }); - Self::shrink_non_migrating_entries(entries); - entries.is_empty().then_some(*vmid) - }) - .collect(); - for vmid in empty_vmids { - data_cache.remove(&vmid); - } - } - - fn remove_agent_placements( - agent_id: ActorId, - manifests: &mut AHashMap, - placements: &mut AHashMap>, - data_cache: &mut AHashMap>, - ) { - let empty_vmids: Vec<_> = placements - .iter_mut() - .filter_map(|(vmid, entries)| { - entries.retain(|entry| entry.agent_id != agent_id); - Self::shrink_non_migrating_entries(entries); - entries.is_empty().then_some(*vmid) - }) - .collect(); - - for vmid in empty_vmids { - Self::remove_vm_state(vmid, manifests, placements, data_cache); - } - } - - fn rollback_failed_create( - vmid: Ulid, - actor_exists: bool, - actor_id: Option, - actor_map: &mut AHashMap, - manifests: &mut AHashMap, - placements: &mut AHashMap>, - data_cache: &mut AHashMap>, - ) { - if !actor_exists { - if let Some(actor_id) = actor_id - && actor_map.get(&actor_id) == Some(&vmid) - { - actor_map.remove(&actor_id); - } - Self::remove_vm_state(vmid, manifests, placements, data_cache); - } - } - - fn cleanup_agent_actor(&mut self, actor_id: ActorId) { - if let Some(keepalive_task) = self.agent_keepalive_tasks.remove(&actor_id) { - trace!(?actor_id, "Aborting agent keepalive task"); - keepalive_task.abort(); - } - self.agent_data_cache.remove(&actor_id); - self.agent_vm_index.remove(&actor_id); - self.invalidate_pending_resources(); - Self::remove_agent_placements( - actor_id, - &mut self.vm_manifests, - &mut self.vm_placements, - &mut self.vm_data_cache, - ); - } - - fn cleanup_vm_actor(&mut self, actor_id: ActorId) { - if let Some(keepalive_task) = self.vm_keepalive_tasks.remove(&actor_id) { - trace!(?actor_id, "Aborting VM keepalive task"); - keepalive_task.abort(); - } - let vmid = self.vm_actorid_ulid_map.remove(&actor_id); - self.invalidate_pending_resources(); - Self::remove_vm_actor(actor_id, &mut self.vm_data_cache); - if let Some(vmid) = vmid - && self - .vm_data_cache - .get(&vmid) - .is_none_or(|entries| entries.iter().all(|entry| entry.actor_ref.is_none())) - { - Self::remove_vm_state( - vmid, - &mut self.vm_manifests, - &mut self.vm_placements, - &mut self.vm_data_cache, - ); - } - } - - fn reconcile_agent_delta( - agent_id: ActorId, - added: &[Ulid], - _removed: &[Ulid], - manifests: &AHashMap, - placements: &mut AHashMap>, - ) { - let now = Instant::now(); - for vmid in added { - if !manifests.contains_key(vmid) { - continue; - } - let entries = placements.entry(*vmid).or_default(); - if let Some(entry) = entries.iter_mut().find(|entry| entry.agent_id == agent_id) { - entry.lifecycle = VmLifecycle::Running; - entry.last_confirmed_at = Some(now); - } else { - entries.push(VmPlacement { - agent_id, - lifecycle: VmLifecycle::Running, - created_at: now, - last_confirmed_at: Some(now), - }); - } - } - - // A removal is an observation about the agent, not a change to the - // scheduler's desired state. Keep the placement so reconciliation can - // schedule the VM again. The full status path performs the same - // distinction for snapshots. - } - - fn reconcile_agent_placements( - agent_id: ActorId, - status: &AgentStatus, - manifests: &AHashMap, - placements: &mut AHashMap>, - ) { - let now = Instant::now(); - let observed: AHashSet<_> = status.vms.iter().copied().collect(); - let missing_agent_placements: Vec<_> = observed - .iter() - .filter(|vmid| manifests.contains_key(vmid)) - .filter(|vmid| { - placements - .get(vmid) - .is_none_or(|entries| !entries.iter().any(|entry| entry.agent_id == agent_id)) - }) - .copied() - .collect(); - for vmid in missing_agent_placements { - placements.entry(vmid).or_default().push(VmPlacement { - agent_id, - lifecycle: VmLifecycle::Running, - created_at: now, - last_confirmed_at: Some(now), - }); - } - - let empty_vmids: Vec<_> = placements - .iter_mut() - .filter_map(|(vmid, entries)| { - for entry in entries - .iter_mut() - .filter(|entry| entry.agent_id == agent_id) - { - if observed.contains(vmid) { - entry.lifecycle = VmLifecycle::Running; - entry.last_confirmed_at = Some(now); - } - } - Self::shrink_non_migrating_entries(entries); - entries.is_empty().then_some(*vmid) - }) - .collect(); - for vmid in empty_vmids { - placements.remove(&vmid); - } - } - - #[expect(dead_code, reason = "reserved for explicit placement by actor id")] - fn lookup_agent_by_actor_id(&self, actor_id: &ActorId) -> Option> { - self.agent_data_cache - .get(actor_id) - .map(|data| data.actor_ref.clone()) - } - - #[expect(dead_code, reason = "reserved for explicit placement by hostname")] - fn lookup_agent_by_hostname(&self, hostname: &str) -> Option> { - self.agent_data_cache - .values() - .find(|data| data.data.hostname == hostname) - .map(|data| data.actor_ref.clone()) - } - - async fn vm_actor_finder(parent_actor_ref: ActorRef) -> Result<(), Report> { - trace!("running vm_actor_finder"); - - let mut vm_actor_stream = RemoteActorRef::::lookup_all(VM); - - while let Some(vm_actor) = vm_actor_stream.try_next().await? { - parent_actor_ref - .tell(VmActorDiscovered { - actor_ref: vm_actor, - }) - .send() - .await?; - } - - Ok(()) - } - - async fn vm_updater_task(scheduler: ActorRef, actor_ref: RemoteActorRef) { - let mut interval = tokio::time::interval(Duration::from_secs(1)); - let mut fails: u8 = 0; - let mut initialized = false; - - loop { - if !initialized { - if let Ok(data) = actor_ref.ask(&GetVMInfo { vmid: None }).await { - let send_result = scheduler - .tell(VmUpdated { - actor_ref: actor_ref.clone(), - data, - }) - .send() - .await - .map_err(|error| eyre!("failed to send VM update: {error}")); - if let Err(error) = send_result { - warn!(?error, "VM updater could not notify scheduler"); - return; - } - initialized = true; - fails = 0; - } else { - fails = fails.saturating_add(1); - } - } else if actor_ref.ask(&GetVMHeartbeat).await.is_ok() { - fails = 0; - } else { - fails = fails.saturating_add(1); - } - - if fails > 5 { - warn!( - ?actor_ref, - "can no longer reach vm actor, cleaning up cache entries" - ); - - let send_result = scheduler - .tell(VmUpdaterStopped { - actor_id: actor_ref.id(), - }) - .send() - .await - .map_err(|error| eyre!("failed to send VM stop: {error}")); - if let Err(error) = send_result { - warn!(?error, "VM updater could not notify scheduler"); - } - return; - } - - interval.tick().await; - } - } - - async fn agent_actor_finder(parent_actor_ref: ActorRef) -> Result<(), Report> { - trace!("running agent_actor_finder"); - - let mut agent_actor_stream = RemoteActorRef::::lookup_all(AGENT); - - while let Some(agent_actor) = agent_actor_stream.try_next().await? { - parent_actor_ref - .tell(AgentActorDiscovered { - actor_ref: agent_actor, - }) - .send() - .await?; - } - - Ok(()) - } - - async fn agent_updater_task(scheduler: ActorRef, actor_ref: RemoteActorRef) { - let mut interval = tokio::time::interval(Duration::from_secs(1)); - let mut status_revision = 0; - let mut initial_status = true; - let mut fails: u8 = 0; - loop { - if let Ok(update) = actor_ref - .ask(&GetAgentStatus { - since_revision: status_revision, - initial: initial_status, - }) - .await - { - status_revision = match &update { - AgentStatusUpdate::Full { revision, .. } - | AgentStatusUpdate::Delta { revision, .. } => *revision, - }; - initial_status = false; - let send_result = scheduler - .tell(AgentUpdated { - actor_id: actor_ref.id(), - actor_ref: actor_ref.clone(), - update, - }) - .send() - .await - .map_err(|error| eyre!("failed to send agent update: {error}")); - if let Err(error) = send_result { - warn!(?error, "agent updater could not notify scheduler"); - return; - } - fails = 0; - } else { - fails = fails.saturating_add(1); - } - - if fails > 5 { - warn!( - ?actor_ref, - "can no longer reach agent actor, stopping updater" - ); - let send_result = scheduler - .tell(AgentUpdaterStopped { - actor_id: actor_ref.id(), - }) - .send() - .await - .map_err(|error| eyre!("failed to send agent stop: {error}")); - if let Err(error) = send_result { - warn!(?error, "agent updater could not notify scheduler"); - } - return; - } - - interval.tick().await; - } - } - - fn start_actor_finder(&mut self, actor_ref: ActorRef) { - self.cache_actor_finder = Some(tokio::spawn(async move { - let mut interval = tokio::time::interval(Duration::from_secs(5)); - loop { - let vm_result = Self::vm_actor_finder(actor_ref.clone()).await; - let agent_result = Self::agent_actor_finder(actor_ref.clone()).await; - if let Err(error) = vm_result { - warn!(?error, "VM actor discovery failed"); - } - if let Err(error) = agent_result { - warn!(?error, "agent actor discovery failed"); - } - actor_ref.tell(ReconcileVmPlacements).send().await.ok(); - interval.tick().await; - } - })); - } - - fn pending_resources(&mut self) -> &AHashMap { - if self.pending_resources_cache.is_none() { - self.pending_resources_cache = Some(pending_resources_by_agent( - &self.vm_manifests, - &self.vm_placements, - )); - } - self.pending_resources_cache - .as_ref() - .expect("pending resources cache was just initialized") - } - - fn invalidate_pending_resources(&mut self) { - self.pending_resources_cache = None; - } - - /// Determine the best agent to schedule a specific VM creation request to. - /// - /// The scheduler first filters agents by capacity and affinity requirements, - /// then uses affinity and general resource scores to select the best match. - /// Affinity rules are roughly based on - /// . - fn schedule_agent(&mut self, msg: &CreateVM) -> Result, Report> { - self.pending_resources(); - let pending_resources = self - .pending_resources_cache - .as_ref() - .expect("pending resources cache was just initialized"); - let mut best_agent = None; - let mut best_score = AgentScore::REJECTED; - - for agent in self.agent_data_cache.values() { - let score = self.score_agent(msg, agent, pending_resources); - - if score > best_score { - best_agent = Some(agent.actor_ref.clone()); - best_score = score; - } - } - - info!(?best_score, ?best_agent, "best agent"); - - best_agent.ok_or_eyre("No valid agents found.") - } - - #[expect(dead_code, reason = "reserved for a future batch create message")] - fn schedule_agents( - &mut self, - msgs: &[CreateVM], - ) -> Vec, Report>> { - self.pending_resources(); - let pending_resources = self - .pending_resources_cache - .as_ref() - .expect("pending resources cache was just initialized"); - msgs.iter() - .map(|msg| { - let mut best_agent = None; - let mut best_score = AgentScore::REJECTED; - for agent in self.agent_data_cache.values() { - let score = self.score_agent(msg, agent, pending_resources); - if score > best_score { - best_agent = Some(agent.actor_ref.clone()); - best_score = score; - } - } - best_agent.ok_or_eyre("No valid agents found.") - }) - .collect() - } - - // this function intentionally only checks against the cache. this has some positives and negatives: - // positive: it will never trigger any network requests so its very fast, and having to do network requests for scoring whenever we want to schedule a vm is likely a bad idea - // negative: it technically has a delayed view of the cluster, meaning that some things that happened in the future, may not exist yet. so we need to be careful about how this is done so affinity rules are not accidentally broken. mostly this means, if we do anything that could affect the outcome of an affinity rule (ex: network request to an agent), we need to update the cache, before we do the action. - fn score_agent( - &self, - msg: &CreateVM, - agent: &CachedAgentActor, - pending_resources: &AHashMap, - ) -> AgentScore { - let mut score = AgentScore::default(); - - let agent_max_vcpus = agent - .data - .vcpus - .saturating_mul(VCPU_OVERPROVISIONMENT_NUMERATOR) - .checked_div(VCPU_OVERPROVISIONMENT_DENOMINATOR) - .unwrap_or(u32::MAX); - // todo: do we care about VMData.max_vcpus? - let (pending_vcpus, pending_ram) = pending_resources - .get(&agent.actor_ref.id()) - .copied() - .unwrap_or_default(); - let used_vcpus = agent.data.used_vcpus.saturating_add(pending_vcpus); - let used_ram = agent.data.used_ram.as_u64().saturating_add(pending_ram); - let requested_vcpus = msg.config.desired.compute.vcpus; - let requested_memory = msg.config.desired.compute.memory_bytes; - let agent_used_vcpus = used_vcpus.saturating_add(requested_vcpus); - - if !has_capacity( - agent_max_vcpus, - used_vcpus, - requested_vcpus, - agent.data.ram.as_u64(), - used_ram, - requested_memory, - ) { - return AgentScore::REJECTED; - } - - #[expect( - clippy::cast_precision_loss, - reason = "the scheduler score intentionally uses f32 ratios" - )] - #[expect( - clippy::arithmetic_side_effects, - reason = "the preceding capacity check guarantees non-negative subtraction" - )] - let vcpu_headroom = (agent_max_vcpus - agent_used_vcpus) as f32 / agent_max_vcpus as f32; - score.general += vcpu_headroom; - - // todo: add ram overprovisionment. not adding this to scheduler until it works on the hypervisor side. - let agent_max_ram = agent.data.ram; - let agent_used_ram = bytesize::ByteSize::b(used_ram.saturating_add(requested_memory)); - - #[expect( - clippy::cast_precision_loss, - reason = "the scheduler score intentionally uses f32 ratios" - )] - let ram_headroom = agent_max_ram - .as_u64() - .saturating_sub(agent_used_ram.as_u64()) as f32 - / agent_max_ram.as_u64() as f32; - score.general += ram_headroom; - - // Roughly based on . - if !msg.config.desired.placement.affinity.is_empty() { - let affinity_rules = &msg.config.desired.placement.affinity; - for rule in affinity_rules { - let mut metadata_tables = Vec::with_capacity(1); - match rule.affinity_type { - AffinityType::VirtualMachine => { - metadata_tables.extend( - Self::placement_vm_ids( - &self.vm_placements, - self.agent_vm_index.get(&agent.actor_ref.id()), - agent.actor_ref.id(), - &agent.data.vms, - ) - .into_iter() - .filter_map(|vmid| self.vm_manifests.get(&vmid)) - .map(|manifest| { - ( - &manifest.desired.metadata.labels, - &manifest.desired.metadata.annotations, - ) - }), - ); - } - AffinityType::Agent => { - metadata_tables.push(( - &agent.data.metadata.labels, - &agent.data.metadata.annotations, - )); - } - } - - let follows_rule = evaluate_affinity_rule(&metadata_tables, rule); - - let Some(affinity_delta) = affinity_delta(&rule.strictness, follows_rule) else { - return AgentScore::REJECTED; - }; - score.affinity = score.affinity.saturating_add(affinity_delta); - } - } - - // todo (future): possibly keep a percent of agents completely empty, to be able to be converted to dedis automatically. - // they would have their agent score set to like f32::MIN, so they can be scheduled to if there is no other available agents. - // rough pseudo code to implement this: - // if agent.metadata.vms.len() == 0 && hash(agent.config.hostname) % total_chance < threshold { - // agent_score = 1; - // } - - score - } -} - -const fn has_capacity( - max_vcpus: u32, - used_vcpus: u32, - requested_vcpus: u32, - max_ram: u64, - used_ram: u64, - requested_ram: u64, -) -> bool { - used_vcpus.saturating_add(requested_vcpus) <= max_vcpus - && used_ram.saturating_add(requested_ram) <= max_ram -} - -fn pending_resources_by_agent( - manifests: &AHashMap, - placements: &AHashMap>, -) -> AHashMap { - let mut resources = AHashMap::new(); - for (vmid, entries) in placements { - let Some(manifest) = manifests.get(vmid) else { - continue; - }; - for entry in entries { - if entry.lifecycle == VmLifecycle::Pending { - let totals = resources.entry(entry.agent_id).or_insert((0u32, 0u64)); - totals.0 = totals.0.saturating_add(manifest.desired.compute.vcpus); - totals.1 = totals - .1 - .saturating_add(manifest.desired.compute.memory_bytes); - } - } - } - resources -} - -#[cfg(test)] -fn pending_resources_for_agent( - manifests: &AHashMap, - placements: &AHashMap>, - agent_id: ActorId, -) -> (u32, u64) { - pending_resources_by_agent(manifests, placements) - .get(&agent_id) - .copied() - .unwrap_or_default() -} - -fn affinity_delta(strictness: &AffinityStrictness, follows_rule: bool) -> Option { - match strictness { - AffinityStrictness::Required if !follows_rule => None, - AffinityStrictness::Required => Some(0), - AffinityStrictness::Preferred { weight } => { - Some(i64::from(follows_rule).saturating_mul(*weight)) - } - } -} - -fn evaluate_affinity_rule( - metadata_tables: &MetadataTables<'_>, - rule: &crate::manifest::AffinityRule, -) -> bool { - let mut follows_rule = false; - - for requirement in &rule.requirements { - let mut requirement_outcome = !metadata_tables.is_empty(); - - for object_metadata in metadata_tables { - let table = match requirement.table { - MetadataTable::Label => object_metadata.0, - MetadataTable::Annotation => object_metadata.1, - }; - - if !evaluate_table_value(table.get(&requirement.key), requirement) { - requirement_outcome = false; - break; - } - } - - if requirement_outcome { - follows_rule = true; - break; - } - } - - follows_rule ^ rule.inverse -} - -fn evaluate_table_value(value_option: Option<&String>, requirement: &AffinityRequirement) -> bool { - let Some(value) = value_option else { - return matches!(requirement.operator, Operator::NotIn); - }; - - match requirement.operator { - Operator::In => requirement.values.contains(value), - Operator::NotIn => !requirement.values.contains(value), - Operator::Lt | Operator::Gt => { - let [requirement_value] = &requirement.values[..] else { - return false; - }; - - let Ok(value_number): Result = value.parse() else { - return false; - }; - - let Ok(requirement_value_number): Result = requirement_value.parse() else { - return false; - }; - - if requirement.operator == Operator::Lt { - value_number < requirement_value_number - } else { - value_number > requirement_value_number - } - } - } -} - -#[cfg(test)] -mod tests { - use super::{ - CachedActorKind, CachedVMActor, SchedulerActor, VmLifecycle, VmPlacement, affinity_delta, - evaluate_affinity_rule, evaluate_table_value, has_capacity, pending_resources_for_agent, - }; - - use crate::manifest::{ - AffinityRequirement, AffinityRule, AffinityStrictness, AffinityType, Boot, Compute, - DesiredState, Metadata, MetadataTable, Operator, VmManifest, - }; - use crate::messages::agent::AgentStatus; - use crate::types::ObjectMetadata; - use ahash::AHashMap; - use bytesize::ByteSize; - use std::collections::BTreeMap; - use std::time::{Duration, Instant}; - use ulid::Ulid; - - fn test_manifest(vcpus: u32, memory_bytes: u64) -> VmManifest { - VmManifest { - api_version: crate::manifest::MANIFEST_VERSION, - id: Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"), - desired: DesiredState { - metadata: Metadata { - name: "test".to_owned(), - ..Default::default() - }, - compute: Compute { - vcpus, - memory_bytes, - ..Default::default() - }, - boot: Boot::default(), - ..Default::default() - }, - observed: None, - } - } - - fn requirement(operator: Operator, values: &[&str]) -> AffinityRequirement { - AffinityRequirement { - key: "tier".to_owned(), - table: MetadataTable::Label, - operator, - values: values.iter().map(|value| (*value).to_owned()).collect(), - } - } - - #[test] - fn removes_expired_unresolved_vm_placeholders() { - let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); - let agent_id = super::ActorId::new(1); - let mut placements: AHashMap> = AHashMap::new(); - placements.insert( - vmid, - vec![VmPlacement { - agent_id, - lifecycle: VmLifecycle::Pending, - created_at: Instant::now() - .checked_sub(Duration::from_secs(31)) - .expect("test timestamp should be representable"), - last_confirmed_at: None, - }], - ); - let mut manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); - let mut data_cache: AHashMap> = AHashMap::new(); - - SchedulerActor::cleanup_unresolved_vm_cache( - &mut manifests, - &mut placements, - &mut data_cache, - ); - - assert!(!manifests.contains_key(&vmid)); - assert!(!placements.contains_key(&vmid)); - assert!(!data_cache.contains_key(&vmid)); - } - - #[test] - fn reserves_pending_vm_resources_until_agent_status_confirms_them() { - let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); - let agent_id = super::ActorId::new(1); - let config = test_manifest(4, ByteSize::gib(8).as_u64()); - let mut placements: AHashMap> = AHashMap::new(); - placements.insert( - vmid, - vec![VmPlacement { - agent_id, - lifecycle: VmLifecycle::Pending, - created_at: Instant::now(), - last_confirmed_at: None, - }], - ); - - let manifests = AHashMap::from([(vmid, config)]); - assert_eq!( - pending_resources_for_agent(&manifests, &placements, agent_id), - (4, ByteSize::gib(8).as_u64()) - ); - } - - #[test] - fn reconciling_source_agent_preserves_destination_migration_placement() { - let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); - let source_agent = super::ActorId::new(1); - let destination_agent = super::ActorId::new(2); - let running_placement = |agent_id| VmPlacement { - agent_id, - lifecycle: VmLifecycle::Running, - created_at: Instant::now(), - last_confirmed_at: Some(Instant::now()), - }; - let manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); - let mut placements = AHashMap::from([(vmid, vec![running_placement(source_agent)])]); - let destination_status = AgentStatus { - hostname: "destination".to_owned(), - vcpus: 1, - ram: ByteSize::b(1), - used_vcpus: 0, - used_ram: ByteSize::b(0), - vms: vec![vmid], - metadata: ObjectMetadata::default(), - }; - SchedulerActor::reconcile_agent_placements( - destination_agent, - &destination_status, - &manifests, - &mut placements, - ); - assert_eq!(manifests.len(), 1); - - let source_status = AgentStatus { - hostname: "source".to_owned(), - vcpus: 1, - ram: ByteSize::b(1), - used_vcpus: 0, - used_ram: ByteSize::b(0), - vms: Vec::new(), - metadata: ObjectMetadata::default(), - }; - - SchedulerActor::reconcile_agent_placements( - source_agent, - &source_status, - &manifests, - &mut placements, - ); - - let remaining = placements - .get(&vmid) - .expect("destination placement remains"); - assert_eq!(remaining.len(), 2); - assert!(remaining.iter().any(|entry| entry.agent_id == source_agent)); - assert!( - remaining - .iter() - .any(|entry| entry.agent_id == destination_agent) - ); - } - - #[test] - fn agent_removal_preserves_desired_placement() { - let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); - let agent_id = super::ActorId::new(1); - let manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); - let mut placements = AHashMap::from([( - vmid, - vec![VmPlacement { - agent_id, - lifecycle: VmLifecycle::Running, - created_at: Instant::now(), - last_confirmed_at: Some(Instant::now()), - }], - )]); - - SchedulerActor::reconcile_agent_delta(agent_id, &[], &[vmid], &manifests, &mut placements); - - assert_eq!(placements[&vmid].len(), 1); - assert_eq!(placements[&vmid][0].agent_id, agent_id); - } - - #[test] - fn failed_create_rolls_back_state_without_an_actor() { - let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); - let mut manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); - let mut actor_map = AHashMap::new(); - let mut placements = AHashMap::from([( - vmid, - vec![VmPlacement { - agent_id: super::ActorId::new(1), - lifecycle: VmLifecycle::Pending, - created_at: Instant::now(), - last_confirmed_at: None, - }], - )]); - let mut data_cache = AHashMap::from([(vmid, vec![CachedVMActor { actor_ref: None }])]); - - SchedulerActor::rollback_failed_create( - vmid, - false, - None, - &mut actor_map, - &mut manifests, - &mut placements, - &mut data_cache, - ); - - assert!(!manifests.contains_key(&vmid)); - assert!(!placements.contains_key(&vmid)); - assert!(!data_cache.contains_key(&vmid)); - } - - #[test] - fn failed_create_keeps_state_if_actor_exists() { - let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); - let mut manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); - let mut actor_map = AHashMap::new(); - let mut placements = AHashMap::new(); - let mut data_cache = AHashMap::new(); - - SchedulerActor::rollback_failed_create( - vmid, - true, - None, - &mut actor_map, - &mut manifests, - &mut placements, - &mut data_cache, - ); - - assert!(manifests.contains_key(&vmid)); - } - - #[test] - fn agent_cleanup_does_not_remove_unrelated_vm_state() { - let agent_id = super::ActorId::new(1); - let vm_actor_id = super::ActorId::new(2); - let placement_agent_id = super::ActorId::new(3); - let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); - let mut scheduler = SchedulerActor { - agent_data_cache: AHashMap::new(), - agent_keepalive_tasks: AHashMap::new(), - vm_actorid_ulid_map: AHashMap::from([(vm_actor_id, vmid)]), - vm_manifests: AHashMap::from([(vmid, test_manifest(1, 1))]), - vm_placements: AHashMap::from([( - vmid, - vec![VmPlacement { - agent_id: placement_agent_id, - lifecycle: VmLifecycle::Running, - created_at: Instant::now(), - last_confirmed_at: Some(Instant::now()), - }], - )]), - vm_data_cache: AHashMap::from([(vmid, vec![CachedVMActor { actor_ref: None }])]), - vm_keepalive_tasks: AHashMap::new(), - pending_resources_cache: None, - agent_vm_index: AHashMap::new(), - actor_kinds: AHashMap::from([(agent_id, CachedActorKind::Agent)]), - cache_actor_finder: None, - }; - - scheduler.cleanup_agent_actor(agent_id); - - assert!(scheduler.vm_placements.contains_key(&vmid)); - assert!(scheduler.vm_manifests.contains_key(&vmid)); - assert!(scheduler.vm_actorid_ulid_map.contains_key(&vm_actor_id)); - assert!(scheduler.vm_data_cache.contains_key(&vmid)); - } - - #[test] - fn vm_cleanup_does_not_remove_unrelated_agent_state() { - let agent_id = super::ActorId::new(1); - let vm_actor_id = super::ActorId::new(2); - let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); - let mut scheduler = SchedulerActor { - agent_data_cache: AHashMap::new(), - agent_keepalive_tasks: AHashMap::new(), - vm_actorid_ulid_map: AHashMap::from([(vm_actor_id, vmid)]), - vm_manifests: AHashMap::from([(vmid, test_manifest(1, 1))]), - vm_placements: AHashMap::new(), - vm_data_cache: AHashMap::from([(vmid, vec![CachedVMActor { actor_ref: None }])]), - vm_keepalive_tasks: AHashMap::new(), - pending_resources_cache: None, - agent_vm_index: AHashMap::new(), - actor_kinds: AHashMap::from([(agent_id, CachedActorKind::Agent)]), - cache_actor_finder: None, - }; - - scheduler.cleanup_vm_actor(vm_actor_id); - - assert!(scheduler.actor_kinds.contains_key(&agent_id)); - assert!(!scheduler.vm_manifests.contains_key(&vmid)); - assert!(!scheduler.vm_data_cache.contains_key(&vmid)); - } - - #[test] - fn evaluates_membership_and_missing_keys() { - let metadata = BTreeMap::from([("tier".to_owned(), "frontend".to_owned())]); - assert!(evaluate_table_value( - metadata.get("tier"), - &requirement(Operator::In, &["frontend", "api"]) - )); - assert!(!evaluate_table_value( - metadata.get("tier"), - &requirement(Operator::In, &["backend"]) - )); - assert!(!evaluate_table_value( - None, - &requirement(Operator::In, &["frontend"]) - )); - assert!(evaluate_table_value( - None, - &requirement(Operator::NotIn, &["frontend"]) - )); - } - - #[test] - fn evaluates_not_in_and_numeric_comparisons() { - let metadata = BTreeMap::from([("tier".to_owned(), "4".to_owned())]); - assert!(evaluate_table_value( - metadata.get("tier"), - &requirement(Operator::NotIn, &["5"]) - )); - assert!(evaluate_table_value( - metadata.get("tier"), - &requirement(Operator::Lt, &["5"]) - )); - assert!(evaluate_table_value( - metadata.get("tier"), - &requirement(Operator::Gt, &["3"]) - )); - assert!(!evaluate_table_value( - metadata.get("tier"), - &requirement(Operator::Lt, &["4", "5"]) - )); - assert!(!evaluate_table_value( - metadata.get("tier"), - &requirement(Operator::Gt, &["not-a-number"]) - )); - } - - #[test] - fn evaluates_inverse_and_empty_requirements() { - let metadata = ObjectMetadata { - labels: BTreeMap::from([("tier".to_owned(), "frontend".to_owned())]), - annotations: BTreeMap::new(), - }; - let rule = AffinityRule { - strictness: AffinityStrictness::Required, - affinity_type: AffinityType::Agent, - inverse: true, - requirements: vec![requirement(Operator::In, &["frontend"])], - }; - assert!(!evaluate_affinity_rule( - &[(&metadata.labels, &metadata.annotations)], - &rule, - )); - - let non_matching_rule = AffinityRule { - strictness: AffinityStrictness::Required, - affinity_type: AffinityType::Agent, - inverse: true, - requirements: vec![requirement(Operator::In, &["backend"])], - }; - assert!(evaluate_affinity_rule( - &[(&metadata.labels, &metadata.annotations)], - &non_matching_rule, - )); - - let empty_rule = AffinityRule { - strictness: AffinityStrictness::Required, - affinity_type: AffinityType::Agent, - inverse: false, - requirements: Vec::new(), - }; - assert!(!evaluate_affinity_rule(&[], &empty_rule)); - } - - #[test] - fn evaluates_required_preferred_and_capacity_rules() { - assert_eq!(affinity_delta(&AffinityStrictness::Required, true), Some(0)); - assert_eq!(affinity_delta(&AffinityStrictness::Required, false), None); - assert_eq!( - affinity_delta(&AffinityStrictness::Preferred { weight: 7 }, true), - Some(7) - ); - assert_eq!( - affinity_delta(&AffinityStrictness::Preferred { weight: 7 }, false), - Some(0) - ); - assert!(has_capacity(8, 2, 2, 16, 4, 4)); - assert!(has_capacity(8, 6, 2, 16, 4, 4)); - assert!(has_capacity(8, 2, 2, 16, 12, 4)); - assert!(!has_capacity(8, 7, 2, 16, 4, 4)); - assert!(!has_capacity(8, 2, 2, 16, 13, 4)); - } -} - -#[derive(Debug, Clone, Copy, PartialEq)] -struct AgentScore { - general: f32, - affinity: i64, -} - -impl AgentScore { - pub const REJECTED: Self = Self { - general: f32::NEG_INFINITY, - affinity: i64::MIN, - }; -} - -impl Default for AgentScore { - fn default() -> Self { - Self { - general: 0.0, - affinity: 0, - } - } -} - -impl PartialOrd for AgentScore { - fn partial_cmp(&self, other: &Self) -> Option { - let affinity_cmp = self.affinity.cmp(&other.affinity); - - if affinity_cmp != Ordering::Equal { - return Some(affinity_cmp); - } - - self.general.partial_cmp(&other.general) - } -} - -#[allow(clippy::unused_async_trait_impl)] -impl Actor for SchedulerActor { - type Args = (); - type Error = Report; - - async fn on_start(_state: Self::Args, actor_ref: ActorRef) -> Result { - let peer_id = *actor_ref.id().peer_id().unwrap(); - - info!(?peer_id, "Scheduler Actor started!"); - - let mut scheduler_actor = Self { - agent_data_cache: AHashMap::new(), - agent_keepalive_tasks: AHashMap::new(), - vm_actorid_ulid_map: AHashMap::new(), - vm_manifests: AHashMap::new(), - vm_placements: AHashMap::new(), - vm_data_cache: AHashMap::new(), - vm_keepalive_tasks: AHashMap::new(), - pending_resources_cache: None, - agent_vm_index: AHashMap::new(), - actor_kinds: AHashMap::new(), - cache_actor_finder: None, - }; - - scheduler_actor.start_actor_finder(actor_ref); - - Ok(scheduler_actor) - } - - async fn on_link_died( - &mut self, - actor_ref: WeakActorRef, - id: ActorId, - reason: ActorStopReason, - ) -> Result, Self::Error> { - warn!(?id, ?reason, "Linked actor died"); - - // check that scheduler actor is still alive. - let Some(_) = actor_ref.upgrade() else { - return Ok(ControlFlow::Break(ActorStopReason::Killed)); - }; - - match self.actor_kinds.remove(&id) { - Some(CachedActorKind::Agent) => self.cleanup_agent_actor(id), - Some(CachedActorKind::Vm) => self.cleanup_vm_actor(id), - None => {} - } - - // todo: attempt vm restarts if necessary. - - Ok(ControlFlow::Continue(())) - } -} - -impl Message for SchedulerActor { - type Reply = (); - - async fn handle(&mut self, msg: VmActorDiscovered, ctx: &mut Context) { - let actor_id = msg.actor_ref.id(); - let updater_is_running = self - .vm_keepalive_tasks - .get(&actor_id) - .is_some_and(|task| !task.is_finished()); - if updater_is_running { - return; - } - self.vm_keepalive_tasks.remove(&actor_id); - if let Err(error) = ctx.actor_ref().link_remote(&msg.actor_ref).await { - warn!(?error, ?actor_id, "failed to link VM actor"); - return; - } - self.actor_kinds.insert(actor_id, CachedActorKind::Vm); - let scheduler = ctx.actor_ref().clone(); - let actor_ref = msg.actor_ref; - let task = tokio::spawn(async move { - Self::vm_updater_task(scheduler, actor_ref).await; - }); - self.vm_keepalive_tasks.insert(actor_id, task); - } -} - -impl Message for SchedulerActor { - type Reply = (); - - async fn handle(&mut self, msg: AgentActorDiscovered, ctx: &mut Context) { - let actor_id = msg.actor_ref.id(); - let updater_is_running = self - .agent_keepalive_tasks - .get(&actor_id) - .is_some_and(|task| !task.is_finished()); - if updater_is_running { - return; - } - self.agent_keepalive_tasks.remove(&actor_id); - if let Err(error) = ctx.actor_ref().link_remote(&msg.actor_ref).await { - warn!(?error, ?actor_id, "failed to link agent actor"); - return; - } - self.actor_kinds.insert(actor_id, CachedActorKind::Agent); - let scheduler = ctx.actor_ref().clone(); - let actor_ref = msg.actor_ref; - let task = tokio::spawn(async move { - Self::agent_updater_task(scheduler, actor_ref).await; - }); - self.agent_keepalive_tasks.insert(actor_id, task); - } -} - -impl Message for SchedulerActor { - type Reply = (); - - async fn handle(&mut self, msg: VmUpdated, _ctx: &mut Context) { - let vmid = msg.data.vmid; - let actor_id = msg.actor_ref.id(); - self.vm_actorid_ulid_map.insert(actor_id, vmid); - if let Some(manifest) = msg.data.config { - self.vm_manifests.insert(vmid, manifest); - } - let cached_vm = CachedVMActor { - actor_ref: Some(msg.actor_ref), - }; - let entries = self.vm_data_cache.entry(vmid).or_default(); - Self::update_cached_vm_entry(entries, actor_id, cached_vm); - } -} - -#[allow(clippy::unused_async_trait_impl)] -impl Message for SchedulerActor { - type Reply = (); - - async fn handle(&mut self, msg: VmUpdaterStopped, _ctx: &mut Context) { - self.vm_keepalive_tasks.remove(&msg.actor_id); - self.actor_kinds.remove(&msg.actor_id); - let vmid = self.vm_actorid_ulid_map.remove(&msg.actor_id); - Self::remove_vm_actor(msg.actor_id, &mut self.vm_data_cache); - - if let Some(vmid) = vmid - && self - .vm_data_cache - .get(&vmid) - .is_none_or(|entries| entries.iter().all(|entry| entry.actor_ref.is_none())) - { - Self::remove_vm_state( - vmid, - &mut self.vm_manifests, - &mut self.vm_placements, - &mut self.vm_data_cache, - ); - } - } -} - -#[allow(clippy::unused_async_trait_impl)] -impl Message for SchedulerActor { - type Reply = (); - - async fn handle(&mut self, msg: AgentUpdated, _ctx: &mut Context) { - let Some(cached) = self.agent_data_cache.get_mut(&msg.actor_id) else { - if let AgentStatusUpdate::Full { revision, status } = msg.update { - self.agent_vm_index - .insert(msg.actor_id, status.vms.iter().copied().collect()); - self.invalidate_pending_resources(); - Self::reconcile_agent_placements( - msg.actor_id, - &status, - &self.vm_manifests, - &mut self.vm_placements, - ); - self.invalidate_pending_resources(); - self.agent_data_cache.insert( - msg.actor_id, - CachedAgentActor { - actor_ref: msg.actor_ref, - data: status, - status_revision: revision, - }, - ); - } - return; - }; - - let revision = match &msg.update { - AgentStatusUpdate::Full { revision, .. } - | AgentStatusUpdate::Delta { revision, .. } => *revision, - }; - if revision <= cached.status_revision { - return; - } - if let AgentStatusUpdate::Delta { added, removed, .. } = &msg.update { - let added = added.clone(); - let removed = removed.clone(); - cached.status_revision = apply_status_update(&mut cached.data, msg.update); - cached.actor_ref = msg.actor_ref; - self.agent_vm_index - .entry(msg.actor_id) - .or_default() - .extend(added.iter().copied()); - if let Some(index) = self.agent_vm_index.get_mut(&msg.actor_id) { - for vmid in &removed { - index.remove(vmid); - } - } - self.invalidate_pending_resources(); - Self::reconcile_agent_delta( - msg.actor_id, - &added, - &removed, - &self.vm_manifests, - &mut self.vm_placements, - ); - } else { - cached.status_revision = apply_status_update(&mut cached.data, msg.update); - cached.actor_ref = msg.actor_ref; - self.agent_vm_index - .insert(msg.actor_id, cached.data.vms.iter().copied().collect()); - Self::reconcile_agent_placements( - msg.actor_id, - &cached.data, - &self.vm_manifests, - &mut self.vm_placements, - ); - self.invalidate_pending_resources(); - } - } -} - -impl Message for SchedulerActor { - type Reply = (); - - async fn handle(&mut self, msg: AgentUpdaterStopped, _ctx: &mut Context) { - self.agent_keepalive_tasks.remove(&msg.actor_id); - self.actor_kinds.remove(&msg.actor_id); - self.agent_data_cache.remove(&msg.actor_id); - self.agent_vm_index.remove(&msg.actor_id); - Self::remove_agent_placements( - msg.actor_id, - &mut self.vm_manifests, - &mut self.vm_placements, - &mut self.vm_data_cache, - ); - self.invalidate_pending_resources(); - } -} - -impl Message for SchedulerActor { - type Reply = (); - - async fn handle(&mut self, _msg: ReconcileVmPlacements, _ctx: &mut Context) { - Self::cleanup_unresolved_vm_cache( - &mut self.vm_manifests, - &mut self.vm_placements, - &mut self.vm_data_cache, - ); - self.invalidate_pending_resources(); - } -} - -impl Message for SchedulerActor { - type Reply = Result; - - async fn handle( - &mut self, - msg: CreateVM, - _ctx: &mut Context, - ) -> Self::Reply { - let target_agent = self.schedule_agent(&msg)?; - - self.vm_manifests.insert(msg.vmid, msg.config.clone()); - self.invalidate_pending_resources(); - self.vm_placements - .entry(msg.vmid) - .or_default() - .push(VmPlacement { - agent_id: target_agent.id(), - lifecycle: VmLifecycle::Pending, - created_at: Instant::now(), - last_confirmed_at: None, - }); - self.vm_data_cache - .entry(msg.vmid) - .or_default() - .push(CachedVMActor { actor_ref: None }); - - let reply = target_agent.ask(&msg).await; - - if let Ok(reply) = &reply - && let Some(actor_id_bytes) = &reply.actor_id - && let Ok(actor_id) = ActorId::from_bytes(actor_id_bytes) - { - self.vm_actorid_ulid_map.insert(actor_id, msg.vmid); - } - - if reply.is_err() { - let actor_exists = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)) - .await - .ok() - .flatten() - .is_some(); - Self::rollback_failed_create( - msg.vmid, - actor_exists, - reply.as_ref().ok().and_then(|reply| { - reply - .actor_id - .as_deref() - .and_then(|bytes| ActorId::from_bytes(bytes).ok()) - }), - &mut self.vm_actorid_ulid_map, - &mut self.vm_manifests, - &mut self.vm_placements, - &mut self.vm_data_cache, - ); - } - - Ok(reply?) - } -} - -impl Message for SchedulerActor { - type Reply = Result; - - async fn handle( - &mut self, - msg: GetConsoleHistory, - _ctx: &mut Context, - ) -> Self::Reply { - let vm = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)).await?; - tracing::trace!(?vm, vmid = %msg.vmid, "GetConsoleHistory"); - if let Some(vm) = vm { - Ok(vm.ask(&msg).await?) - } else { - Err(eyre!("VM not found")) - } - } -} - -impl Message for SchedulerActor { - type Reply = Result; - - async fn handle( - &mut self, - msg: SendConsoleInput, - _ctx: &mut Context, - ) -> Self::Reply { - let vm = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)).await?; - tracing::trace!(?vm, vmid = %msg.vmid, bytes = msg.input.len(), "SendConsoleInput"); - if let Some(vm) = vm { - Ok(vm.ask(&msg).await?) - } else { - Err(eyre!("VM not found")) - } - } -} - -impl Message for SchedulerActor { - type Reply = Result; - - async fn handle( - &mut self, - msg: DeleteVM, - _ctx: &mut Context, - ) -> Self::Reply { - let vm = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)).await?; - tracing::trace!(?vm, "DeleteVM"); - if let Some(vm) = vm { - // don't update cache, because we rely on link dying and updater task to remove from cache once the VM is fully down. - vm.tell(&msg).send()?; - Ok(DeleteVMReply) - } else { - Err(eyre!("VM not found")) - } - } -} - -impl Message for SchedulerActor { - type Reply = Result; - - async fn handle( - &mut self, - msg: ShutdownVM, - _ctx: &mut Context, - ) -> Self::Reply { - let vm = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)).await?; - tracing::trace!(?vm, "ShutdownVM"); - if let Some(vm) = vm { - // don't update cache, because we rely on link dying and updater task to remove from cache once the VM is fully down. - vm.tell(&msg).send()?; - Ok(ShutdownVMReply) - } else { - Err(eyre!("VM not found")) - } - } -} - -/// this only gets data from the cache from agents -/// we may need a different message that actually forcibly runs/updates everything. -/// and/or messages that get data directly from the `VMActors`. -#[allow(clippy::unused_async_trait_impl)] -impl Message for SchedulerActor { - type Reply = Result; - - async fn handle( - &mut self, - _msg: AgentListVMs, - _ctx: &mut Context, - ) -> Self::Reply { - let total_vms = self - .agent_data_cache - .values() - .map(|agent| agent.data.vms.len()) - .sum(); - let mut vms = Vec::with_capacity(total_vms); - - for agent in self.agent_data_cache.values() { - vms.extend_from_slice(agent.data.vms.as_slice()); - } - - Ok(AgentListVMsReply { vms }) - } -} - -#[allow(clippy::unused_async_trait_impl)] -impl Message for SchedulerActor { - type Reply = Pong; - - async fn handle(&mut self, _msg: Ping, _ctx: &mut Context) -> Self::Reply { - Pong - } -} diff --git a/odorobo/src/actors/scheduler_actor/cache.rs b/odorobo/src/actors/scheduler_actor/cache.rs new file mode 100644 index 0000000..00c26dd --- /dev/null +++ b/odorobo/src/actors/scheduler_actor/cache.rs @@ -0,0 +1,310 @@ +//! Cache maintenance, placement reconciliation, and actor cleanup. + +use std::time::{Duration, Instant}; + +use crate::messages::agent::AgentStatus; +use ahash::{AHashMap, AHashSet}; +use kameo::prelude::*; +use tracing::trace; +use ulid::Ulid; + +use crate::actors::agent_actor::AgentActor; +use crate::manifest::VmManifest; + +use super::{CachedVMActor, SchedulerActor, VmLifecycle, VmPlacement}; + +const UNRESOLVED_VM_CACHE_TIMEOUT: Duration = Duration::from_secs(30); + +impl SchedulerActor { + pub(super) fn shrink_non_migrating_entries(entries: &mut Vec) { + if entries.len() <= 1 && entries.capacity() >= 4 { + entries.shrink_to_fit(); + } + } + + pub(super) fn update_cached_vm_entry( + entries: &mut Vec, + actor_id: ActorId, + cached_vm: CachedVMActor, + ) { + if let Some(entry) = entries.iter_mut().find(|entry| { + entry + .actor_ref + .as_ref() + .is_some_and(|actor| actor.id() == actor_id) + }) { + *entry = cached_vm; + } else if let Some(entry) = entries.iter_mut().find(|entry| entry.actor_ref.is_none()) { + *entry = cached_vm; + } else { + entries.push(cached_vm); + } + Self::shrink_non_migrating_entries(entries); + } + + pub(super) fn cleanup_unresolved_vm_cache( + manifests: &mut AHashMap, + placements: &mut AHashMap>, + data_cache: &mut AHashMap>, + ) { + let now = Instant::now(); + let empty_vmids: Vec<_> = placements + .iter_mut() + .filter_map(|(vmid, entries)| { + entries.retain(|entry| { + entry.lifecycle != VmLifecycle::Pending + || now.duration_since(entry.created_at) < UNRESOLVED_VM_CACHE_TIMEOUT + }); + Self::shrink_non_migrating_entries(entries); + entries.is_empty().then_some(*vmid) + }) + .collect(); + + for vmid in empty_vmids { + Self::remove_vm_state(vmid, manifests, placements, data_cache); + } + } + + /// Returns every VM that could be on an agent, deduplicating each source in + /// constant expected time. + pub(super) fn placement_vm_ids( + placements: &AHashMap>, + indexed: Option<&AHashSet>, + agent_id: ActorId, + observed: &[Ulid], + ) -> Vec { + let indexed_len = indexed.map_or(0, |index| index.len()); + let mut vmids = Vec::with_capacity(observed.len().max(indexed_len)); + let mut seen = AHashSet::with_capacity(observed.len().saturating_add(indexed_len)); + + for vmid in observed { + if seen.insert(*vmid) { + vmids.push(*vmid); + } + } + if let Some(indexed) = indexed { + for vmid in indexed { + if seen.insert(*vmid) { + vmids.push(*vmid); + } + } + } + for vmid in placements.iter().filter_map(|(vmid, entries)| { + entries + .iter() + .any(|entry| entry.agent_id == agent_id) + .then_some(vmid) + }) { + if seen.insert(*vmid) { + vmids.push(*vmid); + } + } + vmids + } + + pub(super) fn remove_vm_state( + vmid: Ulid, + manifests: &mut AHashMap, + placements: &mut AHashMap>, + data_cache: &mut AHashMap>, + ) { + manifests.remove(&vmid); + placements.remove(&vmid); + data_cache.remove(&vmid); + } + + pub(super) fn remove_vm_actor( + actor_id: ActorId, + data_cache: &mut AHashMap>, + ) { + let empty_vmids: Vec<_> = data_cache + .iter_mut() + .filter_map(|(vmid, entries)| { + entries.retain(|entry| { + entry + .actor_ref + .as_ref() + .is_none_or(|actor| actor.id() != actor_id) + }); + Self::shrink_non_migrating_entries(entries); + entries.is_empty().then_some(*vmid) + }) + .collect(); + for vmid in empty_vmids { + data_cache.remove(&vmid); + } + } + + pub(super) fn remove_agent_placements( + agent_id: ActorId, + manifests: &mut AHashMap, + placements: &mut AHashMap>, + data_cache: &mut AHashMap>, + ) { + let empty_vmids: Vec<_> = placements + .iter_mut() + .filter_map(|(vmid, entries)| { + entries.retain(|entry| entry.agent_id != agent_id); + Self::shrink_non_migrating_entries(entries); + entries.is_empty().then_some(*vmid) + }) + .collect(); + + for vmid in empty_vmids { + Self::remove_vm_state(vmid, manifests, placements, data_cache); + } + } + + pub(super) fn rollback_failed_create( + vmid: Ulid, + actor_exists: bool, + actor_id: Option, + actor_map: &mut AHashMap, + manifests: &mut AHashMap, + placements: &mut AHashMap>, + data_cache: &mut AHashMap>, + ) { + if !actor_exists { + if let Some(actor_id) = actor_id + && actor_map.get(&actor_id) == Some(&vmid) + { + actor_map.remove(&actor_id); + } + Self::remove_vm_state(vmid, manifests, placements, data_cache); + } + } + + pub(super) fn cleanup_agent_actor(&mut self, actor_id: ActorId) { + if let Some(keepalive_task) = self.agent_keepalive_tasks.remove(&actor_id) { + trace!(?actor_id, "Aborting agent keepalive task"); + keepalive_task.abort(); + } + self.agent_data_cache.remove(&actor_id); + self.agent_vm_index.remove(&actor_id); + self.invalidate_pending_resources(); + Self::remove_agent_placements( + actor_id, + &mut self.vm_manifests, + &mut self.vm_placements, + &mut self.vm_data_cache, + ); + } + + pub(super) fn cleanup_vm_actor(&mut self, actor_id: ActorId) { + if let Some(keepalive_task) = self.vm_keepalive_tasks.remove(&actor_id) { + trace!(?actor_id, "Aborting VM keepalive task"); + keepalive_task.abort(); + } + let vmid = self.vm_actorid_ulid_map.remove(&actor_id); + self.invalidate_pending_resources(); + Self::remove_vm_actor(actor_id, &mut self.vm_data_cache); + if let Some(vmid) = vmid + && self + .vm_data_cache + .get(&vmid) + .is_none_or(|entries| entries.iter().all(|entry| entry.actor_ref.is_none())) + { + Self::remove_vm_state( + vmid, + &mut self.vm_manifests, + &mut self.vm_placements, + &mut self.vm_data_cache, + ); + } + } + + pub(super) fn reconcile_agent_delta( + agent_id: ActorId, + added: &[Ulid], + _removed: &[Ulid], + manifests: &AHashMap, + placements: &mut AHashMap>, + ) { + let now = Instant::now(); + for vmid in added { + if !manifests.contains_key(vmid) { + continue; + } + let entries = placements.entry(*vmid).or_default(); + if let Some(entry) = entries.iter_mut().find(|entry| entry.agent_id == agent_id) { + entry.lifecycle = VmLifecycle::Running; + entry.last_confirmed_at = Some(now); + } else { + entries.push(VmPlacement { + agent_id, + lifecycle: VmLifecycle::Running, + created_at: now, + last_confirmed_at: Some(now), + }); + } + } + + // A removal is an observation about the agent, not a change to the + // scheduler's desired state. Keep the placement so reconciliation can + // schedule the VM again. The full status path performs the same + // distinction for snapshots. + } + + pub(super) fn reconcile_agent_placements( + agent_id: ActorId, + status: &AgentStatus, + manifests: &AHashMap, + placements: &mut AHashMap>, + ) { + let now = Instant::now(); + let observed: AHashSet<_> = status.vms.iter().copied().collect(); + let missing_agent_placements: Vec<_> = observed + .iter() + .filter(|vmid| manifests.contains_key(vmid)) + .filter(|vmid| { + placements + .get(vmid) + .is_none_or(|entries| !entries.iter().any(|entry| entry.agent_id == agent_id)) + }) + .copied() + .collect(); + for vmid in missing_agent_placements { + placements.entry(vmid).or_default().push(VmPlacement { + agent_id, + lifecycle: VmLifecycle::Running, + created_at: now, + last_confirmed_at: Some(now), + }); + } + + let empty_vmids: Vec<_> = placements + .iter_mut() + .filter_map(|(vmid, entries)| { + for entry in entries + .iter_mut() + .filter(|entry| entry.agent_id == agent_id) + { + if observed.contains(vmid) { + entry.lifecycle = VmLifecycle::Running; + entry.last_confirmed_at = Some(now); + } + } + Self::shrink_non_migrating_entries(entries); + entries.is_empty().then_some(*vmid) + }) + .collect(); + for vmid in empty_vmids { + placements.remove(&vmid); + } + } + + #[expect(dead_code, reason = "reserved for explicit placement by actor id")] + fn lookup_agent_by_actor_id(&self, actor_id: &ActorId) -> Option> { + self.agent_data_cache + .get(actor_id) + .map(|data| data.actor_ref.clone()) + } + + #[expect(dead_code, reason = "reserved for explicit placement by hostname")] + fn lookup_agent_by_hostname(&self, hostname: &str) -> Option> { + self.agent_data_cache + .values() + .find(|data| data.data.hostname == hostname) + .map(|data| data.actor_ref.clone()) + } +} diff --git a/odorobo/src/actors/scheduler_actor/discovery.rs b/odorobo/src/actors/scheduler_actor/discovery.rs new file mode 100644 index 0000000..c1f562a --- /dev/null +++ b/odorobo/src/actors/scheduler_actor/discovery.rs @@ -0,0 +1,387 @@ +//! Remote actor discovery, polling tasks, and their internal messages. + +use std::time::Duration; + +use kameo::prelude::*; +use libp2p::futures::TryStreamExt; +use stable_eyre::{Report, eyre::eyre}; +use tracing::{trace, warn}; + +use crate::actors::agent_actor::AgentActor; +use crate::ch_driver::actor::VMActor; +use crate::messages::agent::{AgentStatusUpdate, GetAgentStatus, apply_status_update}; +use crate::messages::vm::{GetVMHeartbeat, GetVMInfo}; +use crate::utils::actor_names::{AGENT, VM}; + +use super::{ + AgentActorDiscovered, AgentUpdated, AgentUpdaterStopped, CachedActorKind, CachedAgentActor, + CachedVMActor, ReconcileVmPlacements, SchedulerActor, VmActorDiscovered, VmUpdated, + VmUpdaterStopped, +}; + +impl SchedulerActor { + async fn vm_actor_finder(parent_actor_ref: ActorRef) -> Result<(), Report> { + trace!("running vm_actor_finder"); + + let mut vm_actor_stream = RemoteActorRef::::lookup_all(VM); + + while let Some(vm_actor) = vm_actor_stream.try_next().await? { + parent_actor_ref + .tell(VmActorDiscovered { + actor_ref: vm_actor, + }) + .send() + .await?; + } + + Ok(()) + } + + async fn vm_updater_task(scheduler: ActorRef, actor_ref: RemoteActorRef) { + let mut interval = tokio::time::interval(Duration::from_secs(1)); + let mut fails: u8 = 0; + let mut initialized = false; + + loop { + if !initialized { + if let Ok(data) = actor_ref.ask(&GetVMInfo { vmid: None }).await { + let send_result = scheduler + .tell(VmUpdated { + actor_ref: actor_ref.clone(), + data, + }) + .send() + .await + .map_err(|error| eyre!("failed to send VM update: {error}")); + if let Err(error) = send_result { + warn!(?error, "VM updater could not notify scheduler"); + return; + } + initialized = true; + fails = 0; + } else { + fails = fails.saturating_add(1); + } + } else if actor_ref.ask(&GetVMHeartbeat).await.is_ok() { + fails = 0; + } else { + fails = fails.saturating_add(1); + } + + if fails > 5 { + warn!( + ?actor_ref, + "can no longer reach vm actor, cleaning up cache entries" + ); + + let send_result = scheduler + .tell(VmUpdaterStopped { + actor_id: actor_ref.id(), + }) + .send() + .await + .map_err(|error| eyre!("failed to send VM stop: {error}")); + if let Err(error) = send_result { + warn!(?error, "VM updater could not notify scheduler"); + } + return; + } + + interval.tick().await; + } + } + + async fn agent_actor_finder(parent_actor_ref: ActorRef) -> Result<(), Report> { + trace!("running agent_actor_finder"); + + let mut agent_actor_stream = RemoteActorRef::::lookup_all(AGENT); + + while let Some(agent_actor) = agent_actor_stream.try_next().await? { + parent_actor_ref + .tell(AgentActorDiscovered { + actor_ref: agent_actor, + }) + .send() + .await?; + } + + Ok(()) + } + + async fn agent_updater_task(scheduler: ActorRef, actor_ref: RemoteActorRef) { + let mut interval = tokio::time::interval(Duration::from_secs(1)); + let mut status_revision = 0; + let mut initial_status = true; + let mut fails: u8 = 0; + loop { + if let Ok(update) = actor_ref + .ask(&GetAgentStatus { + since_revision: status_revision, + initial: initial_status, + }) + .await + { + status_revision = match &update { + AgentStatusUpdate::Full { revision, .. } + | AgentStatusUpdate::Delta { revision, .. } => *revision, + }; + initial_status = false; + let send_result = scheduler + .tell(AgentUpdated { + actor_id: actor_ref.id(), + actor_ref: actor_ref.clone(), + update, + }) + .send() + .await + .map_err(|error| eyre!("failed to send agent update: {error}")); + if let Err(error) = send_result { + warn!(?error, "agent updater could not notify scheduler"); + return; + } + fails = 0; + } else { + fails = fails.saturating_add(1); + } + + if fails > 5 { + warn!( + ?actor_ref, + "can no longer reach agent actor, stopping updater" + ); + let send_result = scheduler + .tell(AgentUpdaterStopped { + actor_id: actor_ref.id(), + }) + .send() + .await + .map_err(|error| eyre!("failed to send agent stop: {error}")); + if let Err(error) = send_result { + warn!(?error, "agent updater could not notify scheduler"); + } + return; + } + + interval.tick().await; + } + } + + pub(super) fn start_actor_finder(&mut self, actor_ref: ActorRef) { + self.cache_actor_finder = Some(tokio::spawn(async move { + let mut interval = tokio::time::interval(Duration::from_secs(5)); + loop { + let vm_result = Self::vm_actor_finder(actor_ref.clone()).await; + let agent_result = Self::agent_actor_finder(actor_ref.clone()).await; + if let Err(error) = vm_result { + warn!(?error, "VM actor discovery failed"); + } + if let Err(error) = agent_result { + warn!(?error, "agent actor discovery failed"); + } + actor_ref.tell(ReconcileVmPlacements).send().await.ok(); + interval.tick().await; + } + })); + } +} + +impl Message for SchedulerActor { + type Reply = (); + + async fn handle(&mut self, msg: VmActorDiscovered, ctx: &mut Context) { + let actor_id = msg.actor_ref.id(); + let updater_is_running = self + .vm_keepalive_tasks + .get(&actor_id) + .is_some_and(|task| !task.is_finished()); + if updater_is_running { + return; + } + self.vm_keepalive_tasks.remove(&actor_id); + if let Err(error) = ctx.actor_ref().link_remote(&msg.actor_ref).await { + warn!(?error, ?actor_id, "failed to link VM actor"); + return; + } + self.actor_kinds.insert(actor_id, CachedActorKind::Vm); + let scheduler = ctx.actor_ref().clone(); + let actor_ref = msg.actor_ref; + let task = tokio::spawn(async move { + Self::vm_updater_task(scheduler, actor_ref).await; + }); + self.vm_keepalive_tasks.insert(actor_id, task); + } +} + +impl Message for SchedulerActor { + type Reply = (); + + async fn handle(&mut self, msg: AgentActorDiscovered, ctx: &mut Context) { + let actor_id = msg.actor_ref.id(); + let updater_is_running = self + .agent_keepalive_tasks + .get(&actor_id) + .is_some_and(|task| !task.is_finished()); + if updater_is_running { + return; + } + self.agent_keepalive_tasks.remove(&actor_id); + if let Err(error) = ctx.actor_ref().link_remote(&msg.actor_ref).await { + warn!(?error, ?actor_id, "failed to link agent actor"); + return; + } + self.actor_kinds.insert(actor_id, CachedActorKind::Agent); + let scheduler = ctx.actor_ref().clone(); + let actor_ref = msg.actor_ref; + let task = tokio::spawn(async move { + Self::agent_updater_task(scheduler, actor_ref).await; + }); + self.agent_keepalive_tasks.insert(actor_id, task); + } +} + +impl Message for SchedulerActor { + type Reply = (); + + async fn handle(&mut self, msg: VmUpdated, _ctx: &mut Context) { + let vmid = msg.data.vmid; + let actor_id = msg.actor_ref.id(); + self.vm_actorid_ulid_map.insert(actor_id, vmid); + if let Some(manifest) = msg.data.config { + self.vm_manifests.insert(vmid, manifest); + } + let cached_vm = CachedVMActor { + actor_ref: Some(msg.actor_ref), + }; + let entries = self.vm_data_cache.entry(vmid).or_default(); + Self::update_cached_vm_entry(entries, actor_id, cached_vm); + } +} + +impl Message for SchedulerActor { + type Reply = (); + + async fn handle(&mut self, msg: VmUpdaterStopped, _ctx: &mut Context) { + self.vm_keepalive_tasks.remove(&msg.actor_id); + self.actor_kinds.remove(&msg.actor_id); + let vmid = self.vm_actorid_ulid_map.remove(&msg.actor_id); + Self::remove_vm_actor(msg.actor_id, &mut self.vm_data_cache); + + if let Some(vmid) = vmid + && self + .vm_data_cache + .get(&vmid) + .is_none_or(|entries| entries.iter().all(|entry| entry.actor_ref.is_none())) + { + Self::remove_vm_state( + vmid, + &mut self.vm_manifests, + &mut self.vm_placements, + &mut self.vm_data_cache, + ); + } + } +} + +impl Message for SchedulerActor { + type Reply = (); + + async fn handle(&mut self, msg: AgentUpdated, _ctx: &mut Context) { + let Some(cached) = self.agent_data_cache.get_mut(&msg.actor_id) else { + if let AgentStatusUpdate::Full { revision, status } = msg.update { + self.agent_vm_index + .insert(msg.actor_id, status.vms.iter().copied().collect()); + self.invalidate_pending_resources(); + Self::reconcile_agent_placements( + msg.actor_id, + &status, + &self.vm_manifests, + &mut self.vm_placements, + ); + self.invalidate_pending_resources(); + self.agent_data_cache.insert( + msg.actor_id, + CachedAgentActor { + actor_ref: msg.actor_ref, + data: status, + status_revision: revision, + }, + ); + } + return; + }; + + let revision = match &msg.update { + AgentStatusUpdate::Full { revision, .. } + | AgentStatusUpdate::Delta { revision, .. } => *revision, + }; + if revision <= cached.status_revision { + return; + } + if let AgentStatusUpdate::Delta { added, removed, .. } = &msg.update { + let added = added.clone(); + let removed = removed.clone(); + cached.status_revision = apply_status_update(&mut cached.data, msg.update); + cached.actor_ref = msg.actor_ref; + self.agent_vm_index + .entry(msg.actor_id) + .or_default() + .extend(added.iter().copied()); + if let Some(index) = self.agent_vm_index.get_mut(&msg.actor_id) { + for vmid in &removed { + index.remove(vmid); + } + } + self.invalidate_pending_resources(); + Self::reconcile_agent_delta( + msg.actor_id, + &added, + &removed, + &self.vm_manifests, + &mut self.vm_placements, + ); + } else { + cached.status_revision = apply_status_update(&mut cached.data, msg.update); + cached.actor_ref = msg.actor_ref; + self.agent_vm_index + .insert(msg.actor_id, cached.data.vms.iter().copied().collect()); + Self::reconcile_agent_placements( + msg.actor_id, + &cached.data, + &self.vm_manifests, + &mut self.vm_placements, + ); + self.invalidate_pending_resources(); + } + } +} + +impl Message for SchedulerActor { + type Reply = (); + + async fn handle(&mut self, msg: AgentUpdaterStopped, _ctx: &mut Context) { + self.agent_keepalive_tasks.remove(&msg.actor_id); + self.actor_kinds.remove(&msg.actor_id); + self.agent_data_cache.remove(&msg.actor_id); + self.agent_vm_index.remove(&msg.actor_id); + Self::remove_agent_placements( + msg.actor_id, + &mut self.vm_manifests, + &mut self.vm_placements, + &mut self.vm_data_cache, + ); + self.invalidate_pending_resources(); + } +} + +impl Message for SchedulerActor { + type Reply = (); + + async fn handle(&mut self, _msg: ReconcileVmPlacements, _ctx: &mut Context) { + Self::cleanup_unresolved_vm_cache( + &mut self.vm_manifests, + &mut self.vm_placements, + &mut self.vm_data_cache, + ); + self.invalidate_pending_resources(); + } +} diff --git a/odorobo/src/actors/scheduler_actor/handlers.rs b/odorobo/src/actors/scheduler_actor/handlers.rs new file mode 100644 index 0000000..8a9876b --- /dev/null +++ b/odorobo/src/actors/scheduler_actor/handlers.rs @@ -0,0 +1,248 @@ +//! Actor lifecycle implementation and public scheduler message handlers. + +use std::ops::ControlFlow; +use std::time::Instant; + +use ahash::AHashMap; +use kameo::prelude::*; +use stable_eyre::{Report, eyre::eyre}; +use tracing::{info, warn}; + +use crate::ch_driver::actor::VMActor; +use crate::messages::vm::{ + AgentListVMs, AgentListVMsReply, CreateVM, CreateVMReply, DeleteVM, DeleteVMReply, + GetConsoleHistory, GetConsoleHistoryReply, SendConsoleInput, SendConsoleInputReply, ShutdownVM, + ShutdownVMReply, +}; +use crate::messages::{Ping, Pong}; +use crate::utils::actor_names::vm_actor_id; + +use super::{CachedActorKind, CachedVMActor, SchedulerActor, VmLifecycle, VmPlacement}; + +impl Actor for SchedulerActor { + type Args = (); + type Error = Report; + + async fn on_start(_state: Self::Args, actor_ref: ActorRef) -> Result { + let peer_id = *actor_ref.id().peer_id().unwrap(); + + info!(?peer_id, "Scheduler Actor started!"); + + let mut scheduler_actor = Self { + agent_data_cache: AHashMap::new(), + agent_keepalive_tasks: AHashMap::new(), + vm_actorid_ulid_map: AHashMap::new(), + vm_manifests: AHashMap::new(), + vm_placements: AHashMap::new(), + vm_data_cache: AHashMap::new(), + vm_keepalive_tasks: AHashMap::new(), + pending_resources_cache: None, + agent_vm_index: AHashMap::new(), + actor_kinds: AHashMap::new(), + cache_actor_finder: None, + }; + + scheduler_actor.start_actor_finder(actor_ref); + + Ok(scheduler_actor) + } + + async fn on_link_died( + &mut self, + actor_ref: WeakActorRef, + id: ActorId, + reason: ActorStopReason, + ) -> Result, Self::Error> { + warn!(?id, ?reason, "Linked actor died"); + + // check that scheduler actor is still alive. + let Some(_) = actor_ref.upgrade() else { + return Ok(ControlFlow::Break(ActorStopReason::Killed)); + }; + + match self.actor_kinds.remove(&id) { + Some(CachedActorKind::Agent) => self.cleanup_agent_actor(id), + Some(CachedActorKind::Vm) => self.cleanup_vm_actor(id), + None => {} + } + + // todo: attempt vm restarts if necessary. + + Ok(ControlFlow::Continue(())) + } +} +impl Message for SchedulerActor { + type Reply = Result; + + async fn handle( + &mut self, + msg: CreateVM, + _ctx: &mut Context, + ) -> Self::Reply { + let target_agent = self.schedule_agent(&msg)?; + + self.vm_manifests.insert(msg.vmid, msg.config.clone()); + self.invalidate_pending_resources(); + self.vm_placements + .entry(msg.vmid) + .or_default() + .push(VmPlacement { + agent_id: target_agent.id(), + lifecycle: VmLifecycle::Pending, + created_at: Instant::now(), + last_confirmed_at: None, + }); + self.vm_data_cache + .entry(msg.vmid) + .or_default() + .push(CachedVMActor { actor_ref: None }); + + let reply = target_agent.ask(&msg).await; + + if let Ok(reply) = &reply + && let Some(actor_id_bytes) = &reply.actor_id + && let Ok(actor_id) = ActorId::from_bytes(actor_id_bytes) + { + self.vm_actorid_ulid_map.insert(actor_id, msg.vmid); + } + + if reply.is_err() { + let actor_exists = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)) + .await + .ok() + .flatten() + .is_some(); + Self::rollback_failed_create( + msg.vmid, + actor_exists, + reply.as_ref().ok().and_then(|reply| { + reply + .actor_id + .as_deref() + .and_then(|bytes| ActorId::from_bytes(bytes).ok()) + }), + &mut self.vm_actorid_ulid_map, + &mut self.vm_manifests, + &mut self.vm_placements, + &mut self.vm_data_cache, + ); + } + + Ok(reply?) + } +} + +impl Message for SchedulerActor { + type Reply = Result; + + async fn handle( + &mut self, + msg: GetConsoleHistory, + _ctx: &mut Context, + ) -> Self::Reply { + let vm = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)).await?; + tracing::trace!(?vm, vmid = %msg.vmid, "GetConsoleHistory"); + if let Some(vm) = vm { + Ok(vm.ask(&msg).await?) + } else { + Err(eyre!("VM not found")) + } + } +} + +impl Message for SchedulerActor { + type Reply = Result; + + async fn handle( + &mut self, + msg: SendConsoleInput, + _ctx: &mut Context, + ) -> Self::Reply { + let vm = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)).await?; + tracing::trace!( + ?vm, + vmid = %msg.vmid, + bytes = msg.input.len(), + "SendConsoleInput" + ); + if let Some(vm) = vm { + Ok(vm.ask(&msg).await?) + } else { + Err(eyre!("VM not found")) + } + } +} + +impl Message for SchedulerActor { + type Reply = Result; + + async fn handle( + &mut self, + msg: DeleteVM, + _ctx: &mut Context, + ) -> Self::Reply { + let vm = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)).await?; + tracing::trace!(?vm, "DeleteVM"); + if let Some(vm) = vm { + // don't update cache, because we rely on link dying and updater task to remove from cache once the VM is fully down. + vm.tell(&msg).send()?; + Ok(DeleteVMReply) + } else { + Err(eyre!("VM not found")) + } + } +} + +impl Message for SchedulerActor { + type Reply = Result; + + async fn handle( + &mut self, + msg: ShutdownVM, + _ctx: &mut Context, + ) -> Self::Reply { + let vm = RemoteActorRef::::lookup(vm_actor_id(msg.vmid)).await?; + tracing::trace!(?vm, "ShutdownVM"); + if let Some(vm) = vm { + // don't update cache, because we rely on link dying and updater task to remove from cache once the VM is fully down. + vm.tell(&msg).send()?; + Ok(ShutdownVMReply) + } else { + Err(eyre!("VM not found")) + } + } +} + +/// this only gets data from the cache from agents +/// we may need a different message that actually forcibly runs/updates everything. +/// and/or messages that get data directly from the `VMActors`. +impl Message for SchedulerActor { + type Reply = Result; + + async fn handle( + &mut self, + _msg: AgentListVMs, + _ctx: &mut Context, + ) -> Self::Reply { + let total_vms = self + .agent_data_cache + .values() + .map(|agent| agent.data.vms.len()) + .sum(); + let mut vms = Vec::with_capacity(total_vms); + + for agent in self.agent_data_cache.values() { + vms.extend_from_slice(agent.data.vms.as_slice()); + } + + Ok(AgentListVMsReply { vms }) + } +} + +impl Message for SchedulerActor { + type Reply = Pong; + + async fn handle(&mut self, _msg: Ping, _ctx: &mut Context) -> Self::Reply { + Pong + } +} diff --git a/odorobo/src/actors/scheduler_actor/scheduling.rs b/odorobo/src/actors/scheduler_actor/scheduling.rs new file mode 100644 index 0000000..30051e6 --- /dev/null +++ b/odorobo/src/actors/scheduler_actor/scheduling.rs @@ -0,0 +1,371 @@ +//! Cache-only scheduling policy, capacity accounting, and affinity evaluation. + +use std::cmp::Ordering; + +use ahash::AHashMap; +use kameo::prelude::{ActorId, RemoteActorRef}; +use stable_eyre::{Report, eyre::OptionExt}; +use tracing::info; +use ulid::Ulid; + +use crate::actors::agent_actor::AgentActor; +use crate::manifest::{ + AffinityRequirement, AffinityStrictness, AffinityType, MetadataTable, Operator, VmManifest, +}; +use crate::messages::vm::CreateVM; + +use super::{CachedAgentActor, SchedulerActor, VmLifecycle, VmPlacement}; + +const VCPU_OVERPROVISIONMENT_NUMERATOR: u32 = 2; +const VCPU_OVERPROVISIONMENT_DENOMINATOR: u32 = 1; +type MetadataTables<'a> = [( + &'a std::collections::BTreeMap, + &'a std::collections::BTreeMap, +)]; + +impl SchedulerActor { + pub(super) fn pending_resources(&mut self) -> &AHashMap { + if self.pending_resources_cache.is_none() { + self.pending_resources_cache = Some(pending_resources_by_agent( + &self.vm_manifests, + &self.vm_placements, + )); + } + self.pending_resources_cache + .as_ref() + .expect("pending resources cache was just initialized") + } + + pub(super) fn invalidate_pending_resources(&mut self) { + self.pending_resources_cache = None; + } + + /// Determine the best agent to schedule a specific VM creation request to. + /// + /// The scheduler first filters agents by capacity and affinity requirements, + /// then uses affinity and general resource scores to select the best match. + /// Affinity rules are roughly based on + /// . + pub(super) fn schedule_agent( + &mut self, + msg: &CreateVM, + ) -> Result, Report> { + self.pending_resources(); + let pending_resources = self + .pending_resources_cache + .as_ref() + .expect("pending resources cache was just initialized"); + let mut best_agent = None; + let mut best_score = AgentScore::REJECTED; + + for agent in self.agent_data_cache.values() { + let score = self.score_agent(msg, agent, pending_resources); + + if score > best_score { + best_agent = Some(agent.actor_ref.clone()); + best_score = score; + } + } + + info!(?best_score, ?best_agent, "best agent"); + + best_agent.ok_or_eyre("No valid agents found.") + } + + #[expect(dead_code, reason = "reserved for a future batch create message")] + fn schedule_agents( + &mut self, + msgs: &[CreateVM], + ) -> Vec, Report>> { + self.pending_resources(); + let pending_resources = self + .pending_resources_cache + .as_ref() + .expect("pending resources cache was just initialized"); + msgs.iter() + .map(|msg| { + let mut best_agent = None; + let mut best_score = AgentScore::REJECTED; + for agent in self.agent_data_cache.values() { + let score = self.score_agent(msg, agent, pending_resources); + if score > best_score { + best_agent = Some(agent.actor_ref.clone()); + best_score = score; + } + } + best_agent.ok_or_eyre("No valid agents found.") + }) + .collect() + } + + // this function intentionally only checks against the cache. this has some positives and negatives: + // positive: it will never trigger any network requests so its very fast, and having to do network requests for scoring whenever we want to schedule a vm is likely a bad idea + // negative: it technically has a delayed view of the cluster, meaning that some things that happened in the future, may not exist yet. so we need to be careful about how this is done so affinity rules are not accidentally broken. mostly this means, if we do anything that could affect the outcome of an affinity rule (ex: network request to an agent), we need to update the cache, before we do the action. + fn score_agent( + &self, + msg: &CreateVM, + agent: &CachedAgentActor, + pending_resources: &AHashMap, + ) -> AgentScore { + let mut score = AgentScore::default(); + + let agent_max_vcpus = agent + .data + .vcpus + .saturating_mul(VCPU_OVERPROVISIONMENT_NUMERATOR) + .checked_div(VCPU_OVERPROVISIONMENT_DENOMINATOR) + .unwrap_or(u32::MAX); + // todo: do we care about VMData.max_vcpus? + let (pending_vcpus, pending_ram) = pending_resources + .get(&agent.actor_ref.id()) + .copied() + .unwrap_or_default(); + let used_vcpus = agent.data.used_vcpus.saturating_add(pending_vcpus); + let used_ram = agent.data.used_ram.as_u64().saturating_add(pending_ram); + let requested_vcpus = msg.config.desired.compute.vcpus; + let requested_memory = msg.config.desired.compute.memory_bytes; + let agent_used_vcpus = used_vcpus.saturating_add(requested_vcpus); + + if !has_capacity( + agent_max_vcpus, + used_vcpus, + requested_vcpus, + agent.data.ram.as_u64(), + used_ram, + requested_memory, + ) { + return AgentScore::REJECTED; + } + + #[expect( + clippy::cast_precision_loss, + reason = "the scheduler score intentionally uses f32 ratios" + )] + #[expect( + clippy::arithmetic_side_effects, + reason = "the preceding capacity check guarantees non-negative subtraction" + )] + let vcpu_headroom = (agent_max_vcpus - agent_used_vcpus) as f32 / agent_max_vcpus as f32; + score.general += vcpu_headroom; + + // todo: add ram overprovisionment. not adding this to scheduler until it works on the hypervisor side. + let agent_max_ram = agent.data.ram; + let agent_used_ram = bytesize::ByteSize::b(used_ram.saturating_add(requested_memory)); + + #[expect( + clippy::cast_precision_loss, + reason = "the scheduler score intentionally uses f32 ratios" + )] + let ram_headroom = agent_max_ram + .as_u64() + .saturating_sub(agent_used_ram.as_u64()) as f32 + / agent_max_ram.as_u64() as f32; + score.general += ram_headroom; + + // Roughly based on . + if !msg.config.desired.placement.affinity.is_empty() { + let affinity_rules = &msg.config.desired.placement.affinity; + for rule in affinity_rules { + let mut metadata_tables = Vec::with_capacity(1); + match rule.affinity_type { + AffinityType::VirtualMachine => { + metadata_tables.extend( + Self::placement_vm_ids( + &self.vm_placements, + self.agent_vm_index.get(&agent.actor_ref.id()), + agent.actor_ref.id(), + &agent.data.vms, + ) + .into_iter() + .filter_map(|vmid| self.vm_manifests.get(&vmid)) + .map(|manifest| { + ( + &manifest.desired.metadata.labels, + &manifest.desired.metadata.annotations, + ) + }), + ); + } + AffinityType::Agent => { + metadata_tables.push(( + &agent.data.metadata.labels, + &agent.data.metadata.annotations, + )); + } + } + + let follows_rule = evaluate_affinity_rule(&metadata_tables, rule); + + let Some(affinity_delta) = affinity_delta(&rule.strictness, follows_rule) else { + return AgentScore::REJECTED; + }; + score.affinity = score.affinity.saturating_add(affinity_delta); + } + } + + // todo (future): possibly keep a percent of agents completely empty, to be able to be converted to dedis automatically. + // they would have their agent score set to like f32::MIN, so they can be scheduled to if there is no other available agents. + // rough pseudo code to implement this: + // if agent.metadata.vms.len() == 0 && hash(agent.config.hostname) % total_chance < threshold { + // agent_score = 1; + // } + + score + } +} + +pub(super) const fn has_capacity( + max_vcpus: u32, + used_vcpus: u32, + requested_vcpus: u32, + max_ram: u64, + used_ram: u64, + requested_ram: u64, +) -> bool { + used_vcpus.saturating_add(requested_vcpus) <= max_vcpus + && used_ram.saturating_add(requested_ram) <= max_ram +} + +fn pending_resources_by_agent( + manifests: &AHashMap, + placements: &AHashMap>, +) -> AHashMap { + let mut resources = AHashMap::new(); + for (vmid, entries) in placements { + let Some(manifest) = manifests.get(vmid) else { + continue; + }; + for entry in entries { + if entry.lifecycle == VmLifecycle::Pending { + let totals = resources.entry(entry.agent_id).or_insert((0u32, 0u64)); + totals.0 = totals.0.saturating_add(manifest.desired.compute.vcpus); + totals.1 = totals + .1 + .saturating_add(manifest.desired.compute.memory_bytes); + } + } + } + resources +} + +#[cfg(test)] +pub(super) fn pending_resources_for_agent( + manifests: &AHashMap, + placements: &AHashMap>, + agent_id: ActorId, +) -> (u32, u64) { + pending_resources_by_agent(manifests, placements) + .get(&agent_id) + .copied() + .unwrap_or_default() +} + +pub(super) fn affinity_delta(strictness: &AffinityStrictness, follows_rule: bool) -> Option { + match strictness { + AffinityStrictness::Required if !follows_rule => None, + AffinityStrictness::Required => Some(0), + AffinityStrictness::Preferred { weight } => { + Some(i64::from(follows_rule).saturating_mul(*weight)) + } + } +} + +pub(super) fn evaluate_affinity_rule( + metadata_tables: &MetadataTables<'_>, + rule: &crate::manifest::AffinityRule, +) -> bool { + let mut follows_rule = false; + + for requirement in &rule.requirements { + let mut requirement_outcome = !metadata_tables.is_empty(); + + for object_metadata in metadata_tables { + let table = match requirement.table { + MetadataTable::Label => object_metadata.0, + MetadataTable::Annotation => object_metadata.1, + }; + + if !evaluate_table_value(table.get(&requirement.key), requirement) { + requirement_outcome = false; + break; + } + } + + if requirement_outcome { + follows_rule = true; + break; + } + } + + if rule.inverse { + !follows_rule + } else { + follows_rule + } +} + +pub(super) fn evaluate_table_value( + value_option: Option<&String>, + requirement: &AffinityRequirement, +) -> bool { + let Some(value) = value_option else { + return matches!(requirement.operator, Operator::NotIn); + }; + + match requirement.operator { + Operator::In => requirement.values.contains(value), + Operator::NotIn => !requirement.values.contains(value), + Operator::Lt | Operator::Gt => { + let [requirement_value] = &requirement.values[..] else { + return false; + }; + + let Ok(value_number): Result = value.parse() else { + return false; + }; + + let Ok(requirement_value_number): Result = requirement_value.parse() else { + return false; + }; + + if requirement.operator == Operator::Lt { + value_number < requirement_value_number + } else { + value_number > requirement_value_number + } + } + } +} +#[derive(Debug, Clone, Copy, PartialEq)] +struct AgentScore { + general: f32, + affinity: i64, +} + +impl AgentScore { + pub const REJECTED: Self = Self { + general: f32::NEG_INFINITY, + affinity: i64::MIN, + }; +} + +impl Default for AgentScore { + fn default() -> Self { + Self { + general: 0.0, + affinity: 0, + } + } +} + +impl PartialOrd for AgentScore { + fn partial_cmp(&self, other: &Self) -> Option { + let affinity_cmp = self.affinity.cmp(&other.affinity); + + if affinity_cmp != Ordering::Equal { + return Some(affinity_cmp); + } + + self.general.partial_cmp(&other.general) + } +} diff --git a/odorobo/src/actors/scheduler_actor/tests.rs b/odorobo/src/actors/scheduler_actor/tests.rs new file mode 100644 index 0000000..1d934ea --- /dev/null +++ b/odorobo/src/actors/scheduler_actor/tests.rs @@ -0,0 +1,383 @@ +use super::{ + CachedActorKind, CachedVMActor, SchedulerActor, VmLifecycle, VmPlacement, + scheduling::{ + affinity_delta, evaluate_affinity_rule, evaluate_table_value, has_capacity, + pending_resources_for_agent, + }, +}; + +use crate::manifest::{ + AffinityRequirement, AffinityRule, AffinityStrictness, AffinityType, Boot, Compute, + DesiredState, Metadata, MetadataTable, Operator, VmManifest, +}; +use crate::messages::agent::AgentStatus; +use crate::types::ObjectMetadata; +use ahash::AHashMap; +use bytesize::ByteSize; +use std::collections::BTreeMap; +use std::time::{Duration, Instant}; +use ulid::Ulid; + +fn test_manifest(vcpus: u32, memory_bytes: u64) -> VmManifest { + VmManifest { + api_version: crate::manifest::MANIFEST_VERSION, + id: Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"), + desired: DesiredState { + metadata: Metadata { + name: "test".to_owned(), + ..Default::default() + }, + compute: Compute { + vcpus, + memory_bytes, + ..Default::default() + }, + boot: Boot::default(), + ..Default::default() + }, + observed: None, + } +} + +fn requirement(operator: Operator, values: &[&str]) -> AffinityRequirement { + AffinityRequirement { + key: "tier".to_owned(), + table: MetadataTable::Label, + operator, + values: values.iter().map(|value| (*value).to_owned()).collect(), + } +} + +#[test] +fn removes_expired_unresolved_vm_placeholders() { + let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); + let agent_id = super::ActorId::new(1); + let mut placements: AHashMap> = AHashMap::new(); + placements.insert( + vmid, + vec![VmPlacement { + agent_id, + lifecycle: VmLifecycle::Pending, + created_at: Instant::now() + .checked_sub(Duration::from_secs(31)) + .expect("test timestamp should be representable"), + last_confirmed_at: None, + }], + ); + let mut manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); + let mut data_cache: AHashMap> = AHashMap::new(); + + SchedulerActor::cleanup_unresolved_vm_cache(&mut manifests, &mut placements, &mut data_cache); + + assert!(!manifests.contains_key(&vmid)); + assert!(!placements.contains_key(&vmid)); + assert!(!data_cache.contains_key(&vmid)); +} + +#[test] +fn reserves_pending_vm_resources_until_agent_status_confirms_them() { + let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); + let agent_id = super::ActorId::new(1); + let config = test_manifest(4, ByteSize::gib(8).as_u64()); + let mut placements: AHashMap> = AHashMap::new(); + placements.insert( + vmid, + vec![VmPlacement { + agent_id, + lifecycle: VmLifecycle::Pending, + created_at: Instant::now(), + last_confirmed_at: None, + }], + ); + + let manifests = AHashMap::from([(vmid, config)]); + assert_eq!( + pending_resources_for_agent(&manifests, &placements, agent_id), + (4, ByteSize::gib(8).as_u64()) + ); +} + +#[test] +fn reconciling_source_agent_preserves_destination_migration_placement() { + let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); + let source_agent = super::ActorId::new(1); + let destination_agent = super::ActorId::new(2); + let running_placement = |agent_id| VmPlacement { + agent_id, + lifecycle: VmLifecycle::Running, + created_at: Instant::now(), + last_confirmed_at: Some(Instant::now()), + }; + let manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); + let mut placements = AHashMap::from([(vmid, vec![running_placement(source_agent)])]); + let destination_status = AgentStatus { + hostname: "destination".to_owned(), + vcpus: 1, + ram: ByteSize::b(1), + used_vcpus: 0, + used_ram: ByteSize::b(0), + vms: vec![vmid], + metadata: ObjectMetadata::default(), + }; + SchedulerActor::reconcile_agent_placements( + destination_agent, + &destination_status, + &manifests, + &mut placements, + ); + assert_eq!(manifests.len(), 1); + + let source_status = AgentStatus { + hostname: "source".to_owned(), + vcpus: 1, + ram: ByteSize::b(1), + used_vcpus: 0, + used_ram: ByteSize::b(0), + vms: Vec::new(), + metadata: ObjectMetadata::default(), + }; + + SchedulerActor::reconcile_agent_placements( + source_agent, + &source_status, + &manifests, + &mut placements, + ); + + let remaining = placements + .get(&vmid) + .expect("destination placement remains"); + assert_eq!(remaining.len(), 2); + assert!(remaining.iter().any(|entry| entry.agent_id == source_agent)); + assert!( + remaining + .iter() + .any(|entry| entry.agent_id == destination_agent) + ); +} + +#[test] +fn agent_removal_preserves_desired_placement() { + let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); + let agent_id = super::ActorId::new(1); + let manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); + let mut placements = AHashMap::from([( + vmid, + vec![VmPlacement { + agent_id, + lifecycle: VmLifecycle::Running, + created_at: Instant::now(), + last_confirmed_at: Some(Instant::now()), + }], + )]); + + SchedulerActor::reconcile_agent_delta(agent_id, &[], &[vmid], &manifests, &mut placements); + + assert_eq!(placements[&vmid].len(), 1); + assert_eq!(placements[&vmid][0].agent_id, agent_id); +} + +#[test] +fn failed_create_rolls_back_state_without_an_actor() { + let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); + let mut manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); + let mut actor_map = AHashMap::new(); + let mut placements = AHashMap::from([( + vmid, + vec![VmPlacement { + agent_id: super::ActorId::new(1), + lifecycle: VmLifecycle::Pending, + created_at: Instant::now(), + last_confirmed_at: None, + }], + )]); + let mut data_cache = AHashMap::from([(vmid, vec![CachedVMActor { actor_ref: None }])]); + + SchedulerActor::rollback_failed_create( + vmid, + false, + None, + &mut actor_map, + &mut manifests, + &mut placements, + &mut data_cache, + ); + + assert!(!manifests.contains_key(&vmid)); + assert!(!placements.contains_key(&vmid)); + assert!(!data_cache.contains_key(&vmid)); +} + +#[test] +fn failed_create_keeps_state_if_actor_exists() { + let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); + let mut manifests = AHashMap::from([(vmid, test_manifest(1, 1))]); + let mut actor_map = AHashMap::new(); + let mut placements = AHashMap::new(); + let mut data_cache = AHashMap::new(); + + SchedulerActor::rollback_failed_create( + vmid, + true, + None, + &mut actor_map, + &mut manifests, + &mut placements, + &mut data_cache, + ); + + assert!(manifests.contains_key(&vmid)); +} + +#[test] +fn agent_cleanup_does_not_remove_unrelated_vm_state() { + let agent_id = super::ActorId::new(1); + let vm_actor_id = super::ActorId::new(2); + let placement_agent_id = super::ActorId::new(3); + let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); + let mut scheduler = SchedulerActor { + agent_data_cache: AHashMap::new(), + agent_keepalive_tasks: AHashMap::new(), + vm_actorid_ulid_map: AHashMap::from([(vm_actor_id, vmid)]), + vm_manifests: AHashMap::from([(vmid, test_manifest(1, 1))]), + vm_placements: AHashMap::from([( + vmid, + vec![VmPlacement { + agent_id: placement_agent_id, + lifecycle: VmLifecycle::Running, + created_at: Instant::now(), + last_confirmed_at: Some(Instant::now()), + }], + )]), + vm_data_cache: AHashMap::from([(vmid, vec![CachedVMActor { actor_ref: None }])]), + vm_keepalive_tasks: AHashMap::new(), + pending_resources_cache: None, + agent_vm_index: AHashMap::new(), + actor_kinds: AHashMap::from([(agent_id, CachedActorKind::Agent)]), + cache_actor_finder: None, + }; + + scheduler.cleanup_agent_actor(agent_id); + + assert!(scheduler.vm_placements.contains_key(&vmid)); + assert!(scheduler.vm_manifests.contains_key(&vmid)); + assert!(scheduler.vm_actorid_ulid_map.contains_key(&vm_actor_id)); + assert!(scheduler.vm_data_cache.contains_key(&vmid)); +} + +#[test] +fn vm_cleanup_does_not_remove_unrelated_agent_state() { + let agent_id = super::ActorId::new(1); + let vm_actor_id = super::ActorId::new(2); + let vmid = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").expect("valid ulid"); + let mut scheduler = SchedulerActor { + agent_data_cache: AHashMap::new(), + agent_keepalive_tasks: AHashMap::new(), + vm_actorid_ulid_map: AHashMap::from([(vm_actor_id, vmid)]), + vm_manifests: AHashMap::from([(vmid, test_manifest(1, 1))]), + vm_placements: AHashMap::new(), + vm_data_cache: AHashMap::from([(vmid, vec![CachedVMActor { actor_ref: None }])]), + vm_keepalive_tasks: AHashMap::new(), + pending_resources_cache: None, + agent_vm_index: AHashMap::new(), + actor_kinds: AHashMap::from([(agent_id, CachedActorKind::Agent)]), + cache_actor_finder: None, + }; + + scheduler.cleanup_vm_actor(vm_actor_id); + + assert!(scheduler.actor_kinds.contains_key(&agent_id)); + assert!(!scheduler.vm_manifests.contains_key(&vmid)); + assert!(!scheduler.vm_data_cache.contains_key(&vmid)); +} + +#[test] +fn evaluates_membership_and_missing_keys() { + let metadata = BTreeMap::from([("tier".to_owned(), "frontend".to_owned())]); + assert!(evaluate_table_value( + metadata.get("tier"), + &requirement(Operator::In, &["frontend", "api"]) + )); + assert!(!evaluate_table_value( + metadata.get("tier"), + &requirement(Operator::In, &["backend"]) + )); + assert!(!evaluate_table_value( + None, + &requirement(Operator::In, &["frontend"]) + )); + assert!(evaluate_table_value( + None, + &requirement(Operator::NotIn, &["frontend"]) + )); +} + +#[test] +fn evaluates_not_in_and_numeric_comparisons() { + let metadata = BTreeMap::from([("tier".to_owned(), "4".to_owned())]); + assert!(evaluate_table_value( + metadata.get("tier"), + &requirement(Operator::NotIn, &["5"]) + )); + assert!(evaluate_table_value( + metadata.get("tier"), + &requirement(Operator::Lt, &["5"]) + )); + assert!(evaluate_table_value( + metadata.get("tier"), + &requirement(Operator::Gt, &["3"]) + )); + assert!(!evaluate_table_value( + metadata.get("tier"), + &requirement(Operator::Lt, &["4", "5"]) + )); + assert!(!evaluate_table_value( + metadata.get("tier"), + &requirement(Operator::Gt, &["not-a-number"]) + )); +} + +#[test] +fn evaluates_inverse_and_empty_requirements() { + let metadata = ObjectMetadata { + labels: BTreeMap::from([("tier".to_owned(), "frontend".to_owned())]), + annotations: BTreeMap::new(), + }; + let rule = AffinityRule { + strictness: AffinityStrictness::Required, + affinity_type: AffinityType::Agent, + inverse: true, + requirements: vec![requirement(Operator::In, &["frontend"])], + }; + assert!(!evaluate_affinity_rule( + &[(&metadata.labels, &metadata.annotations)], + &rule, + )); + + let empty_rule = AffinityRule { + strictness: AffinityStrictness::Required, + affinity_type: AffinityType::Agent, + inverse: false, + requirements: Vec::new(), + }; + assert!(!evaluate_affinity_rule(&[], &empty_rule)); +} + +#[test] +fn evaluates_required_preferred_and_capacity_rules() { + assert_eq!(affinity_delta(&AffinityStrictness::Required, true), Some(0)); + assert_eq!(affinity_delta(&AffinityStrictness::Required, false), None); + assert_eq!( + affinity_delta(&AffinityStrictness::Preferred { weight: 7 }, true), + Some(7) + ); + assert_eq!( + affinity_delta(&AffinityStrictness::Preferred { weight: 7 }, false), + Some(0) + ); + assert!(has_capacity(8, 2, 2, 16, 4, 4)); + assert!(has_capacity(8, 6, 2, 16, 4, 4)); + assert!(has_capacity(8, 2, 2, 16, 12, 4)); + assert!(!has_capacity(8, 7, 2, 16, 4, 4)); + assert!(!has_capacity(8, 2, 2, 16, 13, 4)); +} From f14320ee974599d874ab5800a3616c2b74243d71 Mon Sep 17 00:00:00 2001 From: Willow C Reed Date: Thu, 27 Aug 2026 17:15:01 -0600 Subject: [PATCH 2/7] fix clippy issues on new version --- odorobo/src/main.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/odorobo/src/main.rs b/odorobo/src/main.rs index 8e5a3f5..5ae79f8 100644 --- a/odorobo/src/main.rs +++ b/odorobo/src/main.rs @@ -1,3 +1,6 @@ +#![allow(unknown_lints)] // Supports Clippy releases before `unused_async_trait_impl` was introduced. +#![allow(clippy::unused_async_trait_impl)] // Kameo requires async trait methods even when a handler has no await points. + pub mod actors; mod ch_driver; pub mod config; From 106c58aca2d6c092b1aac78ed453201116e6fa26 Mon Sep 17 00:00:00 2001 From: Willow C Reed Date: Thu, 27 Aug 2026 17:19:27 -0600 Subject: [PATCH 3/7] improve documentation --- odorobo/src/actors/scheduler_actor.rs | 13 ++++++ odorobo/src/actors/scheduler_actor/cache.rs | 36 ++++++++++++++++ .../src/actors/scheduler_actor/discovery.rs | 20 +++++++++ .../src/actors/scheduler_actor/handlers.rs | 23 +++++++++-- .../src/actors/scheduler_actor/scheduling.rs | 41 ++++++++++++++++--- 5 files changed, 124 insertions(+), 9 deletions(-) diff --git a/odorobo/src/actors/scheduler_actor.rs b/odorobo/src/actors/scheduler_actor.rs index 505cbda..c289349 100644 --- a/odorobo/src/actors/scheduler_actor.rs +++ b/odorobo/src/actors/scheduler_actor.rs @@ -36,33 +36,39 @@ use crate::manifest::VmManifest; use crate::messages::agent::{AgentStatus, AgentStatusUpdate}; use crate::messages::vm::GetVMInfoReply; +/// Internal discovery event that starts VM polling for a newly found actor. #[derive(Debug)] struct VmActorDiscovered { actor_ref: RemoteActorRef, } +/// Internal discovery event that starts status polling for a newly found agent. #[derive(Debug)] struct AgentActorDiscovered { actor_ref: RemoteActorRef, } +/// The cache domain that owns cleanup for a linked remote actor. #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum CachedActorKind { Agent, Vm, } +/// A VM's initial identity and configuration snapshot from its updater task. #[derive(Debug)] struct VmUpdated { actor_ref: RemoteActorRef, data: GetVMInfoReply, } +/// Notification that a VM updater exceeded its reachability-failure budget. #[derive(Debug)] struct VmUpdaterStopped { actor_id: ActorId, } +/// A revisioned agent-status update forwarded by its updater task. #[derive(Debug)] struct AgentUpdated { actor_id: ActorId, @@ -70,11 +76,13 @@ struct AgentUpdated { update: AgentStatusUpdate, } +/// Notification that an agent updater exceeded its reachability-failure budget. #[derive(Debug)] struct AgentUpdaterStopped { actor_id: ActorId, } +/// Periodic maintenance trigger for expiring unresolved VM placements. #[derive(Debug)] struct ReconcileVmPlacements; @@ -142,8 +150,13 @@ pub struct SchedulerActor { pub vm_data_cache: AHashMap>, /// Polling tasks that refresh corresponding VM actor cache entries. pub vm_keepalive_tasks: AHashMap>, + /// Lazily computed resources reserved by `Pending` placements, keyed by agent. + /// Invalidated whenever status or placement state that affects capacity changes. pending_resources_cache: Option>, + /// Status-derived VM membership index used to evaluate VM affinity efficiently. + /// This is an optimization and must never override `vm_placements` intent. agent_vm_index: AHashMap>, + /// Classifies linked actors so link-death cleanup affects the owning cache only. actor_kinds: AHashMap, /// Background discovery and reconciliation task. pub cache_actor_finder: Option>, diff --git a/odorobo/src/actors/scheduler_actor/cache.rs b/odorobo/src/actors/scheduler_actor/cache.rs index 00c26dd..c2046d7 100644 --- a/odorobo/src/actors/scheduler_actor/cache.rs +++ b/odorobo/src/actors/scheduler_actor/cache.rs @@ -16,12 +16,18 @@ use super::{CachedVMActor, SchedulerActor, VmLifecycle, VmPlacement}; const UNRESOLVED_VM_CACHE_TIMEOUT: Duration = Duration::from_secs(30); impl SchedulerActor { + /// Releases excess allocation once an entry list no longer represents migration. pub(super) fn shrink_non_migrating_entries(entries: &mut Vec) { if entries.len() <= 1 && entries.capacity() >= 4 { entries.shrink_to_fit(); } } + /// Replaces a matching discovered VM entry, fulfills one unresolved placeholder, + /// or records an additional entry for a concurrent migration placement. + /// + /// Matching is ordered by actor ID, then the first `None` placeholder, then + /// append. Placement and cache vectors are not positionally correlated. pub(super) fn update_cached_vm_entry( entries: &mut Vec, actor_id: ActorId, @@ -42,6 +48,12 @@ impl SchedulerActor { Self::shrink_non_migrating_entries(entries); } + /// Expires pending placements that were never confirmed by agent status. + /// + /// Pending entries expire after 30 seconds. The five-second discovery loop + /// triggers this maintenance, so a slow-to-report create can be forgotten. + /// When the expired placement was the last placement for a VM, all correlated + /// manifest, placement, and actor-cache state is removed. pub(super) fn cleanup_unresolved_vm_cache( manifests: &mut AHashMap, placements: &mut AHashMap>, @@ -67,6 +79,10 @@ impl SchedulerActor { /// Returns every VM that could be on an agent, deduplicating each source in /// constant expected time. + /// + /// The union includes current status, the status-derived index, and desired + /// placements to avoid affinity decisions missing in-flight changes. Output is + /// observation-first, then hash-map iteration order; callers must not rely on it. pub(super) fn placement_vm_ids( placements: &AHashMap>, indexed: Option<&AHashSet>, @@ -102,6 +118,7 @@ impl SchedulerActor { vmids } + /// Removes all scheduler state correlated with a VM identifier. pub(super) fn remove_vm_state( vmid: Ulid, manifests: &mut AHashMap, @@ -113,6 +130,7 @@ impl SchedulerActor { data_cache.remove(&vmid); } + /// Removes a departed actor from every VM cache entry without altering placement intent. pub(super) fn remove_vm_actor( actor_id: ActorId, data_cache: &mut AHashMap>, @@ -135,6 +153,8 @@ impl SchedulerActor { } } + /// Removes placements assigned to a departed agent and drops VM state only + /// when no placement remains. pub(super) fn remove_agent_placements( agent_id: ActorId, manifests: &mut AHashMap, @@ -155,6 +175,10 @@ impl SchedulerActor { } } + /// Rolls back optimistic create state only when no VM actor was created. + /// + /// An actor may exist even if the create request failed or its reply was lost; + /// retaining the state in that case lets normal discovery reconcile it. pub(super) fn rollback_failed_create( vmid: Ulid, actor_exists: bool, @@ -174,6 +198,7 @@ impl SchedulerActor { } } + /// Aborts polling and removes all cache state owned by a departed agent. pub(super) fn cleanup_agent_actor(&mut self, actor_id: ActorId) { if let Some(keepalive_task) = self.agent_keepalive_tasks.remove(&actor_id) { trace!(?actor_id, "Aborting agent keepalive task"); @@ -190,6 +215,8 @@ impl SchedulerActor { ); } + /// Aborts VM polling and removes actor state, retaining a VM only when + /// another discovered actor or unresolved placement can still represent it. pub(super) fn cleanup_vm_actor(&mut self, actor_id: ActorId) { if let Some(keepalive_task) = self.vm_keepalive_tasks.remove(&actor_id) { trace!(?actor_id, "Aborting VM keepalive task"); @@ -213,6 +240,10 @@ impl SchedulerActor { } } + /// Incorporates additions from a status delta into placement observations. + /// + /// Removals deliberately do not delete desired placements: they may be + /// transient observations and reconciliation must be able to recreate the VM. pub(super) fn reconcile_agent_delta( agent_id: ActorId, added: &[Ulid], @@ -245,6 +276,11 @@ impl SchedulerActor { // distinction for snapshots. } + /// Reconciles scheduler placement observations with a complete agent snapshot. + /// + /// Known, reported VMs gain or refresh `Running` placements. Unknown VMs are + /// ignored because the scheduler has no retained intent for them. Absent VMs + /// leave existing desired placements intact so future reconciliation can act. pub(super) fn reconcile_agent_placements( agent_id: ActorId, status: &AgentStatus, diff --git a/odorobo/src/actors/scheduler_actor/discovery.rs b/odorobo/src/actors/scheduler_actor/discovery.rs index c1f562a..a7c7fc6 100644 --- a/odorobo/src/actors/scheduler_actor/discovery.rs +++ b/odorobo/src/actors/scheduler_actor/discovery.rs @@ -20,6 +20,7 @@ use super::{ }; impl SchedulerActor { + /// Enumerates currently discoverable VM actors and forwards each to the scheduler. async fn vm_actor_finder(parent_actor_ref: ActorRef) -> Result<(), Report> { trace!("running vm_actor_finder"); @@ -37,6 +38,9 @@ impl SchedulerActor { Ok(()) } + /// Resolves a VM's identity once, then heartbeats it every second. + /// + /// Six consecutive failed requests stop the task and request cache cleanup. async fn vm_updater_task(scheduler: ActorRef, actor_ref: RemoteActorRef) { let mut interval = tokio::time::interval(Duration::from_secs(1)); let mut fails: u8 = 0; @@ -91,6 +95,7 @@ impl SchedulerActor { } } + /// Enumerates currently discoverable agents and forwards each to the scheduler. async fn agent_actor_finder(parent_actor_ref: ActorRef) -> Result<(), Report> { trace!("running agent_actor_finder"); @@ -108,6 +113,10 @@ impl SchedulerActor { Ok(()) } + /// Polls revisioned agent status every second and forwards accepted responses. + /// + /// The task begins with revision zero so its first response establishes a + /// full snapshot. Six consecutive failed requests stop it and request cleanup. async fn agent_updater_task(scheduler: ActorRef, actor_ref: RemoteActorRef) { let mut interval = tokio::time::interval(Duration::from_secs(1)); let mut status_revision = 0; @@ -166,6 +175,10 @@ impl SchedulerActor { } } + /// Starts the periodic discovery and pending-placement maintenance loop. + /// + /// Each five-second pass discovers both actor types, then asks the scheduler + /// to expire unresolved pending placements. pub(super) fn start_actor_finder(&mut self, actor_ref: ActorRef) { self.cache_actor_finder = Some(tokio::spawn(async move { let mut interval = tokio::time::interval(Duration::from_secs(5)); @@ -185,6 +198,7 @@ impl SchedulerActor { } } +/// Registers a discovered VM actor and starts its single polling task. impl Message for SchedulerActor { type Reply = (); @@ -212,6 +226,7 @@ impl Message for SchedulerActor { } } +/// Registers a discovered agent and starts its single status-polling task. impl Message for SchedulerActor { type Reply = (); @@ -239,6 +254,7 @@ impl Message for SchedulerActor { } } +/// Caches a VM's canonical ID, manifest when supplied, and discovered actor reference. impl Message for SchedulerActor { type Reply = (); @@ -257,6 +273,7 @@ impl Message for SchedulerActor { } } +/// Removes cache state after a VM updater determines its actor is unreachable. impl Message for SchedulerActor { type Reply = (); @@ -282,6 +299,7 @@ impl Message for SchedulerActor { } } +/// Applies a newer agent update, requiring a full snapshot before accepting deltas. impl Message for SchedulerActor { type Reply = (); @@ -355,6 +373,7 @@ impl Message for SchedulerActor { } } +/// Removes cache state and placements owned by an unreachable agent. impl Message for SchedulerActor { type Reply = (); @@ -373,6 +392,7 @@ impl Message for SchedulerActor { } } +/// Performs periodic pending-placement expiry and refreshes resource accounting. impl Message for SchedulerActor { type Reply = (); diff --git a/odorobo/src/actors/scheduler_actor/handlers.rs b/odorobo/src/actors/scheduler_actor/handlers.rs index 8a9876b..8611ece 100644 --- a/odorobo/src/actors/scheduler_actor/handlers.rs +++ b/odorobo/src/actors/scheduler_actor/handlers.rs @@ -19,6 +19,10 @@ use crate::utils::actor_names::vm_actor_id; use super::{CachedActorKind, CachedVMActor, SchedulerActor, VmLifecycle, VmPlacement}; +/// Owns scheduler initialization and cleanup for linked remote actors. +/// +/// Losing an agent removes its placements. Losing a VM removes only its actor +/// cache entry unless no discovered actor or placeholder remains for that VM. impl Actor for SchedulerActor { type Args = (); type Error = Report; @@ -71,6 +75,11 @@ impl Actor for SchedulerActor { Ok(ControlFlow::Continue(())) } } +/// Optimistically reserves a placement, then forwards VM creation to the chosen agent. +/// +/// A successful reply confirms agent acceptance, not observed VM execution. A failed +/// request is rolled back only when discovery cannot find a VM actor, preserving state +/// for eventual reconciliation when the request result was lost or delayed. impl Message for SchedulerActor { type Reply = Result; @@ -173,6 +182,9 @@ impl Message for SchedulerActor { } } +/// Looks up a VM actor and forwards deletion without eagerly altering scheduler caches. +/// +/// Cache cleanup waits for actor link death or updater reachability failure. impl Message for SchedulerActor { type Reply = Result; @@ -193,6 +205,9 @@ impl Message for SchedulerActor { } } +/// Looks up a VM actor and forwards shutdown without eagerly altering scheduler caches. +/// +/// Cache cleanup waits for actor link death or updater reachability failure. impl Message for SchedulerActor { type Reply = Result; @@ -213,9 +228,10 @@ impl Message for SchedulerActor { } } -/// this only gets data from the cache from agents -/// we may need a different message that actually forcibly runs/updates everything. -/// and/or messages that get data directly from the `VMActors`. +/// Returns the concatenated VM IDs from cached agent status snapshots. +/// +/// This is a potentially stale, non-deduplicated observation rather than an +/// authoritative inventory; it performs neither polling nor direct VM-actor queries. impl Message for SchedulerActor { type Reply = Result; @@ -239,6 +255,7 @@ impl Message for SchedulerActor { } } +/// Provides scheduler actor liveness only; it does not imply cache freshness or readiness. impl Message for SchedulerActor { type Reply = Pong; diff --git a/odorobo/src/actors/scheduler_actor/scheduling.rs b/odorobo/src/actors/scheduler_actor/scheduling.rs index 30051e6..be07cce 100644 --- a/odorobo/src/actors/scheduler_actor/scheduling.rs +++ b/odorobo/src/actors/scheduler_actor/scheduling.rs @@ -24,6 +24,10 @@ type MetadataTables<'a> = [( )]; impl SchedulerActor { + /// Returns resources reserved by `Pending` placements, computing them lazily. + /// + /// Running usage comes from cached agent status; only unconfirmed placements + /// are added here to prevent concurrent create requests from overcommitting. pub(super) fn pending_resources(&mut self) -> &AHashMap { if self.pending_resources_cache.is_none() { self.pending_resources_cache = Some(pending_resources_by_agent( @@ -36,15 +40,19 @@ impl SchedulerActor { .expect("pending resources cache was just initialized") } + /// Marks pending-placement resource totals stale after a relevant cache change. pub(super) fn invalidate_pending_resources(&mut self) { self.pending_resources_cache = None; } /// Determine the best agent to schedule a specific VM creation request to. /// - /// The scheduler first filters agents by capacity and affinity requirements, - /// then uses affinity and general resource scores to select the best match. - /// Affinity rules are roughly based on + /// The scheduler first filters agents by capacity and required affinity, then + /// ranks preferred-affinity totals before CPU/RAM headroom. Equal scores have + /// no deterministic tie-breaker because the agent cache is a hash map. + /// + /// Scoring performs no network I/O and therefore uses an eventually consistent + /// cache. Affinity rules are roughly based on /// . pub(super) fn schedule_agent( &mut self, @@ -98,9 +106,13 @@ impl SchedulerActor { .collect() } - // this function intentionally only checks against the cache. this has some positives and negatives: - // positive: it will never trigger any network requests so its very fast, and having to do network requests for scoring whenever we want to schedule a vm is likely a bad idea - // negative: it technically has a delayed view of the cluster, meaning that some things that happened in the future, may not exist yet. so we need to be careful about how this is done so affinity rules are not accidentally broken. mostly this means, if we do anything that could affect the outcome of an affinity rule (ex: network request to an agent), we need to update the cache, before we do the action. + /// Scores one agent from cached state without performing network I/O. + /// + /// vCPUs may use 2× overprovisioning; RAM is not overprovisioned. Pending + /// placements reserve their requested resources until agent status confirms + /// them, while running usage comes directly from the cached agent snapshot. + /// Cache-affecting actions must update cache state before a later affinity + /// decision depends on them. fn score_agent( &self, msg: &CreateVM, @@ -214,6 +226,9 @@ impl SchedulerActor { } } +/// Returns whether a request fits within supplied vCPU and RAM limits. +/// +/// Arguments must already include applicable overprovisioning and pending reservations. pub(super) const fn has_capacity( max_vcpus: u32, used_vcpus: u32, @@ -260,6 +275,9 @@ pub(super) fn pending_resources_for_agent( .unwrap_or_default() } +/// Converts an affinity result into a score contribution or rejection. +/// +/// Failed required rules reject an agent; unmet preferred rules contribute zero. pub(super) fn affinity_delta(strictness: &AffinityStrictness, follows_rule: bool) -> Option { match strictness { AffinityStrictness::Required if !follows_rule => None, @@ -270,6 +288,12 @@ pub(super) fn affinity_delta(strictness: &AffinityStrictness, follows_rule: bool } } +/// Evaluates a rule against the metadata objects in its selected affinity scope. +/// +/// Requirements are OR-ed. For a requirement to match, every supplied metadata +/// object must satisfy it; empty metadata therefore does not match. `inverse` +/// negates the aggregate result, and a rule without requirements is false before +/// inversion. pub(super) fn evaluate_affinity_rule( metadata_tables: &MetadataTables<'_>, rule: &crate::manifest::AffinityRule, @@ -304,6 +328,10 @@ pub(super) fn evaluate_affinity_rule( } } +/// Evaluates one metadata value against a requirement. +/// +/// Missing keys match only `NotIn`. Numeric `Lt` and `Gt` comparisons require +/// exactly one parseable numeric requirement value; malformed comparisons are false. pub(super) fn evaluate_table_value( value_option: Option<&String>, requirement: &AffinityRequirement, @@ -336,6 +364,7 @@ pub(super) fn evaluate_table_value( } } } +/// Lexicographic scheduling score: affinity takes precedence over resource headroom. #[derive(Debug, Clone, Copy, PartialEq)] struct AgentScore { general: f32, From b98e11bca36112db184d0d73456fed02875a814e Mon Sep 17 00:00:00 2001 From: Willow C Reed Date: Thu, 27 Aug 2026 17:21:41 -0600 Subject: [PATCH 4/7] add todos --- odorobo/src/actors/scheduler_actor/cache.rs | 4 ++++ odorobo/src/actors/scheduler_actor/discovery.rs | 7 +++++++ odorobo/src/actors/scheduler_actor/handlers.rs | 2 ++ 3 files changed, 13 insertions(+) diff --git a/odorobo/src/actors/scheduler_actor/cache.rs b/odorobo/src/actors/scheduler_actor/cache.rs index c2046d7..5230832 100644 --- a/odorobo/src/actors/scheduler_actor/cache.rs +++ b/odorobo/src/actors/scheduler_actor/cache.rs @@ -155,6 +155,8 @@ impl SchedulerActor { /// Removes placements assigned to a departed agent and drops VM state only /// when no placement remains. + // TODO: Preserve VM intent and enqueue replacement placement or recreation + // when an agent disappears instead of dropping the last VM state. pub(super) fn remove_agent_placements( agent_id: ActorId, manifests: &mut AHashMap, @@ -311,6 +313,8 @@ impl SchedulerActor { let empty_vmids: Vec<_> = placements .iter_mut() .filter_map(|(vmid, entries)| { + // TODO: Expire or repair `Running` placements that remain absent + // from repeated full snapshots; they are retained indefinitely now. for entry in entries .iter_mut() .filter(|entry| entry.agent_id == agent_id) diff --git a/odorobo/src/actors/scheduler_actor/discovery.rs b/odorobo/src/actors/scheduler_actor/discovery.rs index a7c7fc6..9174876 100644 --- a/odorobo/src/actors/scheduler_actor/discovery.rs +++ b/odorobo/src/actors/scheduler_actor/discovery.rs @@ -280,6 +280,8 @@ impl Message for SchedulerActor { async fn handle(&mut self, msg: VmUpdaterStopped, _ctx: &mut Context) { self.vm_keepalive_tasks.remove(&msg.actor_id); self.actor_kinds.remove(&msg.actor_id); + // TODO: Invalidate pending-resource accounting after this cleanup, as + // link-death cleanup already does. let vmid = self.vm_actorid_ulid_map.remove(&msg.actor_id); Self::remove_vm_actor(msg.actor_id, &mut self.vm_data_cache); @@ -380,6 +382,8 @@ impl Message for SchedulerActor { async fn handle(&mut self, msg: AgentUpdaterStopped, _ctx: &mut Context) { self.agent_keepalive_tasks.remove(&msg.actor_id); self.actor_kinds.remove(&msg.actor_id); + // TODO: Delegate to `cleanup_agent_actor`, or perform its equivalent + // index and resource-accounting cleanup, before allowing rediscovery. self.agent_data_cache.remove(&msg.actor_id); self.agent_vm_index.remove(&msg.actor_id); Self::remove_agent_placements( @@ -403,5 +407,8 @@ impl Message for SchedulerActor { &mut self.vm_data_cache, ); self.invalidate_pending_resources(); + // TODO: Reconcile desired placements absent from agent status by choosing + // a healthy agent and issuing `CreateVM`; this currently only expires + // unconfirmed pending reservations. } } diff --git a/odorobo/src/actors/scheduler_actor/handlers.rs b/odorobo/src/actors/scheduler_actor/handlers.rs index 8611ece..48fde18 100644 --- a/odorobo/src/actors/scheduler_actor/handlers.rs +++ b/odorobo/src/actors/scheduler_actor/handlers.rs @@ -90,6 +90,8 @@ impl Message for SchedulerActor { ) -> Self::Reply { let target_agent = self.schedule_agent(&msg)?; + // TODO: Define duplicate VM-ID semantics before overwriting intent and + // appending another pending placement; reject conflicts or make retries idempotent. self.vm_manifests.insert(msg.vmid, msg.config.clone()); self.invalidate_pending_resources(); self.vm_placements From 3203788ee80ef7902fa4221f8fcd8e2bf89aedc1 Mon Sep 17 00:00:00 2001 From: Willow C Reed Date: Thu, 27 Aug 2026 18:23:03 -0600 Subject: [PATCH 5/7] update split scheduler for manifest contract --- odorobo/src/actors/scheduler_actor/scheduling.rs | 8 ++++---- odorobo/src/actors/scheduler_actor/tests.rs | 8 ++++---- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/odorobo/src/actors/scheduler_actor/scheduling.rs b/odorobo/src/actors/scheduler_actor/scheduling.rs index be07cce..9fea089 100644 --- a/odorobo/src/actors/scheduler_actor/scheduling.rs +++ b/odorobo/src/actors/scheduler_actor/scheduling.rs @@ -291,9 +291,9 @@ pub(super) fn affinity_delta(strictness: &AffinityStrictness, follows_rule: bool /// Evaluates a rule against the metadata objects in its selected affinity scope. /// /// Requirements are OR-ed. For a requirement to match, every supplied metadata -/// object must satisfy it; empty metadata therefore does not match. `inverse` -/// negates the aggregate result, and a rule without requirements is false before -/// inversion. +/// object must satisfy it; empty metadata therefore does not match. The `anti` +/// direction negates the aggregate result, and a rule without requirements is false +/// before direction is applied. pub(super) fn evaluate_affinity_rule( metadata_tables: &MetadataTables<'_>, rule: &crate::manifest::AffinityRule, @@ -321,7 +321,7 @@ pub(super) fn evaluate_affinity_rule( } } - if rule.inverse { + if matches!(rule.direction, crate::manifest::AffinityDirection::Anti) { !follows_rule } else { follows_rule diff --git a/odorobo/src/actors/scheduler_actor/tests.rs b/odorobo/src/actors/scheduler_actor/tests.rs index 1d934ea..309e1df 100644 --- a/odorobo/src/actors/scheduler_actor/tests.rs +++ b/odorobo/src/actors/scheduler_actor/tests.rs @@ -7,8 +7,8 @@ use super::{ }; use crate::manifest::{ - AffinityRequirement, AffinityRule, AffinityStrictness, AffinityType, Boot, Compute, - DesiredState, Metadata, MetadataTable, Operator, VmManifest, + AffinityDirection, AffinityRequirement, AffinityRule, AffinityStrictness, AffinityType, Boot, + Compute, DesiredState, Metadata, MetadataTable, Operator, VmManifest, }; use crate::messages::agent::AgentStatus; use crate::types::ObjectMetadata; @@ -346,7 +346,7 @@ fn evaluates_inverse_and_empty_requirements() { let rule = AffinityRule { strictness: AffinityStrictness::Required, affinity_type: AffinityType::Agent, - inverse: true, + direction: AffinityDirection::Anti, requirements: vec![requirement(Operator::In, &["frontend"])], }; assert!(!evaluate_affinity_rule( @@ -357,7 +357,7 @@ fn evaluates_inverse_and_empty_requirements() { let empty_rule = AffinityRule { strictness: AffinityStrictness::Required, affinity_type: AffinityType::Agent, - inverse: false, + direction: AffinityDirection::Normal, requirements: Vec::new(), }; assert!(!evaluate_affinity_rule(&[], &empty_rule)); From 7238c7cc368f457711c943c496e9d7cdc77206fc Mon Sep 17 00:00:00 2001 From: Willow C Reed Date: Thu, 27 Aug 2026 19:37:33 -0600 Subject: [PATCH 6/7] preserve affinity inversion in split scheduler --- .../src/actors/scheduler_actor/scheduling.rs | 6 +----- odorobo/src/actors/scheduler_actor/tests.rs | 19 +++++++++++++++---- 2 files changed, 16 insertions(+), 9 deletions(-) diff --git a/odorobo/src/actors/scheduler_actor/scheduling.rs b/odorobo/src/actors/scheduler_actor/scheduling.rs index 9fea089..ace48da 100644 --- a/odorobo/src/actors/scheduler_actor/scheduling.rs +++ b/odorobo/src/actors/scheduler_actor/scheduling.rs @@ -321,11 +321,7 @@ pub(super) fn evaluate_affinity_rule( } } - if matches!(rule.direction, crate::manifest::AffinityDirection::Anti) { - !follows_rule - } else { - follows_rule - } + follows_rule ^ rule.inverse } /// Evaluates one metadata value against a requirement. diff --git a/odorobo/src/actors/scheduler_actor/tests.rs b/odorobo/src/actors/scheduler_actor/tests.rs index 309e1df..62baeee 100644 --- a/odorobo/src/actors/scheduler_actor/tests.rs +++ b/odorobo/src/actors/scheduler_actor/tests.rs @@ -7,8 +7,8 @@ use super::{ }; use crate::manifest::{ - AffinityDirection, AffinityRequirement, AffinityRule, AffinityStrictness, AffinityType, Boot, - Compute, DesiredState, Metadata, MetadataTable, Operator, VmManifest, + AffinityRequirement, AffinityRule, AffinityStrictness, AffinityType, Boot, Compute, + DesiredState, Metadata, MetadataTable, Operator, VmManifest, }; use crate::messages::agent::AgentStatus; use crate::types::ObjectMetadata; @@ -346,7 +346,7 @@ fn evaluates_inverse_and_empty_requirements() { let rule = AffinityRule { strictness: AffinityStrictness::Required, affinity_type: AffinityType::Agent, - direction: AffinityDirection::Anti, + inverse: true, requirements: vec![requirement(Operator::In, &["frontend"])], }; assert!(!evaluate_affinity_rule( @@ -354,10 +354,21 @@ fn evaluates_inverse_and_empty_requirements() { &rule, )); + let non_matching_rule = AffinityRule { + strictness: AffinityStrictness::Required, + affinity_type: AffinityType::Agent, + inverse: true, + requirements: vec![requirement(Operator::In, &["backend"])], + }; + assert!(evaluate_affinity_rule( + &[(&metadata.labels, &metadata.annotations)], + &non_matching_rule, + )); + let empty_rule = AffinityRule { strictness: AffinityStrictness::Required, affinity_type: AffinityType::Agent, - direction: AffinityDirection::Normal, + inverse: false, requirements: Vec::new(), }; assert!(!evaluate_affinity_rule(&[], &empty_rule)); From 66db29c77c2054ec31553ca4bffa60b2a5c4b44a Mon Sep 17 00:00:00 2001 From: Cypress Reed Date: Mon, 31 Aug 2026 10:00:39 -0600 Subject: [PATCH 7/7] use proper cleanup routines --- .../src/actors/scheduler_actor/discovery.rs | 33 ++----------------- 1 file changed, 2 insertions(+), 31 deletions(-) diff --git a/odorobo/src/actors/scheduler_actor/discovery.rs b/odorobo/src/actors/scheduler_actor/discovery.rs index 9174876..4a79fb4 100644 --- a/odorobo/src/actors/scheduler_actor/discovery.rs +++ b/odorobo/src/actors/scheduler_actor/discovery.rs @@ -278,26 +278,8 @@ impl Message for SchedulerActor { type Reply = (); async fn handle(&mut self, msg: VmUpdaterStopped, _ctx: &mut Context) { - self.vm_keepalive_tasks.remove(&msg.actor_id); self.actor_kinds.remove(&msg.actor_id); - // TODO: Invalidate pending-resource accounting after this cleanup, as - // link-death cleanup already does. - let vmid = self.vm_actorid_ulid_map.remove(&msg.actor_id); - Self::remove_vm_actor(msg.actor_id, &mut self.vm_data_cache); - - if let Some(vmid) = vmid - && self - .vm_data_cache - .get(&vmid) - .is_none_or(|entries| entries.iter().all(|entry| entry.actor_ref.is_none())) - { - Self::remove_vm_state( - vmid, - &mut self.vm_manifests, - &mut self.vm_placements, - &mut self.vm_data_cache, - ); - } + self.cleanup_vm_actor(msg.actor_id); } } @@ -380,19 +362,8 @@ impl Message for SchedulerActor { type Reply = (); async fn handle(&mut self, msg: AgentUpdaterStopped, _ctx: &mut Context) { - self.agent_keepalive_tasks.remove(&msg.actor_id); self.actor_kinds.remove(&msg.actor_id); - // TODO: Delegate to `cleanup_agent_actor`, or perform its equivalent - // index and resource-accounting cleanup, before allowing rediscovery. - self.agent_data_cache.remove(&msg.actor_id); - self.agent_vm_index.remove(&msg.actor_id); - Self::remove_agent_placements( - msg.actor_id, - &mut self.vm_manifests, - &mut self.vm_placements, - &mut self.vm_data_cache, - ); - self.invalidate_pending_resources(); + self.cleanup_agent_actor(msg.actor_id); } }