use std::collections::BTreeMap;
use std::path::Path;
use serde::{Deserialize, Serialize};
use crate::error::{Error, Result};
use crate::event::Labels;
use crate::executor::{CapacityReport, Executor, LocalExecutor};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ExecutorRules {
pub executors: Vec<ExecutorEntry>,
pub rules: Vec<Rule>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ExecutorEntry {
pub name: String,
#[serde(rename = "type")]
pub kind: ExecutorKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub max_load1: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub min_free_mem: Option<String>,
}
impl ExecutorEntry {
pub fn has_capacity(&self, report: &CapacityReport) -> bool {
if report.slots_free == 0 {
return false;
}
if let Some(max) = self.max_load1 {
if report.load1 > max {
return false;
}
}
if let Some(min) = self.min_free_mem.as_deref().and_then(bytes_of) {
if report.mem_free_bytes < min {
return false;
}
}
true
}
}
pub fn bytes_of(text: &str) -> Option<u64> {
let text = text.trim();
let split = text
.find(|c: char| !c.is_ascii_digit() && c != '.')
.unwrap_or(text.len());
let (number, unit) = text.split_at(split);
let number: f64 = number.parse().ok()?;
if !number.is_finite() || number < 0.0 {
return None;
}
let scale: u64 = match unit.trim() {
"" | "B" => 1,
"KiB" => 1 << 10,
"MiB" => 1 << 20,
"GiB" => 1 << 30,
"TiB" => 1 << 40,
_ => return None,
};
let bytes = number * scale as f64;
(bytes.is_finite() && bytes >= 0.0).then_some(bytes as u64)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
#[non_exhaustive]
pub enum ExecutorKind {
Local,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Rule {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub when: Option<Predicate>,
#[serde(rename = "use")]
pub use_executor: String,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Predicate {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub executor_has_capacity: Option<String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub node_label: BTreeMap<String, String>,
}
impl Predicate {
fn is_empty(&self) -> bool {
self.executor_has_capacity.is_none() && self.node_label.is_empty()
}
fn labels_match(&self, labels: &Labels) -> bool {
self.node_label
.iter()
.all(|(key, value)| label_of(labels, key).is_some_and(|actual| &actual == value))
}
}
fn label_of(labels: &Labels, key: &str) -> Option<String> {
match key {
"run_id" => labels.run_id.clone(),
"node" => labels.node.clone(),
"persona" => labels.persona.clone(),
_ => None,
}
}
pub const SELECTABLE_LABELS: &[&str] = &["run_id", "node", "persona"];
impl ExecutorRules {
pub fn shipped_default() -> Self {
Self {
executors: vec![ExecutorEntry {
name: "local".into(),
kind: ExecutorKind::Local,
max_load1: None,
min_free_mem: None,
}],
rules: vec![
Rule {
when: Some(Predicate {
executor_has_capacity: Some("local".into()),
..Predicate::default()
}),
use_executor: "local".into(),
},
Rule {
when: None,
use_executor: "local".into(),
},
],
}
}
pub fn load(path: &Path) -> Result<Self> {
let text = std::fs::read_to_string(path).map_err(|e| Error::Ledger {
path: path.to_path_buf(),
source: e,
})?;
let rules: Self = serde_norway::from_str(&text)
.map_err(|e| Error::Invalid(format!("{}: {e}", path.display())))?;
rules.validate()?;
Ok(rules)
}
pub fn validate(&self) -> Result<()> {
if self.executors.is_empty() {
return Err(Error::Invalid(
"a rules file needs at least one executor".into(),
));
}
for entry in &self.executors {
if let Some(limit) = &entry.min_free_mem {
if bytes_of(limit).is_none() {
return Err(Error::Invalid(format!(
"executor '{}' sets min_free_mem '{limit}', which is not a size: \
write a byte count or one of B, KiB, MiB, GiB, TiB",
entry.name
)));
}
}
}
for rule in &self.rules {
if !self.executors.iter().any(|e| e.name == rule.use_executor) {
return Err(Error::Invalid(format!(
"rule uses executor '{}', which is not declared",
rule.use_executor
)));
}
if let Some(when) = &rule.when {
if when.is_empty() {
return Err(Error::Invalid(
"a rule's `when` names no condition: write \
`executor_has_capacity`, `node_label`, or no `when` at all"
.into(),
));
}
if let Some(tested) = &when.executor_has_capacity {
if !self.executors.iter().any(|e| &e.name == tested) {
return Err(Error::Invalid(format!(
"rule tests executor '{tested}', which is not declared"
)));
}
}
for key in when.node_label.keys() {
if !SELECTABLE_LABELS.contains(&key.as_str()) {
return Err(Error::Invalid(format!(
"rule tests node label '{key}', which is not one of {}",
SELECTABLE_LABELS.join(", ")
)));
}
}
}
}
Ok(())
}
pub fn select(
&self,
pinned: Option<&str>,
labels: &Labels,
report: &dyn Fn(&str) -> CapacityReport,
) -> Result<String> {
if let Some(pinned) = pinned {
if !self.executors.iter().any(|e| e.name == pinned) {
return Err(Error::Invalid(format!(
"node pins executor '{pinned}', which the rules do not declare"
)));
}
return Ok(pinned.to_string());
}
for rule in &self.rules {
let holds = match &rule.when {
None => true,
Some(when) => {
let capacity = when.executor_has_capacity.as_ref().is_none_or(|tested| {
self.executors
.iter()
.find(|e| &e.name == tested)
.is_some_and(|entry| entry.has_capacity(&report(&entry.name)))
});
capacity && when.labels_match(labels)
}
};
if holds {
return Ok(rule.use_executor.clone());
}
}
Err(Error::Invalid(
"no rule matched and none is a fallback: nothing can dispatch".into(),
))
}
}
pub fn executor_for(entry: &ExecutorEntry) -> Box<dyn Executor> {
match entry.kind {
ExecutorKind::Local => Box::new(LocalExecutor),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn report(slots_free: u32, load1: f64, mem_free_bytes: u64) -> CapacityReport {
CapacityReport {
slots_free,
load1,
mem_free_bytes,
}
}
#[test]
fn the_contracts_own_size_syntax_reads_as_bytes() {
assert_eq!(bytes_of("2GiB"), Some(2 * (1 << 30)));
assert_eq!(bytes_of("512MiB"), Some(512 * (1 << 20)));
assert_eq!(bytes_of("1KiB"), Some(1024));
assert_eq!(bytes_of("1TiB"), Some(1u64 << 40));
assert_eq!(bytes_of("4096"), Some(4096));
assert_eq!(bytes_of("4096B"), Some(4096));
assert_eq!(bytes_of(" 1.5GiB "), Some(1_610_612_736));
assert_eq!(bytes_of("2GB"), None);
assert_eq!(bytes_of("lots"), None);
assert_eq!(bytes_of("-1GiB"), None);
assert_eq!(bytes_of(""), None);
}
#[test]
fn an_executor_is_out_of_capacity_when_any_limit_it_declares_is_exceeded() {
let entry = ExecutorEntry {
name: "local".into(),
kind: ExecutorKind::Local,
max_load1: Some(8.0),
min_free_mem: Some("2GiB".into()),
};
assert!(entry.has_capacity(&report(4, 2.0, 8 << 30)));
assert!(!entry.has_capacity(&report(0, 2.0, 8 << 30)), "no slots");
assert!(
!entry.has_capacity(&report(4, 9.0, 8 << 30)),
"over max_load1"
);
assert!(
!entry.has_capacity(&report(4, 2.0, 1 << 30)),
"under min_free_mem"
);
let vague = ExecutorEntry {
min_free_mem: Some("some".into()),
..entry
};
assert!(vague.has_capacity(&report(4, 2.0, 1)));
}
#[test]
fn a_memory_limit_in_a_unit_this_build_cannot_read_is_refused_by_name() {
let unreadable = ExecutorRules {
executors: vec![ExecutorEntry {
name: "local".into(),
kind: ExecutorKind::Local,
max_load1: None,
min_free_mem: Some("2GB".into()),
}],
rules: vec![Rule {
when: None,
use_executor: "local".into(),
}],
};
let refused = unreadable.validate().expect_err("an unreadable unit");
let said = refused.to_string();
assert!(said.contains("min_free_mem"), "{said}");
assert!(said.contains("2GB"), "{said}");
assert!(
said.contains("GiB"),
"the refusal did not say what to write: {said}"
);
for good in ["2GiB", "512MiB", "1024", "2048B", "1KiB", "1TiB"] {
let entry = ExecutorEntry {
min_free_mem: Some(good.to_string()),
..unreadable.executors[0].clone()
};
ExecutorRules {
executors: vec![entry],
..unreadable.clone()
}
.validate()
.unwrap_or_else(|e| panic!("{good} was refused: {e}"));
}
}
#[test]
fn the_first_rule_whose_predicate_holds_decides() {
let rules = ExecutorRules {
executors: vec![
ExecutorEntry {
name: "fast".into(),
kind: ExecutorKind::Local,
max_load1: Some(1.0),
min_free_mem: None,
},
ExecutorEntry {
name: "slow".into(),
kind: ExecutorKind::Local,
max_load1: None,
min_free_mem: None,
},
],
rules: vec![
Rule {
when: Some(Predicate {
executor_has_capacity: Some("fast".into()),
..Predicate::default()
}),
use_executor: "fast".into(),
},
Rule {
when: None,
use_executor: "slow".into(),
},
],
};
rules
.validate()
.expect("both rules name declared executors");
let idle = |_: &str| report(4, 0.5, u64::MAX);
assert_eq!(
rules
.select(None, &Labels::default(), &idle)
.expect("a rule matched"),
"fast"
);
let busy = |_: &str| report(4, 9.0, u64::MAX);
assert_eq!(
rules
.select(None, &Labels::default(), &busy)
.expect("the fallback"),
"slow"
);
}
fn by_persona() -> ExecutorRules {
ExecutorRules {
executors: vec![
ExecutorEntry {
name: "review-pool".into(),
kind: ExecutorKind::Local,
max_load1: None,
min_free_mem: None,
},
ExecutorEntry {
name: "local".into(),
kind: ExecutorKind::Local,
max_load1: None,
min_free_mem: None,
},
],
rules: vec![
Rule {
when: Some(Predicate {
executor_has_capacity: Some("review-pool".into()),
node_label: BTreeMap::from([("persona".into(), "reviewer".into())]),
}),
use_executor: "review-pool".into(),
},
Rule {
when: None,
use_executor: "local".into(),
},
],
}
}
fn labelled(node: &str, persona: &str) -> Labels {
Labels {
run_id: Some("demo".into()),
round: Some(2),
node: Some(node.into()),
persona: Some(persona.into()),
..Labels::default()
}
}
#[test]
fn a_node_label_predicate_matches_the_nodes_own_labels_exactly() {
let rules = by_persona();
rules.validate().expect("both families are legal");
let idle = |_: &str| report(4, 0.0, u64::MAX);
assert_eq!(
rules
.select(None, &labelled("audit", "reviewer"), &idle)
.expect("the label rule holds"),
"review-pool"
);
assert_eq!(
rules
.select(None, &labelled("audit", "reviewer-2"), &idle)
.expect("the fallback"),
"local"
);
assert_eq!(
rules
.select(None, &Labels::default(), &idle)
.expect("the fallback"),
"local"
);
}
#[test]
fn the_conditions_in_one_when_conjoin() {
let mut rules = by_persona();
let exhausted = |_: &str| report(0, 0.0, u64::MAX);
assert_eq!(
rules
.select(None, &labelled("audit", "reviewer"), &exhausted)
.expect("the fallback"),
"local"
);
rules.rules[0].when = Some(Predicate {
executor_has_capacity: Some("review-pool".into()),
node_label: BTreeMap::from([("node".into(), "audit".into())]),
});
let idle = |_: &str| report(4, 0.0, u64::MAX);
assert_eq!(
rules
.select(None, &labelled("build", "reviewer"), &idle)
.expect("the fallback"),
"local"
);
}
#[test]
fn a_when_that_names_nothing_or_names_a_label_that_is_not_one_is_refused() {
let mut rules = by_persona();
rules.rules[0].when = Some(Predicate::default());
let message = rules.validate().unwrap_err().to_string();
assert!(message.contains("names no condition"), "{message}");
let mut rules = by_persona();
rules.rules[0].when = Some(Predicate {
executor_has_capacity: None,
node_label: BTreeMap::from([("presona".into(), "reviewer".into())]),
});
let message = rules.validate().unwrap_err().to_string();
assert!(message.contains("presona"), "{message}");
}
#[test]
fn a_label_predicate_round_trips_through_the_rules_file_syntax() {
let yaml = "executors:\n - {name: local, type: local}\nrules:\n \
- when: {node_label: {persona: reviewer}}\n use: local\n \
- use: local\n";
let rules: ExecutorRules = serde_norway::from_str(yaml).expect("it parses");
rules.validate().expect("it validates");
assert_eq!(
rules.rules[0]
.when
.as_ref()
.expect("a predicate")
.node_label,
BTreeMap::from([("persona".to_string(), "reviewer".to_string())])
);
let again: ExecutorRules =
serde_norway::from_str(&serde_norway::to_string(&rules).expect("serializes"))
.expect("re-parses");
assert_eq!(again, rules);
}
#[test]
fn a_node_that_pins_an_executor_gets_it_or_is_refused_by_name() {
let rules = ExecutorRules::shipped_default();
let idle = |_: &str| report(4, 0.0, u64::MAX);
assert_eq!(
rules
.select(Some("local"), &Labels::default(), &idle)
.expect("pinned"),
"local"
);
let message = rules
.select(Some("kubernetes"), &Labels::default(), &idle)
.unwrap_err()
.to_string();
assert!(message.contains("kubernetes"), "{message}");
}
#[test]
fn a_rules_file_naming_an_undeclared_executor_is_refused() {
let undeclared_use = ExecutorRules {
executors: vec![ExecutorEntry {
name: "local".into(),
kind: ExecutorKind::Local,
max_load1: None,
min_free_mem: None,
}],
rules: vec![Rule {
when: None,
use_executor: "elsewhere".into(),
}],
};
assert!(undeclared_use
.validate()
.unwrap_err()
.to_string()
.contains("not declared"));
let undeclared_test = ExecutorRules {
executors: vec![ExecutorEntry {
name: "local".into(),
kind: ExecutorKind::Local,
max_load1: None,
min_free_mem: None,
}],
rules: vec![Rule {
when: Some(Predicate {
executor_has_capacity: Some("elsewhere".into()),
..Predicate::default()
}),
use_executor: "local".into(),
}],
};
assert!(undeclared_test
.validate()
.unwrap_err()
.to_string()
.contains("tests executor"));
let empty = ExecutorRules {
executors: vec![],
rules: vec![],
};
assert!(empty
.validate()
.unwrap_err()
.to_string()
.contains("at least one"));
}
#[test]
fn a_rules_file_with_no_matching_rule_and_no_fallback_says_so() {
let rules = ExecutorRules {
executors: vec![ExecutorEntry {
name: "local".into(),
kind: ExecutorKind::Local,
max_load1: Some(0.0),
min_free_mem: None,
}],
rules: vec![Rule {
when: Some(Predicate {
executor_has_capacity: Some("local".into()),
..Predicate::default()
}),
use_executor: "local".into(),
}],
};
let busy = |_: &str| report(4, 5.0, u64::MAX);
let message = rules
.select(None, &Labels::default(), &busy)
.unwrap_err()
.to_string();
assert!(message.contains("nothing can dispatch"), "{message}");
}
#[test]
fn the_shipped_example_file_loads_and_validates() {
let path = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("examples/executors.yaml");
let rules = ExecutorRules::load(&path).expect("the shipped example is legal");
assert_eq!(rules.executors[0].name, "local");
assert_eq!(rules.executors[0].min_free_mem.as_deref(), Some("2GiB"));
assert_eq!(rules.rules.len(), 2);
let idle = |_: &str| report(4, 0.0, u64::MAX);
assert_eq!(
rules
.select(None, &Labels::default(), &idle)
.expect("a rule matched"),
"local"
);
}
#[test]
fn a_missing_or_malformed_rules_file_is_refused_at_its_boundary() {
let missing = std::path::Path::new("no/such/executors.yaml");
assert!(ExecutorRules::load(missing).is_err());
let dir = std::env::temp_dir().join(format!("onepipeline-rules-{}", std::process::id()));
std::fs::create_dir_all(&dir).expect("a scratch directory");
let path = dir.join("executors.yaml");
std::fs::write(
&path,
"executors: [{name: local, type: local, typo: 1}]\nrules: []\n",
)
.expect("written");
let message = ExecutorRules::load(&path).unwrap_err().to_string();
assert!(message.contains("typo"), "{message}");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_declared_local_entry_resolves_to_the_local_executor() {
let entry = ExecutorEntry {
name: "local".into(),
kind: ExecutorKind::Local,
max_load1: None,
min_free_mem: None,
};
assert_eq!(executor_for(&entry).name(), "local");
}
}