diff --git a/crates/common/src/telemetry/alerts.rs b/crates/common/src/telemetry/alerts.rs new file mode 100644 index 0000000..68999fd --- /dev/null +++ b/crates/common/src/telemetry/alerts.rs @@ -0,0 +1,265 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Alerts (monitoring spec MON-25 to MON-30). Each enabled `x:Alert` is read +//! from the registry at evaluation, so a change needs no reload (MON-3), and +//! fires when its condition goes from false to true (MON-26). + +use crate::{ + Server, + expr::{functions::EmptyResolver, if_block::BootstrapExprExt}, +}; +use ahash::AHashSet; +use mail_builder::{ + MessageBuilder, + headers::{HeaderType, address::Address}, +}; +use registry::{ + schema::{ + prelude::{ExpressionContext, ObjectType}, + structs::{Alert, AlertEmail, AlertEvent, Expression, ExpressionMatch}, + }, + types::{id::ObjectId, list::List}, +}; +use std::sync::Mutex; +use store::registry::{RegistryQuery, bootstrap::Bootstrap}; +use trc::{Collector, MetricType, TelemetryEvent}; +use types::id::Id; + +/// An alert email, ready to queue. +#[derive(Debug, Clone)] +pub struct AlertMessage { + pub from: String, + pub to: Vec, + pub body: Vec, +} + +/// The alerts whose condition held at the last evaluation (MON-26). In +/// memory, so a restart while a condition holds fires once more. +static FIRING: Mutex>> = Mutex::new(None); + +fn is_ident_char(c: char) -> bool { + c.is_ascii_alphanumeric() || c == '_' +} + +/// The metric an underscore name stands for (`queue_count`), MON-25. +fn underscore_metric(name: &str) -> Option { + static NAMES: std::sync::OnceLock> = + std::sync::OnceLock::new(); + if !name.contains('_') { + return None; + } + NAMES + .get_or_init(|| { + (0..=u16::MAX) + .filter_map(MetricType::from_id) + .map(|metric| (metric.as_str().replace(['.', '-'], "_"), metric)) + .collect() + }) + .get(name) + .copied() +} + +/// Rewrites bare underscore metric names into `metric('dotted.name')`, +/// leaving quoted text and function names alone (MON-25). +pub fn rewrite(text: &str) -> String { + let mut out = String::with_capacity(text.len()); + let chars = text.chars().collect::>(); + let mut i = 0; + while i < chars.len() { + let c = chars[i]; + if c == '"' || c == '\'' { + let quote = c; + out.push(c); + i += 1; + while i < chars.len() { + out.push(chars[i]); + if chars[i] == quote { + break; + } + i += 1; + } + i += 1; + } else if c.is_ascii_alphabetic() || c == '_' { + let start = i; + while i < chars.len() && is_ident_char(chars[i]) { + i += 1; + } + let word = chars[start..i].iter().collect::(); + let is_call = chars[i..].iter().find(|c| !c.is_whitespace()) == Some(&'('); + match underscore_metric(&word) { + Some(metric) if !is_call => { + out.push_str(&format!("metric('{}')", metric.as_str())); + } + _ => out.push_str(&word), + } + } else { + out.push(c); + i += 1; + } + } + out +} + +/// An alert condition with underscore names rewritten. +pub fn rewrite_condition(condition: &Expression) -> Expression { + Expression { + match_: List::from_iter(condition.match_.iter().map(|m| ExpressionMatch { + if_: rewrite(&m.if_), + then: rewrite(&m.then), + })), + else_: rewrite(&condition.else_), + } +} + +/// `%{metric.name}%` replaced by the metric's value: whole numbers without +/// decimals, others with at most two (MON-27). Unknown names stay. +pub fn render(template: &str) -> String { + let mut out = String::with_capacity(template.len()); + let mut rest = template; + while let Some(start) = rest.find("%{") { + out.push_str(&rest[..start]); + let after = &rest[start + 2..]; + match after.find("}%") { + Some(end) => { + let name = &after[..end]; + match MetricType::parse(name) { + Some(metric) => { + let value = Collector::read_metric(metric); + if value.fract() == 0.0 { + out.push_str(&format!("{}", value as i64)); + } else { + let text = format!("{value:.2}"); + out.push_str(text.trim_end_matches('0').trim_end_matches('.')); + } + } + None => out.push_str(&rest[start..start + 2 + end + 2]), + } + rest = &after[end + 2..]; + } + None => { + out.push_str(&rest[start..]); + rest = ""; + } + } + } + out.push_str(rest); + out +} + +fn build_email(email: ®istry::schema::structs::AlertEmailProperties) -> AlertMessage { + let from = match &email.from_name { + Some(name) => Address::new_address(Some(name.clone()), email.from_address.clone()), + None => Address::new_address(None::, email.from_address.clone()), + }; + let to = email.to.iter().cloned().collect::>(); + let body = MessageBuilder::new() + .from(from) + .to(to + .iter() + .map(|addr| Address::new_address(None::, addr.clone())) + .collect::>()) + .subject(render(&email.subject)) + .header("Auto-Submitted", HeaderType::Text("auto-generated".into())) + .text_body(render(&email.body)) + .write_to_vec() + .unwrap_or_default(); + AlertMessage { + from: email.from_address.clone(), + to, + body, + } +} + +impl Server { + /// Evaluates every enabled alert once (MON-25, MON-26), emits the event + /// of each that fires (MON-28), and returns the emails to queue + /// (MON-29). A failing alert is logged and skipped (MON-37). + pub async fn process_alerts(&self) -> trc::Result> { + let registry = self.registry(); + let ids = registry + .query::>(RegistryQuery::new(ObjectType::Alert)) + .await?; + let mut messages = Vec::new(); + let mut holding = AHashSet::new(); + for id in ids { + let Some(alert) = registry.object::(id).await? else { + continue; + }; + if !alert.enable { + continue; + } + let condition = rewrite_condition(&alert.condition); + let mut bp = Bootstrap::new_uninitialized(registry.clone()); + let if_block = bp.compile_expr( + ObjectId::new(ObjectType::Alert, id), + &ExpressionContext { + expr: &condition, + ..alert.ctx_condition() + }, + ); + if !bp.errors.is_empty() || if_block.is_empty() { + trc::event!( + Registry(trc::RegistryEvent::BuildWarning), + Id = id.id(), + Details = "The alert's condition can't be evaluated", + ); + continue; + } + let holds = self + .eval_if::(&if_block, &EmptyResolver, 0) + .await + .unwrap_or(false); + if !holds { + continue; + } + holding.insert(id.id()); + let was_firing = FIRING + .lock() + .unwrap() + .as_ref() + .is_some_and(|firing| firing.contains(&id.id())); + if was_firing { + continue; + } + + if let AlertEvent::Enabled(event) = &alert.event_alert { + trc::event!( + Telemetry(TelemetryEvent::AlertEvent), + Id = id.id(), + Details = render(event.event_message.as_deref().unwrap_or("Alert triggered")), + ); + } + if let AlertEmail::Enabled(email) = &alert.email_alert { + messages.push(build_email(email)); + } + } + *FIRING.lock().unwrap() = Some(holding); + Ok(messages) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn rewrites_underscore_names() { + assert_eq!(rewrite("domain_count > 1"), "metric('domain.count') > 1"); + assert_eq!( + rewrite("metric('queue.count') > 5 && queue_count < 9"), + "metric('queue.count') > 5 && metric('queue.count') < 9" + ); + assert_eq!(rewrite("'domain_count' == x"), "'domain_count' == x"); + assert_eq!(rewrite("unknown_thing > 1"), "unknown_thing > 1"); + } + + #[test] + fn renders_placeholders() { + assert_eq!(render("no placeholders"), "no placeholders"); + assert_eq!(render("%{no.such-metric}% left"), "%{no.such-metric}% left"); + } +} diff --git a/crates/common/src/telemetry/mod.rs b/crates/common/src/telemetry/mod.rs index 1f96589..bc49e28 100644 --- a/crates/common/src/telemetry/mod.rs +++ b/crates/common/src/telemetry/mod.rs @@ -4,6 +4,7 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ +pub mod alerts; // inbuxa: monitoring (MON-25 to MON-30) pub mod metrics; pub mod tracers; pub mod webhooks; diff --git a/crates/jmap/src/registry/set.rs b/crates/jmap/src/registry/set.rs index fabbb0a..8dce027 100644 --- a/crates/jmap/src/registry/set.rs +++ b/crates/jmap/src/registry/set.rs @@ -606,7 +606,23 @@ impl RegistrySet for Server { } // Validate expressions - if let Some(expressions) = new_object.inner.expression_ctxs() { + // inbuxa: MON-25: an alert condition may name metrics with underscores + let alert_condition = match &new_object.inner { + ObjectInner::Alert(alert) => Some( + common::telemetry::alerts::rewrite_condition(&alert.condition), + ), + _ => None, + }; + let expressions = match (&new_object.inner, &alert_condition) { + (ObjectInner::Alert(alert), Some(condition)) => { + Some(vec![registry::schema::prelude::ExpressionContext { + expr: condition, + ..alert.ctx_condition() + }]) + } + _ => new_object.inner.expression_ctxs(), + }; + if let Some(expressions) = expressions { let mut bp = Bootstrap::new_uninitialized(self.registry().clone()); for expression in expressions { diff --git a/crates/services/src/task_manager/scheduler.rs b/crates/services/src/task_manager/scheduler.rs index dd2637a..0116a92 100644 --- a/crates/services/src/task_manager/scheduler.rs +++ b/crates/services/src/task_manager/scheduler.rs @@ -42,6 +42,8 @@ enum Event { RenewNodeIdLease, // inbuxa: MON-4: metric history StoreMetrics, + // inbuxa: MON-25: alert evaluation + EvaluateAlerts, } /// When the next metric-history tick is due (MON-4), read from the registry @@ -126,6 +128,9 @@ pub fn spawn_task_scheduler(inner: Arc) { // Calculate expensive metrics queue.schedule(Instant::now(), Event::CalculateMetrics); + // inbuxa: MON-25: alerts every minute + queue.schedule(Instant::now() + Duration::from_secs(60), Event::EvaluateAlerts); + // inbuxa: MON-4: metric history on its own schedule queue.schedule( Instant::now() + metrics_collection_delay(&server).await, @@ -325,6 +330,43 @@ pub fn spawn_task_scheduler(inner: Arc) { ); }); } + // inbuxa: MON-25 to MON-30: on the node that calculates metrics; + // a failure is logged and tried again next minute (MON-37) + Event::EvaluateAlerts => { + queue.schedule( + Instant::now() + Duration::from_secs(60), + Event::EvaluateAlerts, + ); + if server.core.network.roles.metrics_calculate { + let server = server.clone(); + tokio::spawn(async move { + match server.process_alerts().await { + Ok(messages) => { + use smtp::reporting::send::MtaReportSend; + for message in messages { + server + .send_autogenerated( + message.from.clone(), + message.to.iter(), + message.body, + None, + 0, + ) + .await; + trc::event!( + Telemetry(TelemetryEvent::AlertMessage), + From = message.from, + To = message.to, + ); + } + } + Err(err) => { + trc::error!(err.details("Failed to evaluate alerts")); + } + } + }); + } + } // inbuxa: MON-4: every node writes its own samples Event::StoreMetrics => { queue.schedule( @@ -431,6 +473,7 @@ impl Event { Event::TrainSpamClassifier => "trainSpamClassifier", Event::RenewNodeIdLease => "renewNodeIdLease", Event::StoreMetrics => "storeMetrics", + Event::EvaluateAlerts => "evaluateAlerts", } } } diff --git a/tests/src/telemetry/mod.rs b/tests/src/telemetry/mod.rs index bf5b97c..86201fe 100644 --- a/tests/src/telemetry/mod.rs +++ b/tests/src/telemetry/mod.rs @@ -4,7 +4,6 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -#[cfg(feature = "pending-rebuild")] // inbuxa: pending-rebuild, see docs/spec/features/ pub mod alerts; pub mod metrics; pub mod tracing; @@ -56,7 +55,6 @@ pub async fn telemetry_tests() { .await; test.insert_account(admin); - #[cfg(feature = "pending-rebuild")] alerts::test(&test).await; metrics::test(&test).await; tracing::test(&test).await;