use std::collections::HashMap;
use std::sync::Arc;
use futures_util::Stream;
use reqwest::Method;
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use crate::client::{CallOptions, Client};
use crate::common::{string_enum, MessageResponse};
use crate::error::Result;
use crate::pagination::{auto_page, paginate, ListParams, Page, PageFetcher};
use crate::query::QueryBuilder;
use crate::resources::escape;
const PATH: &str = "/v2/alertRules";
string_enum! {
AlertSeverity {
INFO => "info",
WARNING => "warning",
CRITICAL => "critical",
}
}
string_enum! {
AlertDomain {
RTC => "rtc",
AGENT => "agent",
}
}
string_enum! {
AlertRuleStatus {
ACTIVE => "active",
PAUSED => "paused",
}
}
string_enum! {
AlertConditionOperator {
GREATER_THAN => ">",
LESS_THAN => "<",
EQUAL => "=",
NOT_EQUAL => "!=",
}
}
string_enum! {
AlertMatchType {
AT_LEAST_ONCE => "at_least_once",
ALL_THE_TIME => "all_the_time",
ON_AVERAGE => "on_average",
IN_TOTAL => "in_total",
LAST => "last",
}
}
string_enum! {
AlertReduceTo {
LAST => "last",
SUM => "sum",
AVG => "avg",
MIN => "min",
MAX => "max",
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AbsentDataAlert {
pub enabled: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub for_seconds: Option<u32>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AlertRuleFilter {
pub category: String,
pub condition: String,
pub value: Value,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AlertRuleGroupBy {
pub key: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub data_type: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub key_type: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AlertRuleHaving {
pub column_name: String,
pub operator: String,
pub value: f64,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AlertRuleQuery {
pub metric: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub metric_category: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub time_aggregation: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub space_aggregation: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub data_type: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub step_interval_seconds: Option<u32>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub filter: Vec<AlertRuleFilter>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub group_by: Vec<AlertRuleGroupBy>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub having: Vec<AlertRuleHaving>,
}
impl AlertRuleQuery {
pub fn metric(metric: impl Into<String>) -> Self {
Self {
metric: metric.into(),
..Default::default()
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AlertRuleCondition {
pub operator: AlertConditionOperator,
pub threshold: f64,
#[serde(skip_serializing_if = "Option::is_none")]
pub time_aggregation: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub space_aggregation: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub minimum_session_gate: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub unit: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub match_type: Option<AlertMatchType>,
#[serde(skip_serializing_if = "Option::is_none")]
pub require_min_points: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub required_num_points: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub frequency_seconds: Option<u32>,
}
impl AlertRuleCondition {
pub fn new(operator: AlertConditionOperator, threshold: f64) -> Self {
Self {
operator,
threshold,
time_aggregation: None,
space_aggregation: None,
minimum_session_gate: None,
unit: None,
match_type: None,
require_min_points: None,
required_num_points: None,
frequency_seconds: None,
}
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AlertRuleEvaluation {
#[serde(skip_serializing_if = "Option::is_none")]
pub frequency_seconds: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub window_seconds: Option<u32>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AlertRuleConfig {
pub query: AlertRuleQuery,
pub condition: AlertRuleCondition,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub labels: HashMap<String, String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub evaluation: Option<AlertRuleEvaluation>,
#[serde(skip_serializing_if = "Option::is_none")]
pub absent_data_alert: Option<AbsentDataAlert>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
pub additional_queries: HashMap<String, AlertRuleQuery>,
#[serde(skip_serializing_if = "Option::is_none")]
pub formula: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub reduce_to: Option<AlertReduceTo>,
}
impl AlertRuleConfig {
pub fn new(query: AlertRuleQuery, condition: AlertRuleCondition) -> Self {
Self {
query,
condition,
labels: HashMap::new(),
evaluation: None,
absent_data_alert: None,
additional_queries: HashMap::new(),
formula: None,
reduce_to: None,
}
}
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AlertRule {
pub id: Option<String>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub alert_rule_id: String,
pub name: Option<String>,
pub severity: Option<AlertSeverity>,
pub description: Option<String>,
pub summary: Option<String>,
pub domain: Option<AlertDomain>,
pub rule_config: Option<AlertRuleConfig>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub notification_channel_ids: Vec<String>,
pub status: Option<AlertRuleStatus>,
pub version: Option<i64>,
pub last_triggered: Option<String>,
pub created_at: Option<String>,
pub updated_at: Option<String>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct CreateAlertRuleParams {
pub name: String,
pub severity: AlertSeverity,
pub rule_config: AlertRuleConfig,
#[serde(skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub summary: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub enabled: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub domain: Option<AlertDomain>,
#[serde(skip_serializing_if = "Option::is_none")]
pub time_aggregation: Option<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub notification_channel_ids: Vec<String>,
}
impl CreateAlertRuleParams {
pub fn new(
name: impl Into<String>,
severity: AlertSeverity,
rule_config: AlertRuleConfig,
) -> Self {
Self {
name: name.into(),
severity,
rule_config,
description: None,
summary: None,
enabled: None,
domain: None,
time_aggregation: None,
notification_channel_ids: Vec::new(),
}
}
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct UpdateAlertRuleParams {
#[serde(skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub severity: Option<AlertSeverity>,
#[serde(skip_serializing_if = "Option::is_none")]
pub rule_config: Option<AlertRuleConfig>,
#[serde(skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub summary: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub enabled: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub domain: Option<AlertDomain>,
#[serde(skip_serializing_if = "Option::is_none")]
pub time_aggregation: Option<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub notification_channel_ids: Vec<String>,
}
#[derive(Debug, Clone, Default)]
pub struct ListAlertRulesParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub region: Option<String>,
pub status: Option<AlertRuleStatus>,
pub metric_category: Option<String>,
pub domain: Option<AlertDomain>,
}
impl ListAlertRulesParams {
fn pagination(&self) -> ListParams {
ListParams {
page: self.page,
per_page: self.per_page,
cursor: self.cursor.clone(),
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct AlertRulesResource<'a> {
client: &'a Client,
}
impl<'a> AlertRulesResource<'a> {
pub(crate) fn new(client: &'a Client) -> Self {
Self { client }
}
pub async fn create(&self, params: CreateAlertRuleParams) -> Result<AlertRule> {
self.client
.wrapped(Method::POST, PATH, "alertRule", CallOptions::json(¶ms)?)
.await
}
pub async fn list(&self, params: ListAlertRulesParams) -> Result<Page<AlertRule>> {
paginate(self.fetcher(¶ms), ¶ms.pagination(), "data", None).await
}
pub fn list_stream(
&self,
params: ListAlertRulesParams,
) -> impl Stream<Item = Result<AlertRule>> + Send {
auto_page(self.fetcher(¶ms), params.pagination(), "data", None)
}
pub async fn get(&self, alert_rule_id: &str) -> Result<AlertRule> {
let path = format!("{PATH}/{}", escape(alert_rule_id));
self.client
.wrapped(Method::GET, &path, "alertRule", CallOptions::new())
.await
}
pub async fn update(
&self,
alert_rule_id: &str,
params: UpdateAlertRuleParams,
) -> Result<AlertRule> {
let path = format!("{PATH}/{}", escape(alert_rule_id));
self.client
.wrapped(Method::PUT, &path, "alertRule", CallOptions::json(¶ms)?)
.await
}
pub async fn delete(&self, alert_rule_id: &str) -> Result<MessageResponse> {
let path = format!("{PATH}/{}", escape(alert_rule_id));
self.client
.json(Method::DELETE, &path, CallOptions::new())
.await
}
fn fetcher(&self, params: &ListAlertRulesParams) -> PageFetcher {
let client = self.client.clone();
let params = params.clone();
Arc::new(move |page, per_page| {
let client = client.clone();
let params = params.clone();
Box::pin(async move {
let query = QueryBuilder::new()
.opt("page", page)
.opt("perPage", per_page)
.opt_str("region", params.region.as_deref())
.opt_str(
"status",
params.status.as_ref().map(AlertRuleStatus::as_str),
)
.opt_str("metricCategory", params.metric_category.as_deref())
.opt_str("domain", params.domain.as_ref().map(AlertDomain::as_str))
.into_pairs();
client
.json::<Value>(Method::GET, PATH, CallOptions::new().query(query))
.await
})
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn create_serializes_the_typed_config() {
let params = CreateAlertRuleParams::new(
"High bitrate",
AlertSeverity::WARNING,
AlertRuleConfig::new(
AlertRuleQuery::metric("stat_consumer_video_bitrate"),
AlertRuleCondition::new(AlertConditionOperator::GREATER_THAN, 2500.0),
),
);
assert_eq!(
serde_json::to_value(¶ms).unwrap(),
json!({
"name": "High bitrate",
"severity": "warning",
"ruleConfig": {
"query": {"metric": "stat_consumer_video_bitrate"},
"condition": {"operator": ">", "threshold": 2500.0},
},
})
);
}
#[test]
fn an_unknown_severity_on_a_response_does_not_fail() {
let rule: AlertRule = serde_json::from_value(json!({
"alertRuleId": "ar-1",
"severity": "catastrophic",
"notificationChannelIds": null,
}))
.unwrap();
assert_eq!(rule.severity.unwrap().as_str(), "catastrophic");
assert!(rule.notification_channel_ids.is_empty());
}
}