use crate::prelude::*;
use crate::webterface::AuthBearer;
use std::time::{SystemTime, UNIX_EPOCH};
use matrix_sdk::ruma::events::MessageLikeEventContent;
use axum::{
extract::{Json, State},
http::StatusCode,
response::IntoResponse,
};
use grafana::{Alert, AlertStatus, Alerts};
static FIRING_ALERTS: LazyLock<FiringAlerts> = LazyLock::new(Default::default);
#[derive(Default)]
struct FiringAlerts {
inner: Arc<Mutex<HashMap<String, Vec<Alert>>>>,
}
impl FiringAlerts {
fn fire(&self, name: &str, alerts: Vec<Alert>) -> anyhow::Result<Vec<Alert>> {
trace!("gathering alerts to fire");
let mut inner = match self.inner.lock() {
Ok(i) => i,
Err(e) => bail!("failed locking alerts map: {e}"),
};
let mut changed: Vec<Alert> = vec![];
trace!("listing known alerts");
let known_alerts: Vec<String> = inner.get(name).map_or_else(
|| {
changed.extend(alerts.clone());
vec![]
},
|a| a.iter().map(|a| a.fingerprint.clone()).collect(),
);
trace!("adding unique firing alerts");
inner
.entry(name.to_owned())
.and_modify(|va| {
for a in alerts.clone() {
if !known_alerts.contains(&a.fingerprint) {
va.push(a.clone());
changed.push(a);
};
}
})
.or_insert(alerts);
drop(inner);
Ok(changed)
}
fn resolve(&self, name: &str, alerts: Vec<Alert>) -> anyhow::Result<Vec<Alert>> {
let mut inner = match self.inner.lock() {
Ok(i) => i,
Err(e) => bail!("failed locking alerts map: {e}"),
};
trace!("known instances: {:#?}", inner.keys());
let resolved_fingerprints: Vec<String> =
alerts.iter().map(|a| a.fingerprint.clone()).collect();
inner
.entry(name.to_owned())
.and_modify(|va| va.retain(|a| !resolved_fingerprints.contains(&a.fingerprint)));
drop(inner);
Ok(alerts)
}
fn get(&self, name: &str) -> Option<Vec<Alert>> {
let Ok(inner) = self.inner.lock() else {
return None;
};
inner.get(name).map(std::borrow::ToOwned::to_owned)
}
fn purge(&self) -> anyhow::Result<()> {
if let Ok(mut inner) = self.inner.lock() {
for instance in inner.values_mut() {
instance.truncate(0);
}
} else {
bail!("failed locking alerts map");
};
Ok(())
}
}
#[derive(Clone, Debug, Deserialize)]
pub struct GrafanaConfig {
pub name: String,
pub token: String,
pub rooms: Vec<String>,
}
#[derive(Clone, Debug, Deserialize)]
pub struct ModuleConfig {
pub grafanas: HashMap<String, GrafanaConfig>,
#[serde(default = "keywords_alerting")]
pub keywords_alerting: Vec<String>,
#[serde(default = "keywords_purge")]
pub keywords_purge: Vec<String>,
pub rooms_purge: Vec<String>,
#[serde(default = "no_firing_alerts_responses")]
pub no_firing_alerts_responses: Vec<String>,
}
fn keywords_alerting() -> Vec<String> {
vec!["alerting".s(), "alerts".s()]
}
fn keywords_purge() -> Vec<String> {
vec!["purge".s(), "alerts_purge".s()]
}
fn no_firing_alerts_responses() -> Vec<String> {
vec!["all systems operational".s()]
}
#[axum::debug_handler]
pub async fn receive_alerts(
State(app_state): State<WebAppState>,
AuthBearer(token): AuthBearer,
Json(alerts): Json<Alerts>,
) -> Result<impl IntoResponse, (StatusCode, &'static str)> {
use AlertStatus::{Firing, Resolved};
let module_config: ModuleConfig = {
match app_state.config.typed_module_config(module_path!()) {
Err(_) => return Err((StatusCode::INTERNAL_SERVER_ERROR, "no auth configuration")),
Ok(v) => v,
}
};
let mut maybe_instance: Option<String> = None;
for (name, config) in &module_config.grafanas {
if token == config.token {
maybe_instance = Some(name.to_owned());
break;
}
}
let Some(instance) = maybe_instance else {
return Err((StatusCode::FORBIDDEN, "unknown token"));
};
trace!("received hook body: {:#?}", alerts);
let changed = match alerts.status {
Firing => FIRING_ALERTS
.fire(&instance, alerts.alerts)
.map_err(|_| (StatusCode::INTERNAL_SERVER_ERROR, "failed to fire alerts")),
Resolved => FIRING_ALERTS
.resolve(&instance, alerts.alerts)
.map_err(|_| {
(
StatusCode::INTERNAL_SERVER_ERROR,
"failed to resolve alerts",
)
}),
};
trace!("{changed:#?}");
if let Ok(alerts) = changed {
if alerts.is_empty() {
return Ok(());
};
async {
if let Some(grafana_config) = module_config.grafanas.get(&instance) {
for room in grafana_config.rooms.clone() {
if let Ok(mx_room) = maybe_get_room(&app_state.mx, &room).await {
let mx_message = to_matrix_message(alerts.clone(), &instance);
if let Err(e) = mx_room.send(mx_message).await {
trace!("failed to send room notification: {e}");
}
}
}
};
Ok(())
}
.await
.map_err(|_: anyhow::Error| {
(
StatusCode::INTERNAL_SERVER_ERROR,
"failed to send room notifications",
)
})?;
};
Ok(())
}
pub(crate) fn starter(_: &Client, config: &Config) -> anyhow::Result<Vec<ModuleInfo>> {
info!("registering grafana modules");
let module_config: ModuleConfig = config.typed_module_config(module_path!())?;
let (alerting_tx, alerting_rx) = mpsc::channel::<ConsumerEvent>(1);
let alerting = ModuleInfo {
name: "alerting".s(),
help: "shows which alerts are now firing".s(),
acl: vec![],
trigger: TriggerType::Keyword(module_config.keywords_alerting.clone()),
channel: alerting_tx,
error_prefix: Some("error".s()),
};
alerting.spawn(alerting_rx, module_config.clone(), alerting_processor);
let (purge_tx, purge_rx) = mpsc::channel::<ConsumerEvent>(1);
let purge = ModuleInfo {
name: "alerts_purge".s(),
help: "reset the firing alerts to empty state".s(),
acl: vec![Acl::Room(module_config.rooms_purge.clone())],
trigger: TriggerType::Keyword(module_config.keywords_purge.clone()),
channel: purge_tx,
error_prefix: Some("error purging state".s()),
};
purge.spawn(purge_rx, module_config, purge_processor);
Ok(vec![alerting, purge])
}
pub async fn purge_processor(ev: ConsumerEvent, _: ModuleConfig) -> anyhow::Result<()> {
trace!("purging alerts");
let response = match FIRING_ALERTS.purge() {
Ok(()) => "alerts purged",
Err(e) => return Err(e),
};
ev.room
.send(RoomMessageEventContent::text_plain(response))
.await?;
Ok(())
}
pub async fn alerting_processor(event: ConsumerEvent, config: ModuleConfig) -> anyhow::Result<()> {
let mut grafanas: Vec<GrafanaConfig> = vec![];
let mut sent: bool = false;
if let Some(maybe_grafana_instances) = event.args {
trace!("maybe instances: {maybe_grafana_instances}");
let mut maybe_grafanas: Vec<String> = vec![];
let mut args = maybe_grafana_instances.split_whitespace();
let first = args
.next()
.ok_or_else(|| anyhow!("missing arguments"))?
.to_string();
maybe_grafanas.push(first);
for maybe_grafana in args {
maybe_grafanas.push(maybe_grafana.to_string());
}
for instance_name in maybe_grafanas {
if let Some(grafana) = config.grafanas.get(&instance_name) {
grafanas.push(grafana.clone());
} else {
bail!("provided grafana instance is not known: {instance_name}");
};
}
} else {
grafanas = config.grafanas.values().cloned().collect();
trace!("all instances: {grafanas:#?}");
}
trace!("grafanas to check: {grafanas:#?}");
for grafana in grafanas {
let name = grafana.name.as_str();
let alerts = FIRING_ALERTS.get(name);
match alerts {
None => {
trace!("no alerts known");
}
Some(va) => {
if va.is_empty() {
continue;
};
event.room.send(to_matrix_message(va, name)).await?;
sent = true;
}
};
}
if !sent {
let mut response = String::new();
config
.no_firing_alerts_responses
.first()
.ok_or_else(|| anyhow!("module misconfigured: missing `ok` responses"))?
.clone_into(&mut response);
if let Ok(now) = SystemTime::now().duration_since(UNIX_EPOCH) {
let milis = now.as_millis();
let chosen_idx: usize = milis as usize % config.no_firing_alerts_responses.len();
if let Some(option) = config.no_firing_alerts_responses.get(chosen_idx) {
option.clone_into(&mut response);
};
};
event
.room
.send(RoomMessageEventContent::text_plain(response))
.await?;
};
Ok(())
}
#[must_use]
pub fn to_matrix_message(va: Vec<grafana::Alert>, instance: &str) -> impl MessageLikeEventContent {
let mut response_html = format!("instance: <b>{instance}</b><br />");
let mut response = format!("instance: {instance}\n");
for alert in va {
let mut annotations_html = "".s();
for (key, value) in alert.annotations.clone() {
annotations_html.push_str(format!("{key}: <b>{value}</b><br/>").as_str());
}
response_html.push_str(
format!(
r"{state_emoji}<b>{state}</b><br/>
{annotations}
since: {since}<br />",
state_emoji = alert.status.clone().into_emoji(),
state = alert.status,
annotations = annotations_html,
since = alert.starts_at,
)
.as_str(),
);
let mut annotations = "".s();
for (key, value) in alert.annotations {
annotations.push_str(format!("{key}: {value}\n").as_str());
}
response.push_str(
format!(
"{state_emoji} {state}\n
{annotations}since: {since}\n",
state_emoji = alert.status.clone().into_emoji(),
state = alert.status,
annotations = annotations,
since = alert.starts_at,
)
.as_str(),
);
}
RoomMessageEventContent::text_html(response, response_html)
}
pub mod grafana {
use serde_derive::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::HashMap;
use std::fmt;
#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, Default)]
pub enum AlertStatus {
#[serde(rename = "resolved")]
Resolved,
#[serde(rename = "firing")]
#[default]
Firing,
}
impl fmt::Display for AlertStatus {
fn fmt(&self, fmt: &mut fmt::Formatter) -> fmt::Result {
use AlertStatus::{Firing, Resolved};
match self {
Resolved => write!(fmt, "Resolved"),
Firing => write!(fmt, "Firing"),
}
}
}
impl AlertStatus {
pub(crate) const fn into_emoji(self) -> &'static str {
use AlertStatus::{Firing, Resolved};
match self {
Firing => "🔥",
Resolved => "🩷",
}
}
}
#[allow(dead_code, missing_docs)]
#[derive(Debug, Clone, Deserialize, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct Alerts {
pub receiver: String,
pub status: AlertStatus,
pub org_id: i64,
pub alerts: Vec<Alert>,
pub group_labels: HashMap<String, String>,
pub common_labels: HashMap<String, String>,
pub common_annotations: HashMap<String, String>,
#[serde(rename = "externalURL")]
pub external_url: String,
pub version: String,
pub group_key: String,
pub truncated_alerts: i64,
pub title: String,
pub state: String,
pub message: String,
}
#[allow(dead_code, missing_docs)]
#[derive(Debug, Clone, Deserialize, Serialize, Default)]
#[serde(rename_all = "camelCase")]
pub struct Alert {
pub status: AlertStatus,
pub labels: HashMap<String, String>,
pub annotations: HashMap<String, String>,
pub starts_at: String,
pub ends_at: String,
#[serde(rename = "generatorURL")]
pub generator_url: String,
pub fingerprint: String,
#[serde(rename = "silenceURL")]
pub silence_url: String,
#[serde(rename = "dashboardURL")]
pub dashboard_url: String,
#[serde(rename = "panelURL")]
pub panel_url: String,
pub values: Value,
}
}