use std::fmt;
use thiserror::Error;
#[derive(Debug, Error, PartialEq, Eq)]
pub enum GraphError {
#[error("node capacity exceeded: max={0}")]
NodeCapacityExceeded(usize),
#[error("mapping capacity exceeded: max={0}")]
MappingCapacityExceeded(usize),
#[error("node not found: id={0}")]
NodeNotFound(u64),
#[error("node already exists: id={0}")]
NodeAlreadyExists(u64),
#[error("domain isolation violation: node={node_id}, domain={domain:?}")]
DomainIsolationViolation {
node_id: u64,
domain: ExecutionDomain,
},
#[error("insufficient resources: domain={domain:?}, need_cpu={need_cpu}, have_cpu={have_cpu}, need_mem={need_mem}, have_mem={have_mem}")]
InsufficientResources {
domain: ExecutionDomain,
need_cpu: u32,
have_cpu: u32,
need_mem: u32,
have_mem: u32,
},
#[error("invalid topology: {0}")]
InvalidTopology(&'static str),
#[error("no workers available for domain: {0:?}")]
NoWorkersForDomain(ExecutionDomain),
#[error("mapping not found: queue_id={0}")]
MappingNotFound(u32),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum ExecutionDomain {
DataPlane,
ControlPlane,
Management,
Observability,
Security,
Storage,
}
impl ExecutionDomain {
pub fn all() -> [ExecutionDomain; 6] {
[
ExecutionDomain::DataPlane,
ExecutionDomain::ControlPlane,
ExecutionDomain::Management,
ExecutionDomain::Observability,
ExecutionDomain::Security,
ExecutionDomain::Storage,
]
}
#[inline]
pub fn is_data_plane(&self) -> bool {
matches!(self, ExecutionDomain::DataPlane)
}
#[inline]
pub fn requires_high_cpu(&self) -> bool {
matches!(self, ExecutionDomain::DataPlane | ExecutionDomain::Security)
}
#[inline]
pub fn requires_large_memory(&self) -> bool {
matches!(self, ExecutionDomain::Storage | ExecutionDomain::Observability)
}
pub fn default_requirements(&self) -> (u32, u32, u32) {
match self {
ExecutionDomain::DataPlane => (4, 4096, 8),
ExecutionDomain::ControlPlane => (2, 2048, 4),
ExecutionDomain::Management => (2, 2048, 2),
ExecutionDomain::Observability => (2, 8192, 4),
ExecutionDomain::Security => (4, 4096, 4),
ExecutionDomain::Storage => (2, 16384, 2),
}
}
}
impl fmt::Display for ExecutionDomain {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
ExecutionDomain::DataPlane => write!(f, "DataPlane"),
ExecutionDomain::ControlPlane => write!(f, "ControlPlane"),
ExecutionDomain::Management => write!(f, "Management"),
ExecutionDomain::Observability => write!(f, "Observability"),
ExecutionDomain::Security => write!(f, "Security"),
ExecutionDomain::Storage => write!(f, "Storage"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum NodeStatus {
Online,
Offline,
Degraded,
Maintenance,
}
impl NodeStatus {
#[inline]
pub fn is_schedulable(&self) -> bool {
matches!(self, NodeStatus::Online)
}
}
#[derive(Debug, Clone, Copy)]
pub struct ResourceNode {
pub id: u64,
pub domain: ExecutionDomain,
pub cpu_cores: u32,
pub memory_mb: u32,
pub nic_queues: u32,
pub status: NodeStatus,
pub numa_node: u32,
}
impl ResourceNode {
pub const fn new(
id: u64,
domain: ExecutionDomain,
cpu_cores: u32,
memory_mb: u32,
nic_queues: u32,
) -> Self {
Self {
id,
domain,
cpu_cores,
memory_mb,
nic_queues,
status: NodeStatus::Online,
numa_node: 0,
}
}
pub const fn with_numa(mut self, numa_node: u32) -> Self {
self.numa_node = numa_node;
self
}
pub const fn with_status(mut self, status: NodeStatus) -> Self {
self.status = status;
self
}
#[inline]
pub fn is_available(&self) -> bool {
self.status.is_schedulable()
}
#[inline]
pub fn meets_requirements(&self, cpu: u32, mem: u32, queues: u32) -> bool {
self.cpu_cores >= cpu
&& self.memory_mb >= mem
&& self.nic_queues >= queues
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum AffinityRule {
None,
NumaLocal,
CpuLocal,
Exact,
}
#[derive(Debug, Clone, Copy)]
pub struct QueueMapping {
pub queue_id: u32,
pub worker_id: u32,
pub affinity: AffinityRule,
}
impl QueueMapping {
pub const fn new(queue_id: u32, worker_id: u32) -> Self {
Self {
queue_id,
worker_id,
affinity: AffinityRule::None,
}
}
pub const fn with_affinity(mut self, affinity: AffinityRule) -> Self {
self.affinity = affinity;
self
}
}
#[derive(Debug, Clone)]
pub struct WorkerAssignment {
pub worker_id: u32,
pub node_id: u64,
pub domain: ExecutionDomain,
pub numa_node: u32,
pub queue_ids: Vec<u32>,
}
#[derive(Debug, Clone, Copy)]
pub struct DomainSnapshot {
pub domain: ExecutionDomain,
pub required_cpu_cores: u32,
pub required_memory_mb: u32,
pub required_nic_queues: u32,
pub worker_count: usize,
}
impl DomainSnapshot {
pub fn from_domain(domain: ExecutionDomain, worker_count: usize) -> Self {
let (cpu, mem, queues) = domain.default_requirements();
Self {
domain,
required_cpu_cores: cpu,
required_memory_mb: mem,
required_nic_queues: queues,
worker_count,
}
}
pub const fn with_cpu(mut self, cpu: u32) -> Self {
self.required_cpu_cores = cpu;
self
}
pub const fn with_memory(mut self, mem: u32) -> Self {
self.required_memory_mb = mem;
self
}
pub const fn with_workers(mut self, count: usize) -> Self {
self.worker_count = count;
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum QueueStrategy {
RssHash,
RoundRobin,
NumaLocal,
}
#[derive(Debug)]
pub struct RuntimeGraph<const MAX_NODES: usize, const MAX_MAPPINGS: usize> {
nodes: [Option<ResourceNode>; MAX_NODES],
mappings: [Option<QueueMapping>; MAX_MAPPINGS],
node_count: usize,
mapping_count: usize,
}
impl<const MAX_NODES: usize, const MAX_MAPPINGS: usize> Default
for RuntimeGraph<MAX_NODES, MAX_MAPPINGS>
{
fn default() -> Self {
Self::new()
}
}
impl<const MAX_NODES: usize, const MAX_MAPPINGS: usize>
RuntimeGraph<MAX_NODES, MAX_MAPPINGS>
{
pub const fn new() -> Self {
Self {
nodes: [const { None }; MAX_NODES],
mappings: [const { None }; MAX_MAPPINGS],
node_count: 0,
mapping_count: 0,
}
}
#[inline]
pub fn node_count(&self) -> usize {
self.node_count
}
#[inline]
pub fn mapping_count(&self) -> usize {
self.mapping_count
}
#[inline]
pub const fn max_nodes(&self) -> usize {
MAX_NODES
}
#[inline]
pub const fn max_mappings(&self) -> usize {
MAX_MAPPINGS
}
#[inline]
pub fn is_empty(&self) -> bool {
self.node_count == 0
}
pub fn add_node(&mut self, node: ResourceNode) -> Result<(), GraphError> {
if self.node_count >= MAX_NODES {
return Err(GraphError::NodeCapacityExceeded(MAX_NODES));
}
if self.nodes.iter().any(|n| n.as_ref().is_some_and(|n| n.id == node.id)) {
return Err(GraphError::NodeAlreadyExists(node.id));
}
let slot = self
.nodes
.iter_mut()
.find(|s| s.is_none())
.ok_or(GraphError::NodeCapacityExceeded(MAX_NODES))?;
*slot = Some(node);
self.node_count += 1;
Ok(())
}
pub fn remove_node(&mut self, id: u64) -> Result<(), GraphError> {
let slot = self
.nodes
.iter_mut()
.find(|s| s.as_ref().is_some_and(|n| n.id == id))
.ok_or(GraphError::NodeNotFound(id))?;
*slot = None;
self.node_count -= 1;
Ok(())
}
pub fn add_mapping(&mut self, mapping: QueueMapping) -> Result<(), GraphError> {
if self.mapping_count >= MAX_MAPPINGS {
return Err(GraphError::MappingCapacityExceeded(MAX_MAPPINGS));
}
if self
.mappings
.iter()
.any(|m| m.as_ref().is_some_and(|m| m.queue_id == mapping.queue_id))
{
return Err(GraphError::InvalidTopology(
"duplicate queue_id in mappings",
));
}
let slot = self
.mappings
.iter_mut()
.find(|s| s.is_none())
.ok_or(GraphError::MappingCapacityExceeded(MAX_MAPPINGS))?;
*slot = Some(mapping);
self.mapping_count += 1;
Ok(())
}
pub fn remove_mapping(&mut self, queue_id: u32) -> Result<(), GraphError> {
let slot = self
.mappings
.iter_mut()
.find(|s| s.as_ref().is_some_and(|m| m.queue_id == queue_id))
.ok_or(GraphError::MappingNotFound(queue_id))?;
*slot = None;
self.mapping_count -= 1;
Ok(())
}
#[inline]
pub fn find_node(&self, id: u64) -> Option<&ResourceNode> {
self.nodes.iter().find_map(|n| n.as_ref().filter(|n| n.id == id))
}
#[inline]
pub fn find_mapping(&self, queue_id: u32) -> Option<&QueueMapping> {
self.mappings.iter().find_map(|m| m.as_ref().filter(|m| m.queue_id == queue_id))
}
pub fn nodes_in_domain(&self, domain: ExecutionDomain) -> Vec<&ResourceNode> {
self.nodes
.iter()
.filter_map(|n| n.as_ref())
.filter(|n| n.domain == domain)
.collect()
}
pub fn online_nodes(&self) -> Vec<&ResourceNode> {
self.nodes
.iter()
.filter_map(|n| n.as_ref())
.filter(|n| n.is_available())
.collect()
}
pub fn online_nodes_in_numa(&self, numa_node: u32) -> Vec<&ResourceNode> {
self.nodes
.iter()
.filter_map(|n| n.as_ref())
.filter(|n| n.is_available() && n.numa_node == numa_node)
.collect()
}
pub fn iter_nodes(&self) -> impl Iterator<Item = &ResourceNode> {
self.nodes.iter().filter_map(|n| n.as_ref())
}
pub fn iter_mappings(&self) -> impl Iterator<Item = &QueueMapping> {
self.mappings.iter().filter_map(|m| m.as_ref())
}
pub fn validate(&self) -> Result<(), GraphError> {
for node in self.iter_nodes() {
if node.cpu_cores == 0 {
return Err(GraphError::InvalidTopology("node has zero cpu_cores"));
}
if node.memory_mb == 0 {
return Err(GraphError::InvalidTopology("node has zero memory_mb"));
}
}
let mappings: Vec<&QueueMapping> = self.iter_mappings().collect();
for (i, m) in mappings.iter().enumerate() {
if mappings[..i].iter().any(|prev| prev.queue_id == m.queue_id) {
return Err(GraphError::InvalidTopology(
"duplicate queue_id in mappings",
));
}
}
Ok(())
}
pub fn plan_topology(
&self,
snapshots: &[DomainSnapshot],
strategy: QueueStrategy,
) -> Result<PlannedTopology, GraphError> {
let planner = TopologyPlanner;
planner.plan(snapshots, self, strategy)
}
pub fn find_optimal_layout(
&self,
domain: ExecutionDomain,
cpu_needed: u32,
mem_needed: u32,
queues_needed: u32,
) -> Result<Vec<u64>, GraphError> {
let available: Vec<&ResourceNode> = self
.nodes
.iter()
.filter_map(|n| n.as_ref())
.filter(|n| n.is_available() && n.domain == domain)
.collect();
if available.is_empty() {
return Err(GraphError::NoWorkersForDomain(domain));
}
let mut selected = Vec::new();
let mut total_cpu = 0u32;
let mut total_mem = 0u32;
let mut total_queues = 0u32;
let mut numa_groups: Vec<Vec<&ResourceNode>> = Vec::new();
for node in &available {
let found = numa_groups.iter_mut().find(|g| {
g.first().is_some_and(|n| n.numa_node == node.numa_node)
});
match found {
Some(group) => group.push(node),
None => numa_groups.push(vec![node]),
}
}
'outer: for group in &numa_groups {
for node in group {
if total_cpu >= cpu_needed
&& total_mem >= mem_needed
&& total_queues >= queues_needed
{
break 'outer;
}
if !selected.contains(&node.id) {
selected.push(node.id);
total_cpu = total_cpu.saturating_add(node.cpu_cores);
total_mem = total_mem.saturating_add(node.memory_mb);
total_queues = total_queues.saturating_add(node.nic_queues);
}
}
}
if total_cpu < cpu_needed || total_mem < mem_needed || total_queues < queues_needed {
for node in &available {
if total_cpu >= cpu_needed
&& total_mem >= mem_needed
&& total_queues >= queues_needed
{
break;
}
if !selected.contains(&node.id) {
selected.push(node.id);
total_cpu = total_cpu.saturating_add(node.cpu_cores);
total_mem = total_mem.saturating_add(node.memory_mb);
total_queues = total_queues.saturating_add(node.nic_queues);
}
}
}
if total_cpu < cpu_needed || total_mem < mem_needed || total_queues < queues_needed {
return Err(GraphError::InsufficientResources {
domain,
need_cpu: cpu_needed,
have_cpu: total_cpu,
need_mem: mem_needed,
have_mem: total_mem,
});
}
Ok(selected)
}
}
#[derive(Debug, Clone)]
pub struct PlannedTopology {
pub domain_assignments: Vec<DomainAssignment>,
pub queue_distributions: Vec<QueueDistribution>,
}
impl PlannedTopology {
#[inline]
pub fn domain_count(&self) -> usize {
self.domain_assignments.len()
}
#[inline]
pub fn total_workers(&self) -> usize {
self.domain_assignments.iter().map(|a| a.workers.len()).sum()
}
pub fn find_assignment(&self, domain: ExecutionDomain) -> Option<&DomainAssignment> {
self.domain_assignments.iter().find(|a| a.domain == domain)
}
pub fn validate_against(&self, snapshots: &[DomainSnapshot]) -> Result<(), GraphError> {
for snapshot in snapshots {
let assignment = self
.find_assignment(snapshot.domain)
.ok_or(GraphError::NoWorkersForDomain(snapshot.domain))?;
if assignment.workers.len() < snapshot.worker_count {
let need_cpu = snapshot
.required_cpu_cores
.saturating_mul(snapshot.worker_count as u32);
let have_cpu = snapshot
.required_cpu_cores
.saturating_mul(assignment.workers.len() as u32);
let need_mem = snapshot
.required_memory_mb
.saturating_mul(snapshot.worker_count as u32);
let have_mem = snapshot
.required_memory_mb
.saturating_mul(assignment.workers.len() as u32);
return Err(GraphError::InsufficientResources {
domain: snapshot.domain,
need_cpu,
have_cpu,
need_mem,
have_mem,
});
}
}
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct DomainAssignment {
pub domain: ExecutionDomain,
pub workers: Vec<WorkerAssignment>,
}
#[derive(Debug, Clone)]
pub struct QueueDistribution {
pub strategy: QueueStrategy,
pub mappings: Vec<QueueMapping>,
}
#[inline]
fn rss_hash_worker(queue_id: u32, worker_count: u32) -> u32 {
let mut hash: u32 = 0x811c_9dc5;
for byte in queue_id.to_le_bytes() {
hash = (hash ^ u32::from(byte)).wrapping_mul(0x0100_0193);
}
hash % worker_count
}
#[derive(Debug)]
pub struct TopologyPlanner;
impl TopologyPlanner {
pub fn plan<const MAX_NODES: usize, const MAX_MAPPINGS: usize>(
&self,
snapshots: &[DomainSnapshot],
graph: &RuntimeGraph<MAX_NODES, MAX_MAPPINGS>,
strategy: QueueStrategy,
) -> Result<PlannedTopology, GraphError> {
let mut domain_assignments = Vec::with_capacity(snapshots.len());
let mut queue_distributions = Vec::with_capacity(snapshots.len());
for snapshot in snapshots {
let assignment = self.assign_domain(snapshot, graph)?;
let distribution = self.distribute_queues(&assignment, graph, strategy);
domain_assignments.push(assignment);
queue_distributions.push(distribution);
}
Ok(PlannedTopology {
domain_assignments,
queue_distributions,
})
}
fn assign_domain<const MAX_NODES: usize, const MAX_MAPPINGS: usize>(
&self,
snapshot: &DomainSnapshot,
graph: &RuntimeGraph<MAX_NODES, MAX_MAPPINGS>,
) -> Result<DomainAssignment, GraphError> {
let domain_nodes = graph.nodes_in_domain(snapshot.domain);
let available: Vec<&ResourceNode> = domain_nodes
.into_iter()
.filter(|n| n.is_available())
.collect();
if available.is_empty() {
return Err(GraphError::NoWorkersForDomain(snapshot.domain));
}
let (per_cpu, per_mem, per_queues) = snapshot.domain.default_requirements();
let numa_groups = self.group_by_numa(&available);
let mut workers = Vec::with_capacity(snapshot.worker_count);
let mut worker_id_counter: u32 = 0;
let mut used_cpu: std::collections::HashMap<u64, u32> = std::collections::HashMap::new();
let mut used_mem: std::collections::HashMap<u64, u32> = std::collections::HashMap::new();
let mut queues_global: u32 = available
.iter()
.map(|n| n.nic_queues)
.fold(0u32, |a, b| a.saturating_add(b));
for group in &numa_groups {
if workers.len() >= snapshot.worker_count {
break;
}
let mut queues_available: u32 = group
.iter()
.map(|n| n.nic_queues)
.fold(0u32, |a, b| a.saturating_add(b));
loop {
if workers.len() >= snapshot.worker_count {
break;
}
if queues_available < per_queues {
break;
}
let chosen = group.iter().find(|n| {
let used_c = used_cpu.get(&n.id).copied().unwrap_or(0);
let used_m = used_mem.get(&n.id).copied().unwrap_or(0);
n.cpu_cores.saturating_sub(used_c) >= per_cpu
&& n.memory_mb.saturating_sub(used_m) >= per_mem
});
let Some(&node) = chosen else {
break;
};
*used_cpu.entry(node.id).or_insert(0) += per_cpu;
*used_mem.entry(node.id).or_insert(0) += per_mem;
queues_available = queues_available.saturating_sub(per_queues);
workers.push(WorkerAssignment {
worker_id: worker_id_counter,
node_id: node.id,
domain: snapshot.domain,
numa_node: node.numa_node,
queue_ids: Vec::new(),
});
worker_id_counter = worker_id_counter.saturating_add(1);
}
}
while workers.len() < snapshot.worker_count {
if queues_global < per_queues {
break;
}
let chosen = available.iter().find(|n| {
let used_c = used_cpu.get(&n.id).copied().unwrap_or(0);
let used_m = used_mem.get(&n.id).copied().unwrap_or(0);
n.cpu_cores.saturating_sub(used_c) >= per_cpu
&& n.memory_mb.saturating_sub(used_m) >= per_mem
});
let Some(&node) = chosen else {
break;
};
*used_cpu.entry(node.id).or_insert(0) += per_cpu;
*used_mem.entry(node.id).or_insert(0) += per_mem;
queues_global = queues_global.saturating_sub(per_queues);
workers.push(WorkerAssignment {
worker_id: worker_id_counter,
node_id: node.id,
domain: snapshot.domain,
numa_node: node.numa_node,
queue_ids: Vec::new(),
});
worker_id_counter = worker_id_counter.saturating_add(1);
}
if workers.len() < snapshot.worker_count {
let need_cpu: u32 = (snapshot.worker_count as u32).saturating_mul(per_cpu);
let have_cpu: u32 = (workers.len() as u32).saturating_mul(per_cpu);
let need_mem: u32 = (snapshot.worker_count as u32).saturating_mul(per_mem);
let have_mem: u32 = (workers.len() as u32).saturating_mul(per_mem);
return Err(GraphError::InsufficientResources {
domain: snapshot.domain,
need_cpu,
have_cpu,
need_mem,
have_mem,
});
}
Ok(DomainAssignment {
domain: snapshot.domain,
workers,
})
}
fn group_by_numa<'a>(
&self,
nodes: &[&'a ResourceNode],
) -> Vec<Vec<&'a ResourceNode>> {
let mut groups: Vec<Vec<&ResourceNode>> = Vec::new();
for node in nodes {
let found = groups.iter_mut().find(|g| {
g.first().is_some_and(|n| n.numa_node == node.numa_node)
});
match found {
Some(group) => group.push(node),
None => groups.push(vec![node]),
}
}
groups
}
fn distribute_queues<const MAX_NODES: usize, const MAX_MAPPINGS: usize>(
&self,
assignment: &DomainAssignment,
graph: &RuntimeGraph<MAX_NODES, MAX_MAPPINGS>,
strategy: QueueStrategy,
) -> QueueDistribution {
if assignment.workers.is_empty() {
return QueueDistribution {
strategy,
mappings: Vec::new(),
};
}
let mut mappings = Vec::new();
let mut seen_nodes = std::collections::HashSet::new();
let total_queues: u32 = assignment
.workers
.iter()
.filter_map(|w| graph.find_node(w.node_id))
.filter(|n| seen_nodes.insert(n.id))
.map(|n| n.nic_queues)
.fold(0u32, |a, b| a.saturating_add(b));
let worker_count = assignment.workers.len() as u32;
match strategy {
QueueStrategy::RssHash => {
for queue_id in 0..total_queues {
let worker_id = rss_hash_worker(queue_id, worker_count);
mappings.push(QueueMapping {
queue_id,
worker_id,
affinity: AffinityRule::None,
});
}
}
QueueStrategy::RoundRobin => {
for queue_id in 0..total_queues {
let worker_id = queue_id % worker_count;
mappings.push(QueueMapping {
queue_id,
worker_id,
affinity: AffinityRule::CpuLocal,
});
}
}
QueueStrategy::NumaLocal => {
let mut next_queue_id: u32 = 0;
for (idx, worker) in assignment.workers.iter().enumerate() {
if let Some(node) = graph.find_node(worker.node_id) {
let queues = node.nic_queues;
for _q in 0..queues {
mappings.push(QueueMapping {
queue_id: next_queue_id,
worker_id: idx as u32,
affinity: AffinityRule::NumaLocal,
});
next_queue_id = next_queue_id.saturating_add(1);
}
}
}
}
}
QueueDistribution {
strategy,
mappings,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn make_node(id: u64, domain: ExecutionDomain, cpu: u32, mem: u32, queues: u32) -> ResourceNode {
ResourceNode::new(id, domain, cpu, mem, queues)
}
fn make_numa_node(
id: u64,
domain: ExecutionDomain,
cpu: u32,
mem: u32,
queues: u32,
numa: u32,
) -> ResourceNode {
ResourceNode::new(id, domain, cpu, mem, queues).with_numa(numa)
}
fn small_graph() -> RuntimeGraph<8, 32> {
RuntimeGraph::new()
}
#[test]
fn test_execution_domain_properties() {
let dp = ExecutionDomain::DataPlane;
assert!(dp.is_data_plane());
assert!(dp.requires_high_cpu());
assert!(!dp.requires_large_memory());
let storage = ExecutionDomain::Storage;
assert!(storage.requires_large_memory());
assert!(!storage.requires_high_cpu());
let all = ExecutionDomain::all();
assert_eq!(all.len(), 6);
}
#[test]
fn test_domain_default_requirements() {
let (cpu, mem, queues) = ExecutionDomain::DataPlane.default_requirements();
assert_eq!(cpu, 4);
assert_eq!(mem, 4096);
assert_eq!(queues, 8);
let (_cpu, mem, _queues) = ExecutionDomain::Storage.default_requirements();
assert_eq!(mem, 16384);
}
#[test]
fn test_add_and_remove_node() {
let mut graph = small_graph();
assert_eq!(graph.node_count(), 0);
let node = make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4);
graph.add_node(node).unwrap();
assert_eq!(graph.node_count(), 1);
assert!(graph.find_node(1).is_some());
assert!(graph.find_node(99).is_none());
graph.remove_node(1).unwrap();
assert_eq!(graph.node_count(), 0);
assert!(graph.find_node(1).is_none());
}
#[test]
fn test_duplicate_node_id_rejected() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
let err = graph
.add_node(make_node(1, ExecutionDomain::ControlPlane, 4, 4096, 2))
.unwrap_err();
assert!(matches!(err, GraphError::NodeAlreadyExists(1)));
}
#[test]
fn test_node_capacity_exceeded() {
let mut graph: RuntimeGraph<2, 8> = RuntimeGraph::new();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
graph
.add_node(make_node(2, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
let err = graph
.add_node(make_node(3, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap_err();
assert!(matches!(err, GraphError::NodeCapacityExceeded(2)));
}
#[test]
fn test_queue_mapping_management() {
let mut graph = small_graph();
let m1 = QueueMapping::new(0, 0).with_affinity(AffinityRule::NumaLocal);
graph.add_mapping(m1).unwrap();
assert_eq!(graph.mapping_count(), 1);
assert!(graph.find_mapping(0).is_some());
assert!(graph.find_mapping(99).is_none());
graph.remove_mapping(0).unwrap();
assert_eq!(graph.mapping_count(), 0);
}
#[test]
fn test_nodes_in_domain() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
graph
.add_node(make_node(2, ExecutionDomain::DataPlane, 4, 4096, 2))
.unwrap();
graph
.add_node(make_node(3, ExecutionDomain::ControlPlane, 2, 2048, 1))
.unwrap();
let dp_nodes = graph.nodes_in_domain(ExecutionDomain::DataPlane);
assert_eq!(dp_nodes.len(), 2);
let cp_nodes = graph.nodes_in_domain(ExecutionDomain::ControlPlane);
assert_eq!(cp_nodes.len(), 1);
let mgmt_nodes = graph.nodes_in_domain(ExecutionDomain::Management);
assert_eq!(mgmt_nodes.len(), 0);
}
#[test]
fn test_online_nodes_filter() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
graph
.add_node(
ResourceNode::new(2, ExecutionDomain::DataPlane, 4, 4096, 2)
.with_status(NodeStatus::Offline),
)
.unwrap();
graph
.add_node(make_node(3, ExecutionDomain::ControlPlane, 2, 2048, 1))
.unwrap();
let online = graph.online_nodes();
assert_eq!(online.len(), 2); }
#[test]
fn test_find_optimal_layout_basic() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
graph
.add_node(make_node(2, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
graph
.add_node(make_node(3, ExecutionDomain::ControlPlane, 2, 2048, 1))
.unwrap();
let result = graph
.find_optimal_layout(ExecutionDomain::DataPlane, 4, 4096, 2)
.unwrap();
assert!(!result.is_empty());
let node = graph.find_node(result[0]).unwrap();
assert_eq!(node.domain, ExecutionDomain::DataPlane);
}
#[test]
fn test_insufficient_resources_error() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 2, 1024, 1))
.unwrap();
let err = graph
.find_optimal_layout(ExecutionDomain::DataPlane, 16, 65536, 8)
.unwrap_err();
assert!(matches!(err, GraphError::InsufficientResources { .. }));
}
#[test]
fn test_numa_aware_grouping() {
let mut graph = small_graph();
graph
.add_node(make_numa_node(
1,
ExecutionDomain::DataPlane,
8,
8192,
4,
0,
))
.unwrap();
graph
.add_node(make_numa_node(
2,
ExecutionDomain::DataPlane,
8,
8192,
4,
0,
))
.unwrap();
graph
.add_node(make_numa_node(
3,
ExecutionDomain::DataPlane,
8,
8192,
4,
1,
))
.unwrap();
let online_numa0 = graph.online_nodes_in_numa(0);
assert_eq!(online_numa0.len(), 2);
let online_numa1 = graph.online_nodes_in_numa(1);
assert_eq!(online_numa1.len(), 1);
}
#[test]
fn test_planner_rss_hash_strategy() {
let mut graph = small_graph();
for i in 0..4 {
graph
.add_node(make_numa_node(
i,
ExecutionDomain::DataPlane,
8,
8192,
4,
(i % 2) as u32,
))
.unwrap();
}
let snapshot = DomainSnapshot::from_domain(ExecutionDomain::DataPlane, 2);
let topology = graph
.plan_topology(&[snapshot], QueueStrategy::RssHash)
.unwrap();
assert_eq!(topology.domain_count(), 1);
assert!(topology.total_workers() >= 2);
let assignment = topology
.find_assignment(ExecutionDomain::DataPlane)
.unwrap();
assert_eq!(assignment.workers.len(), 2);
let dist = &topology.queue_distributions[0];
assert!(dist.mappings.len() >= 4); }
#[test]
fn test_planner_round_robin_strategy() {
let mut graph = small_graph();
for i in 0..3 {
graph
.add_node(make_node(i, ExecutionDomain::ControlPlane, 4, 8192, 8))
.unwrap();
}
let snapshot = DomainSnapshot::from_domain(ExecutionDomain::ControlPlane, 3);
let topology = graph
.plan_topology(&[snapshot], QueueStrategy::RoundRobin)
.unwrap();
let dist = &topology.queue_distributions[0];
let w0_count = dist.mappings.iter().filter(|m| m.worker_id == 0).count();
let w1_count = dist.mappings.iter().filter(|m| m.worker_id == 1).count();
assert!(w0_count >= 1);
assert!(w1_count >= 1);
}
#[test]
fn test_planner_numa_local_strategy() {
let mut graph = small_graph();
graph
.add_node(make_numa_node(
0,
ExecutionDomain::DataPlane,
8,
16384,
12,
0,
))
.unwrap();
graph
.add_node(make_numa_node(
1,
ExecutionDomain::DataPlane,
8,
16384,
12,
0,
))
.unwrap();
graph
.add_node(make_numa_node(
2,
ExecutionDomain::DataPlane,
8,
16384,
12,
1,
))
.unwrap();
let snapshot = DomainSnapshot::from_domain(ExecutionDomain::DataPlane, 3);
let topology = graph
.plan_topology(&[snapshot], QueueStrategy::NumaLocal)
.unwrap();
let assignment = topology
.find_assignment(ExecutionDomain::DataPlane)
.unwrap();
assert_eq!(assignment.workers.len(), 3);
let dist = &topology.queue_distributions[0];
for mapping in &dist.mappings {
assert!(matches!(mapping.affinity, AffinityRule::NumaLocal));
}
}
#[test]
fn test_multi_domain_planning() {
let mut graph = small_graph();
graph
.add_node(make_numa_node(
0,
ExecutionDomain::DataPlane,
16,
16384,
8,
0,
))
.unwrap();
graph
.add_node(make_numa_node(
1,
ExecutionDomain::DataPlane,
16,
16384,
8,
1,
))
.unwrap();
graph
.add_node(make_numa_node(
2,
ExecutionDomain::ControlPlane,
4,
8192,
8,
0,
))
.unwrap();
graph
.add_node(make_numa_node(
3,
ExecutionDomain::Observability,
4,
16384,
8,
1,
))
.unwrap();
let snapshots = vec![
DomainSnapshot::from_domain(ExecutionDomain::DataPlane, 2),
DomainSnapshot::from_domain(ExecutionDomain::ControlPlane, 1),
DomainSnapshot::from_domain(ExecutionDomain::Observability, 1),
];
let topology = graph
.plan_topology(&snapshots, QueueStrategy::RssHash)
.unwrap();
assert_eq!(topology.domain_count(), 3);
assert_eq!(topology.total_workers(), 4);
let dp = topology
.find_assignment(ExecutionDomain::DataPlane)
.unwrap();
assert_eq!(dp.workers.len(), 2);
let cp = topology
.find_assignment(ExecutionDomain::ControlPlane)
.unwrap();
assert_eq!(cp.workers.len(), 1);
}
#[test]
fn test_validate_topology() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
graph
.add_node(make_node(2, ExecutionDomain::ControlPlane, 2, 2048, 1))
.unwrap();
assert!(graph.validate().is_ok());
graph
.add_node(ResourceNode::new(3, ExecutionDomain::Management, 0, 1024, 1))
.unwrap();
let err = graph.validate().unwrap_err();
assert!(matches!(err, GraphError::InvalidTopology("node has zero cpu_cores")));
}
#[test]
fn test_domain_isolation() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
graph
.add_node(make_node(2, ExecutionDomain::ControlPlane, 4, 4096, 2))
.unwrap();
let dp_nodes = graph.nodes_in_domain(ExecutionDomain::DataPlane);
assert_eq!(dp_nodes.len(), 1);
assert_eq!(dp_nodes[0].id, 1);
let cp_nodes = graph.nodes_in_domain(ExecutionDomain::ControlPlane);
assert_eq!(cp_nodes.len(), 1);
assert_eq!(cp_nodes[0].id, 2);
let result = graph
.find_optimal_layout(ExecutionDomain::DataPlane, 4, 2048, 2)
.unwrap();
for node_id in &result {
let node = graph.find_node(*node_id).unwrap();
assert_eq!(node.domain, ExecutionDomain::DataPlane);
}
}
#[test]
fn test_resource_node_meets_requirements() {
let node = ResourceNode::new(1, ExecutionDomain::DataPlane, 8, 8192, 4);
assert!(node.meets_requirements(4, 4096, 2));
assert!(node.meets_requirements(8, 8192, 4));
assert!(!node.meets_requirements(16, 8192, 4));
assert!(!node.meets_requirements(8, 16384, 4));
assert!(!node.meets_requirements(8, 8192, 8));
}
#[test]
fn test_planned_topology_validate_against() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 16384, 12))
.unwrap();
let snapshot = DomainSnapshot::from_domain(ExecutionDomain::DataPlane, 1);
let topology = graph
.plan_topology(&[snapshot], QueueStrategy::RssHash)
.unwrap();
assert!(topology.validate_against(&[snapshot]).is_ok());
let insufficient =
DomainSnapshot::from_domain(ExecutionDomain::DataPlane, 10);
let err = topology.validate_against(&[insufficient]).unwrap_err();
assert!(matches!(err, GraphError::InsufficientResources { .. }));
}
#[test]
fn test_node_status_schedulable() {
assert!(NodeStatus::Online.is_schedulable());
assert!(!NodeStatus::Offline.is_schedulable());
assert!(!NodeStatus::Degraded.is_schedulable());
assert!(!NodeStatus::Maintenance.is_schedulable());
let node = ResourceNode::new(1, ExecutionDomain::DataPlane, 8, 8192, 4);
assert!(node.is_available());
let offline = node.with_status(NodeStatus::Offline);
assert!(!offline.is_available());
}
#[test]
fn test_empty_graph_operations() {
let graph: RuntimeGraph<4, 16> = RuntimeGraph::new();
assert!(graph.is_empty());
assert_eq!(graph.node_count(), 0);
assert_eq!(graph.mapping_count(), 0);
assert_eq!(graph.max_nodes(), 4);
assert_eq!(graph.max_mappings(), 16);
let nodes = graph.online_nodes();
assert!(nodes.is_empty());
let err = graph
.find_optimal_layout(ExecutionDomain::DataPlane, 1, 1, 1)
.unwrap_err();
assert!(matches!(err, GraphError::NoWorkersForDomain(_)));
}
#[test]
fn test_node_not_found_error() {
let mut graph = small_graph();
let err = graph.remove_node(999).unwrap_err();
assert!(matches!(err, GraphError::NodeNotFound(999)));
}
#[test]
fn test_domain_snapshot_customization() {
let snap = DomainSnapshot::from_domain(ExecutionDomain::DataPlane, 4)
.with_cpu(16)
.with_memory(32768)
.with_workers(8);
assert_eq!(snap.domain, ExecutionDomain::DataPlane);
assert_eq!(snap.required_cpu_cores, 16);
assert_eq!(snap.required_memory_mb, 32768);
assert_eq!(snap.worker_count, 8);
}
#[test]
fn test_const_constructors() {
const NODE: ResourceNode =
ResourceNode::new(1, ExecutionDomain::DataPlane, 8, 8192, 4);
assert_eq!(NODE.id, 1);
assert_eq!(NODE.cpu_cores, 8);
const MAPPING: QueueMapping = QueueMapping::new(0, 0);
assert_eq!(MAPPING.queue_id, 0);
assert_eq!(MAPPING.worker_id, 0);
const GRAPH: RuntimeGraph<16, 64> = RuntimeGraph::new();
assert_eq!(GRAPH.node_count(), 0);
}
#[test]
fn test_execution_domain_display() {
let domains = ExecutionDomain::all();
for d in domains.iter() {
let display = format!("{}", d);
assert!(!display.is_empty());
}
}
#[test]
fn test_node_status_debug() {
let statuses = [
NodeStatus::Online,
NodeStatus::Offline,
NodeStatus::Degraded,
NodeStatus::Maintenance,
];
for s in statuses.iter() {
let debug_str = format!("{:?}", s);
assert!(!debug_str.is_empty());
}
}
#[test]
fn test_graph_error_display() {
let errors = [
GraphError::NodeCapacityExceeded(10),
GraphError::MappingCapacityExceeded(20),
GraphError::NodeNotFound(42),
GraphError::NodeAlreadyExists(7),
GraphError::InvalidTopology("test error"),
GraphError::NoWorkersForDomain(ExecutionDomain::DataPlane),
];
for e in errors.iter() {
let display = format!("{}", e);
assert!(!display.is_empty());
}
}
#[test]
fn test_queue_strategy_debug() {
let strategies = [
QueueStrategy::RssHash,
QueueStrategy::RoundRobin,
QueueStrategy::NumaLocal,
];
for s in strategies.iter() {
let debug_str = format!("{:?}", s);
assert!(!debug_str.is_empty());
}
}
#[test]
fn test_affinity_rule_variants() {
let rules = [
AffinityRule::None,
AffinityRule::NumaLocal,
AffinityRule::CpuLocal,
AffinityRule::Exact,
];
assert_eq!(rules.len(), 4);
for r in rules.iter() {
let debug_str = format!("{:?}", r);
assert!(!debug_str.is_empty());
}
}
#[test]
fn test_worker_assignment_debug() {
let assignment = WorkerAssignment {
worker_id: 0,
node_id: 1,
domain: ExecutionDomain::DataPlane,
numa_node: 0,
queue_ids: vec![0, 1, 2],
};
let debug_str = format!("{:?}", assignment);
assert!(!debug_str.is_empty());
assert!(debug_str.contains("DataPlane"));
}
#[test]
fn test_planned_topology_domain_count() {
let topology = PlannedTopology {
domain_assignments: vec![],
queue_distributions: vec![],
};
assert_eq!(topology.domain_count(), 0);
assert_eq!(topology.total_workers(), 0);
}
#[test]
fn test_resource_node_with_numa() {
let node = ResourceNode::new(1, ExecutionDomain::DataPlane, 8, 8192, 4)
.with_numa(2);
assert_eq!(node.numa_node, 2);
}
#[test]
fn test_remove_nonexistent_mapping() {
let mut graph: RuntimeGraph<4, 16> = RuntimeGraph::new();
let result = graph.remove_mapping(999);
assert!(result.is_err());
}
#[test]
fn test_mapping_capacity_boundary() {
let mut graph: RuntimeGraph<4, 2> = RuntimeGraph::new();
graph.add_mapping(QueueMapping::new(0, 0)).unwrap();
graph.add_mapping(QueueMapping::new(1, 1)).unwrap();
let result = graph.add_mapping(QueueMapping::new(2, 2));
assert!(result.is_err());
assert!(matches!(result.unwrap_err(), GraphError::MappingCapacityExceeded(2)));
}
#[test]
fn test_rss_hash_differs_from_round_robin() {
let worker_count = 4u32;
let differs = (0..16u32).any(|q| rss_hash_worker(q, worker_count) != q % worker_count);
assert!(differs, "RSS 哈希映射与轮询取模完全相同(语义造假)");
for q in 0..16u32 {
assert_eq!(rss_hash_worker(q, worker_count), rss_hash_worker(q, worker_count));
assert!(rss_hash_worker(q, worker_count) < worker_count);
}
}
#[test]
fn test_numa_local_queue_ids_globally_unique() {
let mut graph = small_graph();
graph
.add_node(make_numa_node(0, ExecutionDomain::DataPlane, 8, 16384, 8, 0))
.unwrap();
graph
.add_node(make_numa_node(1, ExecutionDomain::DataPlane, 8, 16384, 8, 0))
.unwrap();
let snapshot = DomainSnapshot::from_domain(ExecutionDomain::DataPlane, 2);
let topology = graph
.plan_topology(&[snapshot], QueueStrategy::NumaLocal)
.unwrap();
let dist = &topology.queue_distributions[0];
assert_eq!(dist.mappings.len(), 16);
let mut ids: Vec<u32> = dist.mappings.iter().map(|m| m.queue_id).collect();
ids.sort_unstable();
ids.dedup();
assert_eq!(
ids.len(),
dist.mappings.len(),
"NumaLocal 映射的 queue_id 必须全局唯一"
);
}
#[test]
fn test_add_mapping_rejects_duplicate_queue_id() {
let mut graph = small_graph();
graph
.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4))
.unwrap();
graph.add_mapping(QueueMapping::new(0, 0)).unwrap();
let err = graph.add_mapping(QueueMapping::new(0, 1)).unwrap_err();
assert!(matches!(
err,
GraphError::InvalidTopology("duplicate queue_id in mappings")
));
}
#[test]
fn test_iterators() {
let mut graph: RuntimeGraph<8, 16> = RuntimeGraph::new();
graph.add_node(make_node(1, ExecutionDomain::DataPlane, 8, 8192, 4)).unwrap();
graph.add_node(make_node(2, ExecutionDomain::ControlPlane, 4, 4096, 2)).unwrap();
graph.add_mapping(QueueMapping::new(0, 0)).unwrap();
let nodes: Vec<_> = graph.iter_nodes().collect();
assert_eq!(nodes.len(), 2);
let mappings: Vec<_> = graph.iter_mappings().collect();
assert_eq!(mappings.len(), 1);
}
}