use std::collections::HashMap;
use chrono::{DateTime, Utc};
use crate::service::{collect_metric_buckets, resolve_stat};
use crate::state::{AlarmState, CloudWatchState, MetricAlarm, MetricDatum};
pub(crate) fn evaluate_alarms(state: &mut CloudWatchState, region: &str, now: DateTime<Utc>) {
let mut transitions: Vec<(String, String, AlarmState, String)> = Vec::new();
{
let data_map = state.metrics.get(region);
if let Some(alarms) = state.alarms.get_mut(region) {
for alarm in alarms.values_mut() {
let data: Option<&[MetricDatum]> = alarm
.namespace
.as_ref()
.and_then(|ns| data_map.and_then(|m| m.get(ns)))
.map(|v| v.as_slice());
let manually_set = alarm.state_manually_set;
if let Some((new_state, reason, had_data)) =
evaluate_metric_alarm(alarm, data, now, manually_set)
{
if had_data {
alarm.state_manually_set = false;
}
if new_state != alarm.state_value {
let old = alarm.state_value.as_str().to_string();
alarm.state_value = new_state;
alarm.state_reason = reason.clone();
alarm.state_updated_timestamp = now;
transitions.push((alarm.alarm_name.clone(), old, new_state, reason));
}
}
}
}
}
for (name, old, new, reason) in &transitions {
record_transition(state, region, name, "MetricAlarm", old, *new, reason, now);
}
evaluate_composite_alarms(state, region, now);
}
fn evaluate_composite_alarms(state: &mut CloudWatchState, region: &str, now: DateTime<Utc>) {
let composite_names: Vec<String> = match state.composite_alarms.get(region) {
Some(c) if !c.is_empty() => c.keys().cloned().collect(),
_ => return,
};
let mut states: HashMap<String, AlarmState> = HashMap::new();
if let Some(a) = state.alarms.get(region) {
for (n, al) in a.iter() {
states.insert(n.clone(), al.state_value);
}
}
if let Some(c) = state.composite_alarms.get(region) {
for (n, al) in c.iter() {
states.insert(n.clone(), al.state_value);
}
}
for _ in 0..composite_names.len() {
let mut changed = false;
for name in &composite_names {
let Some(rule) = state
.composite_alarms
.get(region)
.and_then(|c| c.get(name))
.map(|a| a.alarm_rule.clone())
else {
continue;
};
if let Some(new_state) = evaluate_rule(&rule, &states) {
if states.get(name) != Some(&new_state) {
states.insert(name.clone(), new_state);
changed = true;
}
}
}
if !changed {
break;
}
}
let mut transitions: Vec<(String, String, AlarmState)> = Vec::new();
if let Some(c) = state.composite_alarms.get_mut(region) {
for (name, al) in c.iter_mut() {
let Some(&new_state) = states.get(name) else {
continue;
};
if new_state != al.state_value {
let old = al.state_value.as_str().to_string();
al.state_value = new_state;
al.state_reason = format!(
"The composite alarm rule evaluated to {}.",
new_state.as_str()
);
al.state_updated_timestamp = now;
transitions.push((name.clone(), old, new_state));
}
}
}
for (name, old, new) in &transitions {
let reason = format!("The composite alarm rule evaluated to {}.", new.as_str());
record_transition(
state,
region,
name,
"CompositeAlarm",
old,
*new,
&reason,
now,
);
}
}
pub(crate) fn evaluate_metric_alarm(
alarm: &MetricAlarm,
data: Option<&[MetricDatum]>,
now: DateTime<Utc>,
manually_set: bool,
) -> Option<(AlarmState, String, bool)> {
if !alarm.metrics.is_empty() || alarm.threshold_metric_id.is_some() {
return None;
}
let threshold = alarm.threshold?;
alarm.namespace.as_ref()?;
let metric_name = alarm.metric_name.as_ref()?;
let stat = alarm
.extended_statistic
.clone()
.or_else(|| alarm.statistic.clone())?;
let period = alarm.period.unwrap_or(60).max(1);
let eval_periods = alarm.evaluation_periods.max(1);
let m = alarm.datapoints_to_alarm.unwrap_or(eval_periods).max(1);
let treat = alarm
.treat_missing_data
.as_deref()
.unwrap_or("missing")
.to_ascii_lowercase();
let now_secs = now.timestamp();
let latest = now_secs - now_secs.rem_euclid(period);
let window_start = latest - (eval_periods - 1) * period;
let window_end = latest + period; let start_ts = DateTime::<Utc>::from_timestamp(window_start, 0)?;
let end_ts = DateTime::<Utc>::from_timestamp(window_end, 0)?;
let buckets = match data {
Some(d) => collect_metric_buckets(
d,
metric_name,
&alarm.dimensions,
alarm.unit.as_deref(),
period,
start_ts,
end_ts,
),
None => Default::default(),
};
let mut present = 0i64;
let mut breaching = 0i64;
let mut latest_value: Option<f64> = None;
for i in 0..eval_periods {
let slot_secs = latest - i * period;
let Some(slot_ts) = DateTime::<Utc>::from_timestamp(slot_secs, 0) else {
continue;
};
if let Some(bucket) = buckets.get(&slot_ts) {
let mut sorted = bucket.samples.clone();
sorted.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
if let Some(v) = resolve_stat(&stat, bucket, &sorted) {
present += 1;
if latest_value.is_none() {
latest_value = Some(v);
}
if breaches(&alarm.comparison_operator, v, threshold) {
breaching += 1;
}
}
}
}
let missing = eval_periods - present;
if present == 0 {
return match treat.as_str() {
"breaching" => Some((
AlarmState::Alarm,
"Insufficient data treated as breaching.".to_string(),
false,
)),
"notbreaching" => Some((
AlarmState::Ok,
"Insufficient data treated as not breaching.".to_string(),
false,
)),
"ignore" => None,
_ if manually_set => None,
_ => Some((
AlarmState::InsufficientData,
format!("Insufficient Data: {eval_periods} datapoints were unknown."),
false,
)),
};
}
let effective_breaching = match treat.as_str() {
"breaching" => breaching + missing,
_ => breaching,
};
if effective_breaching >= m {
let reason = format!(
"Threshold Crossed: {effective_breaching} out of the last {eval_periods} datapoints \
[{}] was {} the threshold ({threshold}).",
latest_value.map(|v| v.to_string()).unwrap_or_default(),
operator_phrase(&alarm.comparison_operator),
);
Some((AlarmState::Alarm, reason, true))
} else {
let reason = format!(
"Threshold not crossed: {breaching} out of the last {eval_periods} datapoints was {} \
the threshold ({threshold}).",
operator_phrase(&alarm.comparison_operator),
);
Some((AlarmState::Ok, reason, true))
}
}
fn breaches(op: &str, value: f64, threshold: f64) -> bool {
match op {
"GreaterThanThreshold" => value > threshold,
"GreaterThanOrEqualToThreshold" => value >= threshold,
"LessThanThreshold" => value < threshold,
"LessThanOrEqualToThreshold" => value <= threshold,
_ => false,
}
}
fn operator_phrase(op: &str) -> &'static str {
match op {
"GreaterThanThreshold" => "greater than",
"GreaterThanOrEqualToThreshold" => "greater than or equal to",
"LessThanThreshold" => "less than",
"LessThanOrEqualToThreshold" => "less than or equal to",
_ => "compared against",
}
}
#[allow(clippy::too_many_arguments)]
fn record_transition(
state: &mut CloudWatchState,
region: &str,
name: &str,
alarm_type: &str,
old: &str,
new: AlarmState,
reason: &str,
_now: DateTime<Utc>,
) {
let new_str = new.as_str();
let summary = format!("Alarm updated from {old} to {new_str}");
let history_data = format!(
"{{\"oldState\":{{\"stateValue\":\"{old}\"}},\"newState\":{{\"stateValue\":\"{new_str}\",\"stateReason\":\"{}\"}}}}",
reason.replace('"', "\\\"")
);
crate::service::push_alarm_history(
state,
region,
name,
alarm_type,
"StateUpdate",
summary,
history_data,
);
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) enum RuleNode {
And(Box<RuleNode>, Box<RuleNode>),
Or(Box<RuleNode>, Box<RuleNode>),
Not(Box<RuleNode>),
State(AlarmState, String),
True,
False,
}
fn evaluate_rule(rule: &str, states: &HashMap<String, AlarmState>) -> Option<AlarmState> {
let node = parse_alarm_rule(rule).ok()?;
Some(if eval_node(&node, states) {
AlarmState::Alarm
} else {
AlarmState::Ok
})
}
fn eval_node(node: &RuleNode, states: &HashMap<String, AlarmState>) -> bool {
match node {
RuleNode::And(a, b) => eval_node(a, states) && eval_node(b, states),
RuleNode::Or(a, b) => eval_node(a, states) || eval_node(b, states),
RuleNode::Not(a) => !eval_node(a, states),
RuleNode::True => true,
RuleNode::False => false,
RuleNode::State(want, name) => states.get(name.as_str()) == Some(want),
}
}
pub(crate) fn parse_alarm_rule(rule: &str) -> Result<RuleNode, String> {
let chars: Vec<char> = rule.chars().collect();
let mut p = RuleParser { chars, pos: 0 };
p.skip_ws();
if p.pos >= p.chars.len() {
return Err("empty alarm rule".to_string());
}
let node = p.parse_or()?;
p.skip_ws();
if p.pos != p.chars.len() {
return Err(format!(
"unexpected trailing input in alarm rule at {}",
p.pos
));
}
Ok(node)
}
struct RuleParser {
chars: Vec<char>,
pos: usize,
}
impl RuleParser {
fn skip_ws(&mut self) {
while self.pos < self.chars.len() && self.chars[self.pos].is_whitespace() {
self.pos += 1;
}
}
fn peek(&self) -> Option<char> {
self.chars.get(self.pos).copied()
}
fn peek_keyword(&self) -> Option<String> {
let mut i = self.pos;
while i < self.chars.len() && self.chars[i].is_whitespace() {
i += 1;
}
let start = i;
while i < self.chars.len() && (self.chars[i].is_ascii_alphabetic() || self.chars[i] == '_')
{
i += 1;
}
if i == start {
None
} else {
Some(
self.chars[start..i]
.iter()
.collect::<String>()
.to_ascii_uppercase(),
)
}
}
fn consume_keyword(&mut self, kw: &str) {
self.skip_ws();
self.pos += kw.len();
}
fn parse_or(&mut self) -> Result<RuleNode, String> {
let mut left = self.parse_and()?;
while self.peek_keyword().as_deref() == Some("OR") {
self.consume_keyword("OR");
let right = self.parse_and()?;
left = RuleNode::Or(Box::new(left), Box::new(right));
}
Ok(left)
}
fn parse_and(&mut self) -> Result<RuleNode, String> {
let mut left = self.parse_unary()?;
while self.peek_keyword().as_deref() == Some("AND") {
self.consume_keyword("AND");
let right = self.parse_unary()?;
left = RuleNode::And(Box::new(left), Box::new(right));
}
Ok(left)
}
fn parse_unary(&mut self) -> Result<RuleNode, String> {
if self.peek_keyword().as_deref() == Some("NOT") {
self.consume_keyword("NOT");
let inner = self.parse_unary()?;
return Ok(RuleNode::Not(Box::new(inner)));
}
self.parse_primary()
}
fn parse_primary(&mut self) -> Result<RuleNode, String> {
self.skip_ws();
match self.peek() {
Some('(') => {
self.pos += 1;
let node = self.parse_or()?;
self.skip_ws();
if self.peek() != Some(')') {
return Err("expected ')' in alarm rule".to_string());
}
self.pos += 1;
Ok(node)
}
Some(_) => {
let kw = self
.peek_keyword()
.ok_or_else(|| "expected a token in alarm rule".to_string())?;
match kw.as_str() {
"TRUE" => {
self.consume_keyword("TRUE");
Ok(RuleNode::True)
}
"FALSE" => {
self.consume_keyword("FALSE");
Ok(RuleNode::False)
}
"ALARM" | "OK" | "INSUFFICIENT_DATA" => self.parse_func(&kw),
other => Err(format!("unexpected token '{other}' in alarm rule")),
}
}
None => Err("unexpected end of alarm rule".to_string()),
}
}
fn parse_func(&mut self, kind: &str) -> Result<RuleNode, String> {
self.skip_ws();
while self.pos < self.chars.len()
&& (self.chars[self.pos].is_ascii_alphabetic() || self.chars[self.pos] == '_')
{
self.pos += 1;
}
self.skip_ws();
if self.peek() != Some('(') {
return Err(format!("expected '(' after {kind} in alarm rule"));
}
self.pos += 1;
let start = self.pos;
while self.pos < self.chars.len() && self.chars[self.pos] != ')' {
self.pos += 1;
}
if self.peek() != Some(')') {
return Err(format!("unterminated {kind}(...) in alarm rule"));
}
let raw: String = self.chars[start..self.pos].iter().collect();
self.pos += 1; let name = normalize_ref(raw.trim());
if name.is_empty() {
return Err(format!("empty reference in {kind}(...)"));
}
let state = match kind {
"ALARM" => AlarmState::Alarm,
"OK" => AlarmState::Ok,
"INSUFFICIENT_DATA" => AlarmState::InsufficientData,
_ => unreachable!(),
};
Ok(RuleNode::State(state, name))
}
}
fn normalize_ref(raw: &str) -> String {
let unquoted = raw
.strip_prefix('"')
.and_then(|s| s.strip_suffix('"'))
.or_else(|| raw.strip_prefix('\'').and_then(|s| s.strip_suffix('\'')))
.unwrap_or(raw)
.trim();
if let Some(idx) = unquoted.find(":alarm:") {
return unquoted[idx + ":alarm:".len()..].to_string();
}
unquoted.to_string()
}
#[cfg(test)]
mod tests {
use super::*;
fn states(pairs: &[(&str, AlarmState)]) -> HashMap<String, AlarmState> {
pairs.iter().map(|(n, s)| (n.to_string(), *s)).collect()
}
#[test]
fn parses_and_evaluates_simple_alarm_predicate() {
let node = parse_alarm_rule("ALARM(a)").unwrap();
assert!(eval_node(&node, &states(&[("a", AlarmState::Alarm)])));
assert!(!eval_node(&node, &states(&[("a", AlarmState::Ok)])));
assert!(!eval_node(&node, &states(&[])));
}
#[test]
fn evaluates_and_or_not_and_parens() {
let node = parse_alarm_rule("(ALARM(a) AND ALARM(b)) OR NOT OK(c)").unwrap();
let s = states(&[
("a", AlarmState::Alarm),
("b", AlarmState::Alarm),
("c", AlarmState::Ok),
]);
assert!(eval_node(&node, &s));
let s2 = states(&[
("a", AlarmState::Ok),
("b", AlarmState::Alarm),
("c", AlarmState::Alarm),
]);
assert!(eval_node(&node, &s2));
}
#[test]
fn quoted_and_arn_references_resolve() {
let node = parse_alarm_rule("ALARM(\"my-alarm\")").unwrap();
assert!(eval_node(
&node,
&states(&[("my-alarm", AlarmState::Alarm)])
));
let node =
parse_alarm_rule("ALARM(arn:aws:cloudwatch:us-east-1:123456789012:alarm:my-alarm)")
.unwrap();
assert!(eval_node(
&node,
&states(&[("my-alarm", AlarmState::Alarm)])
));
}
#[test]
fn insufficient_data_predicate() {
let node = parse_alarm_rule("INSUFFICIENT_DATA(a)").unwrap();
assert!(eval_node(
&node,
&states(&[("a", AlarmState::InsufficientData)])
));
}
#[test]
fn malformed_rules_error() {
assert!(parse_alarm_rule("").is_err());
assert!(parse_alarm_rule("ALARM(").is_err());
assert!(parse_alarm_rule("ALARM(a) AND").is_err());
assert!(parse_alarm_rule("BOGUS(a)").is_err());
assert!(parse_alarm_rule("ALARM(a) garbage").is_err());
}
#[test]
fn evaluate_rule_maps_to_alarm_or_ok() {
let s = states(&[("a", AlarmState::Alarm)]);
assert_eq!(evaluate_rule("ALARM(a)", &s), Some(AlarmState::Alarm));
assert_eq!(evaluate_rule("OK(a)", &s), Some(AlarmState::Ok));
assert_eq!(evaluate_rule("!!!", &s), None);
}
fn sum_alarm(state: AlarmState, treat: Option<&str>) -> MetricAlarm {
MetricAlarm {
alarm_name: "a".into(),
alarm_arn: "arn:aws:cloudwatch:us-east-1:123456789012:alarm:a".into(),
alarm_description: None,
actions_enabled: true,
ok_actions: vec![],
alarm_actions: vec![],
insufficient_data_actions: vec![],
state_value: state,
state_reason: String::new(),
state_updated_timestamp: Utc::now(),
metric_name: Some("M".into()),
namespace: Some("Test/NS".into()),
statistic: Some("Sum".into()),
extended_statistic: None,
dimensions: Default::default(),
period: Some(60),
unit: None,
evaluation_periods: 1,
datapoints_to_alarm: None,
threshold: Some(1.0),
comparison_operator: "GreaterThanThreshold".into(),
treat_missing_data: treat.map(|s| s.to_string()),
evaluate_low_sample_count_percentile: None,
threshold_metric_id: None,
configuration_updated_timestamp: Utc::now(),
alarm_configuration_updated_timestamp: Utc::now(),
metrics: vec![],
state_manually_set: false,
}
}
fn datum(value: f64, ts: DateTime<Utc>) -> MetricDatum {
MetricDatum {
metric_name: "M".into(),
dimensions: Default::default(),
timestamp: ts,
value: Some(value),
statistic_values: None,
unit: None,
storage_resolution: None,
}
}
#[test]
fn missing_data_default_transitions_to_insufficient_data() {
let now = Utc::now();
let alarm = sum_alarm(AlarmState::Alarm, None);
let out = evaluate_metric_alarm(&alarm, None, now, false);
assert!(matches!(
out,
Some((AlarmState::InsufficientData, _, false))
));
}
#[test]
fn missing_data_ignore_keeps_state() {
let now = Utc::now();
let alarm = sum_alarm(AlarmState::Alarm, Some("ignore"));
assert_eq!(evaluate_metric_alarm(&alarm, None, now, false), None);
}
#[test]
fn manual_state_survives_missing_data() {
let now = Utc::now();
let alarm = sum_alarm(AlarmState::Alarm, None);
assert_eq!(evaluate_metric_alarm(&alarm, None, now, true), None);
}
#[test]
fn present_data_drives_state_and_reports_had_data() {
let now = Utc::now();
let alarm = sum_alarm(AlarmState::InsufficientData, None);
let data = [datum(5.0, now)];
let out = evaluate_metric_alarm(&alarm, Some(&data), now, true);
assert!(matches!(out, Some((AlarmState::Alarm, _, true))), "{out:?}");
}
}