use crate::types::Level;
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct FieldPredicate {
pub path: String,
pub value: serde_json::Value,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct TaskEventsRuleFilter {
pub target: Option<String>,
pub min_level: Option<Level>,
pub field_filters: Vec<FieldPredicate>,
}
impl FieldPredicate {
pub fn matches(&self, fields: &serde_json::Value) -> bool {
let mut current = fields;
for key in self.path.split('.') {
match current.get(key) {
Some(v) => current = v,
None => return false,
}
}
current == &self.value
}
}
impl TaskEventsRuleFilter {
pub fn matches_fields(&self, fields: &serde_json::Value) -> bool {
self.field_filters.iter().all(|fp| fp.matches(fields))
}
pub fn parse(expr: &str) -> Result<Self, ParseError> {
let mut tokens = expr.split_ascii_whitespace().peekable();
let mut target: Option<String> = None;
let mut min_level: Option<Level> = None;
let mut field_filters: Vec<FieldPredicate> = Vec::new();
if let Some(first) = tokens.peek() {
if !first.contains('=') {
target = Some(tokens.next().unwrap().to_owned());
}
}
for token in tokens {
let (key, val_str) = token.split_once('=').ok_or_else(|| {
ParseError(format!("expected key=value, got {token:?}"))
})?;
match key {
"level" | "min_level" => {
min_level = Some(parse_level(val_str)?);
}
"" => return Err(ParseError("empty key before '='".into())),
_ => {
field_filters.push(FieldPredicate {
path: key.to_owned(),
value: parse_scalar(val_str),
});
}
}
}
Ok(Self { target, min_level, field_filters })
}
pub fn parse_check_field(check: &str) -> Result<Self, ParseError> {
let inner = check.strip_prefix("task-events:").ok_or_else(|| {
ParseError(format!("check field does not start with 'task-events:': {check:?}"))
})?;
Self::parse(inner)
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ScopeSpec {
CurrentTaskRun,
Service(String),
}
impl Default for ScopeSpec {
fn default() -> Self {
ScopeSpec::CurrentTaskRun
}
}
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct EventsRuleFilter {
#[serde(default)]
pub scope: ScopeSpec,
pub target: Option<String>,
pub min_level: Option<Level>,
pub field_filters: Vec<FieldPredicate>,
}
impl EventsRuleFilter {
pub fn matches_fields(&self, fields: &serde_json::Value) -> bool {
self.field_filters.iter().all(|fp| fp.matches(fields))
}
pub fn parse_check_field(check: &str) -> Result<Self, ParseError> {
let inner = if let Some(rest) = check.strip_prefix("events:") {
rest
} else if let Some(rest) = check.strip_prefix("task-events:") {
rest
} else {
return Err(ParseError(format!(
"check field does not start with 'events:' or 'task-events:': {check:?}"
)));
};
Self::parse(inner)
}
fn parse(expr: &str) -> Result<Self, ParseError> {
let mut tokens = expr.split_ascii_whitespace().peekable();
let mut scope = ScopeSpec::CurrentTaskRun;
let mut target: Option<String> = None;
let mut min_level: Option<Level> = None;
let mut field_filters: Vec<FieldPredicate> = Vec::new();
if let Some(&first) = tokens.peek() {
if let Some(rest) = first.strip_prefix("scope=") {
tokens.next();
scope = parse_scope_spec(rest)?;
if let Some(&next) = tokens.peek() {
if !next.contains('=') {
target = Some(tokens.next().unwrap().to_owned());
}
}
} else if !first.contains('=') {
target = Some(tokens.next().unwrap().to_owned());
}
}
for token in tokens {
let (key, val_str) = token.split_once('=').ok_or_else(|| {
ParseError(format!("expected key=value, got {token:?}"))
})?;
match key {
"level" | "min_level" => {
min_level = Some(parse_level(val_str)?);
}
"" => return Err(ParseError("empty key before '='".into())),
_ => {
field_filters.push(FieldPredicate {
path: key.to_owned(),
value: parse_scalar(val_str),
});
}
}
}
Ok(Self { scope, target, min_level, field_filters })
}
}
fn parse_scope_spec(s: &str) -> Result<ScopeSpec, ParseError> {
if let Some(inner) = s.strip_prefix("taskrun(").and_then(|s| s.strip_suffix(')')) {
let _ = inner;
Ok(ScopeSpec::CurrentTaskRun)
} else if let Some(inner) = s.strip_prefix("service(").and_then(|s| s.strip_suffix(')')) {
if inner.is_empty() {
return Err(ParseError("scope=service() requires a mesh ident".into()));
}
Ok(ScopeSpec::Service(inner.to_owned()))
} else {
Err(ParseError(format!(
"unknown scope: {s:?}; expected taskrun(<uuid>) or service(<mesh-ident>)"
)))
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct ParseError(pub String);
impl std::fmt::Display for ParseError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}
impl std::error::Error for ParseError {}
fn parse_level(s: &str) -> Result<Level, ParseError> {
match s.to_ascii_lowercase().as_str() {
"trace" => Ok(Level::Trace),
"debug" => Ok(Level::Debug),
"info" => Ok(Level::Info),
"warn" | "warning" => Ok(Level::Warn),
"error" => Ok(Level::Error),
"fatal" => Ok(Level::Fatal),
other => Err(ParseError(format!("unknown level {other:?}"))),
}
}
fn parse_scalar(s: &str) -> serde_json::Value {
if s == "true" { return serde_json::Value::Bool(true); }
if s == "false" { return serde_json::Value::Bool(false); }
if let Ok(n) = s.parse::<i64>() { return serde_json::json!(n); }
if let Ok(n) = s.parse::<f64>() { return serde_json::json!(n); }
serde_json::Value::String(s.to_owned())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_target_only() {
let f = TaskEventsRuleFilter::parse("cargo::rustc").unwrap();
assert_eq!(f.target.as_deref(), Some("cargo::rustc"));
assert!(f.min_level.is_none());
assert!(f.field_filters.is_empty());
}
#[test]
fn parse_target_and_field() {
let f = TaskEventsRuleFilter::parse("cargo::rustc error.code=E0308").unwrap();
assert_eq!(f.target.as_deref(), Some("cargo::rustc"));
assert!(f.min_level.is_none());
assert_eq!(f.field_filters.len(), 1);
assert_eq!(f.field_filters[0].path, "error.code");
assert_eq!(f.field_filters[0].value, serde_json::Value::String("E0308".into()));
}
#[test]
fn parse_level_key() {
let f = TaskEventsRuleFilter::parse("clippy::warning level=error").unwrap();
assert_eq!(f.target.as_deref(), Some("clippy::warning"));
assert_eq!(f.min_level, Some(Level::Error));
assert!(f.field_filters.is_empty());
}
#[test]
fn parse_min_level_alias() {
let f = TaskEventsRuleFilter::parse("min_level=warn").unwrap();
assert!(f.target.is_none());
assert_eq!(f.min_level, Some(Level::Warn));
}
#[test]
fn parse_multiple_field_filters() {
let f = TaskEventsRuleFilter::parse(
"cargo::rustc error.code=E0308 file.path=src/lib.rs"
).unwrap();
assert_eq!(f.field_filters.len(), 2);
assert_eq!(f.field_filters[0].path, "error.code");
assert_eq!(f.field_filters[1].path, "file.path");
assert_eq!(
f.field_filters[1].value,
serde_json::Value::String("src/lib.rs".into())
);
}
#[test]
fn parse_boolean_value() {
let f = TaskEventsRuleFilter::parse("build.success=false").unwrap();
assert!(f.target.is_none());
assert_eq!(f.field_filters[0].value, serde_json::Value::Bool(false));
}
#[test]
fn parse_numeric_value() {
let f = TaskEventsRuleFilter::parse("cargo::rustc file.line=42").unwrap();
assert_eq!(f.field_filters[0].value, serde_json::json!(42i64));
}
#[test]
fn parse_empty_is_unconstrained() {
let f = TaskEventsRuleFilter::parse("").unwrap();
assert!(f.target.is_none());
assert!(f.min_level.is_none());
assert!(f.field_filters.is_empty());
}
#[test]
fn parse_check_field_prefix() {
let f = TaskEventsRuleFilter::parse_check_field(
"task-events:cargo::rustc error.code=E0308"
).unwrap();
assert_eq!(f.target.as_deref(), Some("cargo::rustc"));
}
#[test]
fn parse_check_field_wrong_prefix() {
assert!(TaskEventsRuleFilter::parse_check_field("ast:foo").is_err());
}
#[test]
fn parse_missing_eq_returns_err() {
assert!(TaskEventsRuleFilter::parse("cargo::rustc noequalssign").is_err());
}
#[test]
fn parse_level_warning_alias() {
let f = TaskEventsRuleFilter::parse("level=warning").unwrap();
assert_eq!(f.min_level, Some(Level::Warn));
}
#[test]
fn parse_unknown_level_returns_err() {
assert!(TaskEventsRuleFilter::parse("level=critical").is_err());
}
#[test]
fn field_predicate_matches_present_key() {
let fp = FieldPredicate {
path: "error.code".into(),
value: serde_json::Value::String("E0308".into()),
};
let fields = serde_json::json!({"error": {"code": "E0308"}});
assert!(fp.matches(&fields));
}
#[test]
fn field_predicate_rejects_wrong_value() {
let fp = FieldPredicate {
path: "error.code".into(),
value: serde_json::Value::String("E0308".into()),
};
let fields = serde_json::json!({"error": {"code": "E0309"}});
assert!(!fp.matches(&fields));
}
#[test]
fn field_predicate_rejects_missing_key() {
let fp = FieldPredicate {
path: "error.code".into(),
value: serde_json::Value::String("E0308".into()),
};
let fields = serde_json::json!({"error": {}});
assert!(!fp.matches(&fields));
}
#[test]
fn matches_fields_empty_predicates_always_true() {
let f = TaskEventsRuleFilter { target: None, min_level: None, field_filters: vec![] };
assert!(f.matches_fields(&serde_json::json!({})));
}
#[test]
fn matches_fields_all_must_match() {
let f = TaskEventsRuleFilter::parse(
"cargo::rustc error.code=E0308 file.path=src/lib.rs"
).unwrap();
let ok = serde_json::json!({"error": {"code": "E0308"}, "file": {"path": "src/lib.rs"}});
assert!(f.matches_fields(&ok));
let bad = serde_json::json!({"error": {"code": "E0308"}});
assert!(!f.matches_fields(&bad));
}
#[test]
fn matches_fields_boolean_predicate() {
let f = TaskEventsRuleFilter::parse("build.success=false").unwrap();
assert!(f.matches_fields(&serde_json::json!({"build": {"success": false}})));
assert!(!f.matches_fields(&serde_json::json!({"build": {"success": true}})));
}
#[test]
fn events_filter_default_scope_is_current_taskrun() {
let f = EventsRuleFilter::parse_check_field("events:cargo::rustc level=error").unwrap();
assert_eq!(f.scope, ScopeSpec::CurrentTaskRun);
assert_eq!(f.target.as_deref(), Some("cargo::rustc"));
assert_eq!(f.min_level, Some(Level::Error));
}
#[test]
fn events_filter_task_events_alias() {
let alias = EventsRuleFilter::parse_check_field("task-events:cargo::rustc level=error").unwrap();
let canonical = EventsRuleFilter::parse_check_field("events:cargo::rustc level=error").unwrap();
assert_eq!(alias, canonical);
}
#[test]
fn events_filter_service_scope() {
let f = EventsRuleFilter::parse_check_field(
"events:scope=service(noisetable-api.pdx) level=error"
).unwrap();
assert_eq!(f.scope, ScopeSpec::Service("noisetable-api.pdx".into()));
assert!(f.target.is_none());
assert_eq!(f.min_level, Some(Level::Error));
}
#[test]
fn events_filter_service_scope_with_target_and_fields() {
let f = EventsRuleFilter::parse_check_field(
"events:scope=service(api.prod) cargo::rustc error.code=E0308"
).unwrap();
assert_eq!(f.scope, ScopeSpec::Service("api.prod".into()));
assert_eq!(f.target.as_deref(), Some("cargo::rustc"));
assert_eq!(f.field_filters.len(), 1);
assert_eq!(f.field_filters[0].path, "error.code");
}
#[test]
fn events_filter_taskrun_scope_desugars_to_current() {
let f = EventsRuleFilter::parse_check_field(
"events:scope=taskrun(00000000-0000-0000-0000-000000000000)"
).unwrap();
assert_eq!(f.scope, ScopeSpec::CurrentTaskRun);
}
#[test]
fn events_filter_service_empty_ident_is_err() {
assert!(EventsRuleFilter::parse_check_field("events:scope=service()").is_err());
}
#[test]
fn events_filter_unknown_scope_is_err() {
assert!(EventsRuleFilter::parse_check_field("events:scope=forge(abc)").is_err());
}
#[test]
fn events_filter_wrong_prefix_is_err() {
assert!(EventsRuleFilter::parse_check_field("ast:foo").is_err());
}
#[test]
fn events_filter_no_scope_no_target() {
let f = EventsRuleFilter::parse_check_field("events:").unwrap();
assert_eq!(f.scope, ScopeSpec::CurrentTaskRun);
assert!(f.target.is_none());
assert!(f.min_level.is_none());
assert!(f.field_filters.is_empty());
}
#[test]
fn events_filter_matches_fields_service_scope() {
let f = EventsRuleFilter::parse_check_field(
"events:scope=service(api.prod) level=error"
).unwrap();
assert!(f.matches_fields(&serde_json::json!({})));
}
}