use crate::app::node_info::{ClusterView, NodeInfo};
use crate::app::spec::ActorSpec;
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct PlacementDecision {
pub node_id: String,
pub reason: String,
}
pub trait PlacementStrategy: Send + Sync {
fn fitness(&self, node: &NodeInfo, spec: &ActorSpec, view: &ClusterView) -> f64;
fn name(&self) -> &str {
std::any::type_name::<Self>()
}
}
pub fn select_node(
strategy: &dyn PlacementStrategy,
spec: &ActorSpec,
view: &ClusterView,
) -> Option<PlacementDecision> {
let mut best: Option<(String, f64)> = None;
for node in view.alive_nodes() {
if !spec.placement.is_satisfied_by(
&node.node_id(),
&node.class,
&node.metadata,
&node.running_actors,
) {
continue;
}
let score = strategy.fitness(node, spec, view);
if score <= 0.0 {
continue;
}
match &best {
Some((_, best_score)) if score <= *best_score => {}
_ => {
best = Some((node.node_id(), score));
}
}
}
best.map(|(node_id, score)| PlacementDecision {
reason: format!(
"strategy={}, score={score:.2}, node={node_id}",
strategy.name()
),
node_id,
})
}
pub struct LeastLoaded;
impl PlacementStrategy for LeastLoaded {
fn fitness(&self, node: &NodeInfo, _spec: &ActorSpec, _view: &ClusterView) -> f64 {
1.0 / (node.actor_count() as f64 + 1.0)
}
fn name(&self) -> &str {
"LeastLoaded"
}
}
pub struct RandomPlacement {
rng: std::sync::Mutex<rand::rngs::StdRng>,
}
impl RandomPlacement {
pub fn new(seed: u64) -> Self {
use rand::SeedableRng;
Self {
rng: std::sync::Mutex::new(rand::rngs::StdRng::seed_from_u64(seed)),
}
}
}
impl PlacementStrategy for RandomPlacement {
fn fitness(&self, _node: &NodeInfo, _spec: &ActorSpec, _view: &ClusterView) -> f64 {
use rand::Rng;
self.rng.lock().unwrap().random::<f64>()
}
fn name(&self) -> &str {
"Random"
}
}
pub struct Pinned {
pub preferred_node_id: String,
}
impl PlacementStrategy for Pinned {
fn fitness(&self, node: &NodeInfo, _spec: &ActorSpec, _view: &ClusterView) -> f64 {
if node.node_id() == self.preferred_node_id {
1000.0 } else {
1.0 }
}
fn name(&self) -> &str {
"Pinned"
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cluster::config::{NodeClass, NodeIdentity};
use std::collections::HashMap;
fn make_node(name: &str, incarnation: u64, actors: Vec<String>) -> NodeInfo {
NodeInfo {
identity: NodeIdentity::for_test(name, incarnation),
class: NodeClass::Worker,
metadata: HashMap::new(),
running_actors: actors,
is_alive: true,
}
}
fn make_view(nodes: Vec<NodeInfo>) -> ClusterView {
let mut view = ClusterView::new();
for node in nodes {
view.upsert_node(node);
}
view
}
#[test]
fn test_least_loaded_prefers_empty_node() {
let view = make_view(vec![
make_node("alpha", 1, vec!["w/0".into(), "w/1".into()]),
make_node("beta", 2, vec![]),
]);
let spec = ActorSpec::new("w/2", "app::Worker");
let decision = select_node(&LeastLoaded, &spec, &view).unwrap();
assert!(
decision.node_id.contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("beta").to_string()
),
"expected beta, got {}",
decision.node_id
);
}
#[test]
fn test_pinned_prefers_target_node() {
let n1 = make_node("alpha", 1, vec![]);
let n2 = make_node("beta", 2, vec![]);
let beta_id = n2.node_id();
let view = make_view(vec![n1, n2]);
let strategy = Pinned {
preferred_node_id: beta_id.clone(),
};
let spec = ActorSpec::new("w/0", "app::Worker");
let decision = select_node(&strategy, &spec, &view).unwrap();
assert_eq!(decision.node_id, beta_id);
}
#[test]
fn test_no_eligible_nodes_returns_none() {
let view = make_view(vec![make_node("alpha", 1, vec![])]);
let spec = ActorSpec::new("coord/0", "app::Coordinator").with_constraints(
crate::app::spec::PlacementConstraints {
required_classes: vec![NodeClass::Coordinator],
..Default::default()
},
);
let decision = select_node(&LeastLoaded, &spec, &view);
assert!(decision.is_none());
}
#[test]
fn test_dead_nodes_excluded() {
let mut view = make_view(vec![
make_node("alpha", 1, vec![]),
make_node("beta", 2, vec!["w/0".into()]),
]);
let alpha_id = view
.nodes
.keys()
.find(|k| {
k.contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("alpha").to_string(),
)
})
.unwrap()
.clone();
view.mark_failed(&alpha_id);
let spec = ActorSpec::new("w/1", "app::Worker");
let decision = select_node(&LeastLoaded, &spec, &view).unwrap();
assert!(
decision.node_id.contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("beta").to_string()
),
"expected beta, got {}",
decision.node_id
);
}
#[test]
fn test_colocate_with_pins_to_anchor_node() {
let view = make_view(vec![
make_node("alpha", 1, vec![]),
make_node("beta", 2, vec!["writer/c1".into()]),
]);
let spec = ActorSpec::new("reader/c1/0", "app::Reader").with_constraints(
crate::app::spec::PlacementConstraints {
colocate_with: Some("writer/c1".into()),
..Default::default()
},
);
let decision = select_node(&LeastLoaded, &spec, &view).unwrap();
assert!(
decision.node_id.contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("beta").to_string()
),
"colocate_with should pin to the anchor's node (beta), got {}",
decision.node_id
);
}
#[test]
fn test_colocate_with_no_anchor_returns_none() {
let view = make_view(vec![
make_node("alpha", 1, vec![]),
make_node("beta", 2, vec![]),
]);
let spec = ActorSpec::new("reader/c1/0", "app::Reader").with_constraints(
crate::app::spec::PlacementConstraints {
colocate_with: Some("writer/c1".into()),
..Default::default()
},
);
assert!(select_node(&LeastLoaded, &spec, &view).is_none());
}
#[test]
fn test_required_node_id_pins_exactly_or_none() {
let alpha = make_node("alpha", 1, vec![]); let beta = make_node("beta", 2, vec!["x".into()]);
let beta_id = beta.node_id();
let view = make_view(vec![alpha, beta]);
let spec = ActorSpec::new("w/0", "app::Worker").with_constraints(
crate::app::spec::PlacementConstraints {
required_node_id: Some(beta_id.clone()),
..Default::default()
},
);
assert_eq!(
select_node(&LeastLoaded, &spec, &view).unwrap().node_id,
beta_id
);
let spec_missing = ActorSpec::new("w/1", "app::Worker").with_constraints(
crate::app::spec::PlacementConstraints {
required_node_id: Some("ghost@127.0.0.1:9999#0".into()),
..Default::default()
},
);
assert!(select_node(&LeastLoaded, &spec_missing, &view).is_none());
}
#[test]
fn test_default_constraints_unchanged_backward_compat() {
let view = make_view(vec![make_node("alpha", 1, vec![])]);
let spec = ActorSpec::new("w/0", "app::Worker");
assert!(select_node(&LeastLoaded, &spec, &view).is_some());
}
}