use std::collections::HashMap;
use std::time::Duration;
use crate::cluster::config::NodeClass;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub enum CrashStrategy {
#[default]
Redistribute,
WaitForReturn(Duration),
Abandon,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct PlacementConstraints {
pub required_classes: Vec<NodeClass>,
pub required_metadata: HashMap<String, String>,
pub anti_affinity_labels: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub colocate_with: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub required_node_id: Option<String>,
}
impl PlacementConstraints {
pub fn is_satisfied_by(
&self,
node_id: &str,
node_class: &NodeClass,
node_metadata: &HashMap<String, String>,
running_actors: &[String],
) -> bool {
if !self.required_classes.is_empty() && !self.required_classes.contains(node_class) {
return false;
}
for (key, value) in &self.required_metadata {
match node_metadata.get(key) {
Some(v) if v == value => {}
_ => return false,
}
}
for label in &self.anti_affinity_labels {
if running_actors.contains(label) {
return false;
}
}
if let Some(anchor) = &self.colocate_with
&& !running_actors.iter().any(|l| l == anchor)
{
return false;
}
if let Some(required) = &self.required_node_id
&& node_id != required
{
return false;
}
true
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ActorSpec {
pub label: String,
pub actor_type_name: String,
pub initial_state: Vec<u8>,
pub crash_strategy: CrashStrategy,
pub placement: PlacementConstraints,
}
impl ActorSpec {
pub fn new(label: impl Into<String>, actor_type_name: impl Into<String>) -> Self {
Self {
label: label.into(),
actor_type_name: actor_type_name.into(),
initial_state: Vec::new(),
crash_strategy: CrashStrategy::default(),
placement: PlacementConstraints::default(),
}
}
pub fn with_state(mut self, state: Vec<u8>) -> Self {
self.initial_state = state;
self
}
pub fn with_crash_strategy(mut self, strategy: CrashStrategy) -> Self {
self.crash_strategy = strategy;
self
}
pub fn with_constraints(mut self, constraints: PlacementConstraints) -> Self {
self.placement = constraints;
self
}
}
#[cfg(test)]
mod tests {
use super::*;
const ANY_NODE: &str = "node-a@127.0.0.1:7100#1";
#[test]
fn test_placement_constraints_empty_allows_all() {
let constraints = PlacementConstraints::default();
assert!(constraints.is_satisfied_by(ANY_NODE, &NodeClass::Worker, &HashMap::new(), &[]));
}
#[test]
fn test_placement_constraints_required_class() {
let constraints = PlacementConstraints {
required_classes: vec![NodeClass::Worker],
..Default::default()
};
assert!(constraints.is_satisfied_by(ANY_NODE, &NodeClass::Worker, &HashMap::new(), &[]));
assert!(!constraints.is_satisfied_by(ANY_NODE, &NodeClass::Edge, &HashMap::new(), &[]));
}
#[test]
fn test_placement_constraints_required_metadata() {
let constraints = PlacementConstraints {
required_metadata: [("region".into(), "us-west".into())].into(),
..Default::default()
};
let good_meta: HashMap<String, String> = [("region".into(), "us-west".into())].into();
let bad_meta: HashMap<String, String> = [("region".into(), "eu-west".into())].into();
assert!(constraints.is_satisfied_by(ANY_NODE, &NodeClass::Worker, &good_meta, &[]));
assert!(!constraints.is_satisfied_by(ANY_NODE, &NodeClass::Worker, &bad_meta, &[]));
assert!(!constraints.is_satisfied_by(ANY_NODE, &NodeClass::Worker, &HashMap::new(), &[]));
}
#[test]
fn test_placement_constraints_anti_affinity() {
let constraints = PlacementConstraints {
anti_affinity_labels: vec!["replica/0".into()],
..Default::default()
};
assert!(constraints.is_satisfied_by(
ANY_NODE,
&NodeClass::Worker,
&HashMap::new(),
&["worker/1".into()],
));
assert!(!constraints.is_satisfied_by(
ANY_NODE,
&NodeClass::Worker,
&HashMap::new(),
&["replica/0".into(), "worker/1".into()],
));
}
#[test]
fn test_actor_spec_builder() {
let spec = ActorSpec::new("worker/0", "app::Worker")
.with_state(vec![1, 2, 3])
.with_crash_strategy(CrashStrategy::Abandon);
assert_eq!(spec.label, "worker/0");
assert_eq!(spec.actor_type_name, "app::Worker");
assert_eq!(spec.initial_state, vec![1, 2, 3]);
assert!(matches!(spec.crash_strategy, CrashStrategy::Abandon));
}
}