use core::fmt;
use std::collections::BTreeMap;
use serde::Deserialize;
use shep_core::barks::Bark;
use shep_core::protocol::{BusEvent, ProcessEventKind, ProcessInfo};
use shep_core::status::ProcStatus;
use shep_core::values::{MemSize, UpDuration};
use super::sinks::{self, Sink, SinkConfigError};
fn default_debounce() -> UpDuration {
UpDuration::from_millis(5 * 60 * 1_000)
}
#[derive(Debug, Clone, PartialEq, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Rule {
#[serde(flatten)]
pub when: Trigger,
pub sinks: Vec<String>,
#[serde(default = "default_debounce")]
pub debounce: UpDuration,
}
#[derive(Debug, Clone, PartialEq, Deserialize)]
#[serde(tag = "on", rename_all = "snake_case")]
pub enum Trigger {
Event {
kinds: Vec<String>,
},
GaveUp,
RestartRate {
restarts: u32,
within: UpDuration,
},
MemoryAbove {
bytes: MemSize,
},
}
fn trigger_name(when: &Trigger) -> &'static str {
match when {
Trigger::Event { .. } => "event",
Trigger::GaveUp => "gave_up",
Trigger::RestartRate { .. } => "restart_rate",
Trigger::MemoryAbove { .. } => "memory_above",
}
}
fn wire_spelling(kind: ProcessEventKind) -> String {
serde_json::to_value(kind)
.ok()
.and_then(|value| value.as_str().map(str::to_owned))
.unwrap_or_default()
}
fn is_known_kind(kind: &str) -> bool {
serde_json::from_value::<ProcessEventKind>(serde_json::Value::String(kind.to_owned())).is_ok()
}
#[derive(Debug)]
pub enum RulesError {
UnknownSink {
index: usize,
sink: String,
},
NoSinks {
index: usize,
},
UnknownKind {
index: usize,
kind: String,
},
InsecureSink(SinkConfigError),
}
impl fmt::Display for RulesError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::UnknownSink { index, sink } => write!(
f,
"rule {index} routes to sink \"{sink}\", which [dog.bark.sinks] does not define"
),
Self::NoSinks { index } => write!(f, "rule {index} routes to no sink at all"),
Self::UnknownKind { index, kind } => write!(
f,
"rule {index}'s event trigger names \"{kind}\", which is not an event kind on the wire"
),
Self::InsecureSink(source) => write!(f, "{source}"),
}
}
}
impl core::error::Error for RulesError {
fn source(&self) -> Option<&(dyn core::error::Error + 'static)> {
match self {
Self::InsecureSink(source) => Some(source),
Self::UnknownSink { .. } | Self::NoSinks { .. } | Self::UnknownKind { .. } => None,
}
}
}
impl From<SinkConfigError> for RulesError {
fn from(source: SinkConfigError) -> Self {
Self::InsecureSink(source)
}
}
#[derive(Debug, Default)]
struct SubjectState {
last_fired: BTreeMap<usize, u64>,
restart_windows: BTreeMap<usize, (u32, u64)>,
}
#[derive(Debug)]
pub struct Rules {
rules: Vec<Rule>,
subjects: BTreeMap<String, SubjectState>,
}
impl Rules {
pub fn new(rules: Vec<Rule>, sinks: &BTreeMap<String, Sink>) -> Result<Self, RulesError> {
for (name, sink) in sinks {
sinks::require_secure_scheme(name, sink)?;
}
for (index, rule) in rules.iter().enumerate() {
if rule.sinks.is_empty() {
return Err(RulesError::NoSinks { index });
}
for sink in &rule.sinks {
if !sinks.contains_key(sink) {
return Err(RulesError::UnknownSink {
index,
sink: sink.clone(),
});
}
}
if let Trigger::Event { kinds } = &rule.when {
for kind in kinds {
if !is_known_kind(kind) {
return Err(RulesError::UnknownKind {
index,
kind: kind.clone(),
});
}
}
}
}
Ok(Self {
rules,
subjects: BTreeMap::new(),
})
}
#[must_use]
pub fn default_rules(sinks: &BTreeMap<String, Sink>) -> Vec<Rule> {
vec![Rule {
when: Trigger::GaveUp,
sinks: sinks.keys().cloned().collect(),
debounce: default_debounce(),
}]
}
fn try_fire(&mut self, idx: usize, subject: &str, now_ms: u64, debounce: UpDuration) -> bool {
let state = self.subjects.entry(subject.to_owned()).or_default();
let ready = state
.last_fired
.get(&idx)
.is_none_or(|&last| now_ms.saturating_sub(last) >= debounce.as_millis());
if ready {
state.last_fired.insert(idx, now_ms);
}
ready
}
fn restart_window_crossed(
&mut self,
idx: usize,
subject: &str,
current_restarts: u32,
threshold: u32,
within: UpDuration,
now_ms: u64,
) -> bool {
let state = self.subjects.entry(subject.to_owned()).or_default();
let window = state.restart_windows.entry(idx).or_insert((0, now_ms));
if now_ms.saturating_sub(window.1) > within.as_millis() {
*window = (current_restarts, now_ms);
}
current_restarts.saturating_sub(window.0) >= threshold
}
#[must_use]
pub fn on_event(&mut self, event: &BusEvent, now_ms: u64) -> Vec<Firing> {
let BusEvent::Process {
event: kind, info, ..
} = event
else {
return Vec::new();
};
let kind = *kind;
let kind_wire = wire_spelling(kind);
let mut firings = Vec::new();
for idx in 0..self.rules.len() {
let debounce = self.rules[idx].debounce;
let trigger = self.rules[idx].when.clone();
let message = match &trigger {
Trigger::Event { kinds } if kinds.iter().any(|k| k == &kind_wire) => {
Some(format!("{} {kind_wire}", info.name))
}
Trigger::GaveUp if kind == ProcessEventKind::Errored => {
Some(format!("{} gave up: restart budget exhausted", info.name))
}
_ => None,
};
let Some(message) = message else { continue };
if !self.try_fire(idx, &info.name, now_ms, debounce) {
continue;
}
let sinks = self.rules[idx].sinks.clone();
firings.push(Firing {
bark: Bark {
at_ms: now_ms,
rule: trigger_name(&trigger).to_owned(),
subject: info.name.clone(),
message,
sinks: Vec::new(),
},
sinks,
});
}
firings
}
#[must_use]
pub fn on_poll(&mut self, flock: &[ProcessInfo], now_ms: u64) -> Vec<Firing> {
let mut firings = Vec::new();
for info in flock {
for idx in 0..self.rules.len() {
let debounce = self.rules[idx].debounce;
let trigger = self.rules[idx].when.clone();
let message = match &trigger {
Trigger::Event { .. } => None,
Trigger::GaveUp => (info.status == ProcStatus::Errored).then(|| {
format!("{} gave up: restart budget exhausted", info.name)
}),
Trigger::RestartRate { restarts, within } => self
.restart_window_crossed(
idx,
&info.name,
info.restarts,
*restarts,
*within,
now_ms,
)
.then(|| {
format!(
"{} restarted {} times, at or past the {restarts}-within-{within} early warning",
info.name, info.restarts
)
}),
Trigger::MemoryAbove { bytes } => info.memory_bytes.and_then(|used| {
(used >= bytes.bytes()).then(|| {
format!(
"{} memory at {}, at or above the {bytes} limit",
info.name,
MemSize::from_bytes(used)
)
})
}),
};
let Some(message) = message else { continue };
if !self.try_fire(idx, &info.name, now_ms, debounce) {
continue;
}
let sinks = self.rules[idx].sinks.clone();
firings.push(Firing {
bark: Bark {
at_ms: now_ms,
rule: trigger_name(&trigger).to_owned(),
subject: info.name.clone(),
message,
sinks: Vec::new(),
},
sinks,
});
}
}
firings
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Firing {
pub bark: Bark,
pub sinks: Vec<String>,
}
#[cfg(test)]
mod tests {
use shep_core::protocol::ProcessInfo;
use super::*;
fn one_sink(name: &str) -> BTreeMap<String, Sink> {
let mut sinks = BTreeMap::new();
sinks.insert(
name.to_owned(),
Sink::Json {
url: "http://localhost/hook".to_owned(),
body: None,
},
);
sinks
}
fn base_info(name: &str, status: ProcStatus) -> ProcessInfo {
ProcessInfo::builder(1, name, status)
.pid(Some(4242))
.uptime_ms(1_000)
.build()
}
fn errored_info(name: &str) -> ProcessInfo {
base_info(name, ProcStatus::Errored)
}
fn online_info(name: &str) -> ProcessInfo {
base_info(name, ProcStatus::Online)
}
fn process_event(name: &str, kind: ProcessEventKind) -> BusEvent {
BusEvent::Process {
event: kind,
info: base_info(name, ProcStatus::Online),
manually: false,
at_ms: 0,
}
}
fn errored_event(name: &str) -> BusEvent {
process_event(name, ProcessEventKind::Errored)
}
fn restart_event(name: &str) -> BusEvent {
process_event(name, ProcessEventKind::Restart)
}
fn gave_up_rules() -> Rules {
let sinks = one_sink("ops");
Rules::new(
vec![Rule {
when: Trigger::GaveUp,
sinks: vec!["ops".to_owned()],
debounce: default_debounce(),
}],
&sinks,
)
.unwrap()
}
fn restart_rate_rules(restarts: u32, within: UpDuration) -> Rules {
let sinks = one_sink("ops");
Rules::new(
vec![Rule {
when: Trigger::RestartRate { restarts, within },
sinks: vec!["ops".to_owned()],
debounce: default_debounce(),
}],
&sinks,
)
.unwrap()
}
fn rule_to(sink: &str) -> Rule {
Rule {
when: Trigger::GaveUp,
sinks: vec![sink.to_owned()],
debounce: default_debounce(),
}
}
#[test]
fn an_errored_seen_by_both_routes_fires_once() {
let mut rules = gave_up_rules();
let first = rules.on_event(&errored_event("web"), 1_000);
assert_eq!(first.len(), 1);
let second = rules.on_poll(&[errored_info("web")], 2_000);
assert!(second.is_empty(), "the debounce covers the other route");
}
#[test]
fn the_poll_fires_what_the_bus_never_carried() {
let mut rules = gave_up_rules();
let fired = rules.on_poll(&[errored_info("web")], 1_000);
assert_eq!(fired.len(), 1);
assert_eq!(fired[0].bark.subject, "web");
}
#[test]
fn one_flapping_sheep_does_not_mute_another_going_down() {
let mut rules = gave_up_rules();
assert_eq!(rules.on_event(&errored_event("web"), 1_000).len(), 1);
assert_eq!(rules.on_event(&errored_event("api"), 1_100).len(), 1);
assert!(rules.on_event(&errored_event("web"), 1_200).is_empty());
}
#[test]
fn the_early_warning_counts_the_shepherds_restarts_and_not_its_own() {
let mut rules = restart_rate_rules(5, UpDuration::from_millis(60_000));
for at in [1_000, 2_000, 3_000] {
let _ = rules.on_event(&restart_event("web"), at);
}
let mut info = online_info("web");
info.restarts = 9;
let fired = rules.on_poll(&[info], 4_000);
assert_eq!(
fired.len(),
1,
"9 restarts crosses a threshold of 5; 3 does not"
);
}
#[test]
fn a_rule_routed_at_a_sink_that_does_not_exist_is_refused_at_startup() {
let err = Rules::new(vec![rule_to("pager")], &BTreeMap::new()).unwrap_err();
assert!(matches!(err, RulesError::UnknownSink { .. }));
assert!(err.to_string().contains("pager"));
}
#[test]
fn a_bark_with_sinks_and_no_rules_still_alerts_when_the_shepherd_gives_up() {
let sinks = one_sink("ops");
let rules = Rules::default_rules(&sinks);
assert_eq!(rules.len(), 1);
assert_eq!(rules[0].when, Trigger::GaveUp);
assert_eq!(rules[0].sinks, vec!["ops"]);
}
#[test]
fn a_rule_with_no_sinks_at_all_is_refused_at_startup() {
let rule = Rule {
when: Trigger::GaveUp,
sinks: Vec::new(),
debounce: default_debounce(),
};
let err = Rules::new(vec![rule], &one_sink("ops")).unwrap_err();
assert!(matches!(err, RulesError::NoSinks { .. }));
}
#[test]
fn an_event_rule_naming_an_unknown_kind_is_refused_at_startup() {
let rule = Rule {
when: Trigger::Event {
kinds: vec!["exit".to_owned(), "not_a_real_kind".to_owned()],
},
sinks: vec!["ops".to_owned()],
debounce: default_debounce(),
};
let err = Rules::new(vec![rule], &one_sink("ops")).unwrap_err();
assert!(matches!(err, RulesError::UnknownKind { .. }));
assert!(err.to_string().contains("not_a_real_kind"));
}
#[test]
fn an_event_rule_does_not_fire_on_a_kind_it_was_not_given() {
let sinks = one_sink("ops");
let mut rules = Rules::new(
vec![Rule {
when: Trigger::Event {
kinds: vec!["exit".to_owned()],
},
sinks: vec!["ops".to_owned()],
debounce: default_debounce(),
}],
&sinks,
)
.unwrap();
let online = rules.on_event(&process_event("web", ProcessEventKind::Online), 1_000);
assert!(online.is_empty(), "online was not in the configured kinds");
let exit = rules.on_event(&process_event("web", ProcessEventKind::Exit), 1_100);
assert_eq!(exit.len(), 1, "exit was, and should still fire");
}
#[test]
fn restart_rate_fires_at_the_threshold_and_not_one_below_it() {
let mut rules = restart_rate_rules(5, UpDuration::from_millis(60_000));
let mut below = online_info("web");
below.restarts = 4;
assert!(
rules.on_poll(&[below], 1_000).is_empty(),
"4 restarts is below a threshold of 5"
);
let mut at = online_info("web");
at.restarts = 5;
assert_eq!(
rules.on_poll(&[at], 1_100).len(),
1,
"5 restarts meets a threshold of 5"
);
}
#[test]
fn restart_rate_window_slides_once_it_elapses() {
let sinks = one_sink("ops");
let mut rules = Rules::new(
vec![Rule {
when: Trigger::RestartRate {
restarts: 5,
within: UpDuration::from_millis(1_000),
},
sinks: vec!["ops".to_owned()],
debounce: UpDuration::from_millis(0),
}],
&sinks,
)
.unwrap();
let mut info = online_info("web");
info.restarts = 5;
assert_eq!(
rules.on_poll(&[info.clone()], 0).len(),
1,
"5 restarts opens the window past threshold"
);
assert!(
rules.on_poll(&[info.clone()], 2_000).is_empty(),
"no new restarts since the window reset"
);
info.restarts = 10;
assert_eq!(
rules.on_poll(&[info], 2_100).len(),
1,
"5 more restarts inside the new window crosses it again"
);
}
#[test]
fn memory_above_fires_at_the_ceiling_and_not_one_byte_below_it() {
let sinks = one_sink("ops");
let mut rules = Rules::new(
vec![Rule {
when: Trigger::MemoryAbove {
bytes: MemSize::from_bytes(1_000),
},
sinks: vec!["ops".to_owned()],
debounce: default_debounce(),
}],
&sinks,
)
.unwrap();
let mut below = online_info("web");
below.memory_bytes = Some(999);
assert!(
rules.on_poll(&[below], 1_000).is_empty(),
"999 is below 1000"
);
let mut at = online_info("web");
at.memory_bytes = Some(1_000);
assert_eq!(rules.on_poll(&[at], 1_100).len(), 1, "1000 meets 1000");
}
#[test]
fn memory_above_does_not_fire_when_usage_is_unknown() {
let sinks = one_sink("ops");
let mut rules = Rules::new(
vec![Rule {
when: Trigger::MemoryAbove {
bytes: MemSize::from_bytes(1_000),
},
sinks: vec!["ops".to_owned()],
debounce: default_debounce(),
}],
&sinks,
)
.unwrap();
let info = online_info("web");
assert!(info.memory_bytes.is_none());
assert!(rules.on_poll(&[info], 1_000).is_empty());
}
#[test]
fn gave_up_does_not_fire_on_event_for_a_non_errored_kind() {
let mut rules = gave_up_rules();
let online = rules.on_event(&process_event("web", ProcessEventKind::Online), 1_000);
assert!(online.is_empty(), "GaveUp fires on Errored only");
let restart = rules.on_event(&restart_event("web"), 1_100);
assert!(restart.is_empty(), "GaveUp fires on Errored only");
}
#[test]
fn gave_up_does_not_fire_on_poll_for_a_non_errored_status() {
let mut rules = gave_up_rules();
let fired = rules.on_poll(&[online_info("web")], 1_000);
assert!(fired.is_empty(), "GaveUp fires when status is Errored only");
}
#[test]
fn debounce_boundary_is_inclusive_at_exactly_its_own_duration() {
let mut rules = gave_up_rules();
let debounce_ms = default_debounce().as_millis();
assert_eq!(rules.on_event(&errored_event("web"), 0).len(), 1);
assert!(
rules
.on_event(&errored_event("web"), debounce_ms - 1)
.is_empty(),
"one millisecond short of the debounce must still be quiet"
);
assert_eq!(
rules.on_event(&errored_event("web"), debounce_ms).len(),
1,
"exactly at the debounce it may fire again"
);
}
}