diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 8cc05b9d55..838ffcea87 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -04dd9c5c352f04995cd0470400a877261f956561 +78859aab0c6e78283ec8d2b37e8c410963afdc83 diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 00fd6f4d69..a0efeaaa76 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -29112,6 +29112,56 @@ components: - interval_secs - message + TriggerFilter: + description: > + Either a leaf filter, matching a field of the message (parsed as JSON) against a + value by equality (or superset, when the value is an object or array) — addressed + by `key` for a top-level field or `path` for a dotted path into nested objects — + or a group nesting sub-filters under a boolean operator (`none_of` matches when + none of its sub-filters do). + oneOf: + - type: object + properties: + key: + type: string + value: {} + required: + - key + - value + - type: object + properties: + path: + type: string + description: Dotted path into nested objects, e.g. `a.b.c`. Does not traverse arrays. + value: {} + required: + - path + - value + - type: object + properties: + any_of: + type: array + items: + $ref: "#/components/schemas/TriggerFilter" + required: + - any_of + - type: object + properties: + all_of: + type: array + items: + $ref: "#/components/schemas/TriggerFilter" + required: + - all_of + - type: object + properties: + none_of: + type: array + items: + $ref: "#/components/schemas/TriggerFilter" + required: + - none_of + WebsocketTrigger: allOf: - $ref: "#/components/schemas/TriggerExtraProperty" @@ -29132,23 +29182,16 @@ components: description: Last error message if the trigger failed filters: type: array - description: Array of key-value filters to match incoming messages (only matching messages trigger the script) + description: "Filters to match incoming messages (only matching messages trigger the script). Each entry is either a leaf `{key, value}` (top-level field) or `{path, value}` (dotted path into nested objects), or a group `{any_of: [...]}` / `{all_of: [...]}` / `{none_of: [...]}` nesting more entries. Entries at the top level are combined with `filter_logic`." items: - type: object - properties: - key: - type: string - value: {} - required: - - key - - value + $ref: "#/components/schemas/TriggerFilter" filter_logic: type: string enum: - and - or default: and - description: "Logic to apply when evaluating filters. 'and' requires all filters to match, 'or' requires any filter to match." + description: "Logic to apply when evaluating the top-level filters. 'and' requires all of them to match, 'or' requires any of them to match. Nested `any_of`/`all_of`/`none_of` groups carry their own logic." initial_messages: type: array nullable: true @@ -29204,23 +29247,16 @@ components: $ref: "#/components/schemas/TriggerMode" filters: type: array - description: Array of key-value filters to match incoming messages (only matching messages trigger the script) + description: "Filters to match incoming messages (only matching messages trigger the script). Each entry is either a leaf `{key, value}` (top-level field) or `{path, value}` (dotted path into nested objects), or a group `{any_of: [...]}` / `{all_of: [...]}` / `{none_of: [...]}` nesting more entries. Entries at the top level are combined with `filter_logic`." items: - type: object - properties: - key: - type: string - value: {} - required: - - key - - value + $ref: "#/components/schemas/TriggerFilter" filter_logic: type: string enum: - and - or default: and - description: "Logic to apply when evaluating filters. 'and' requires all filters to match, 'or' requires any filter to match." + description: "Logic to apply when evaluating the top-level filters. 'and' requires all of them to match, 'or' requires any of them to match. Nested `any_of`/`all_of`/`none_of` groups carry their own logic." initial_messages: type: array nullable: true @@ -29287,23 +29323,16 @@ components: description: True if script_path points to a flow, false if it points to a script filters: type: array - description: Array of key-value filters to match incoming messages (only matching messages trigger the script) + description: "Filters to match incoming messages (only matching messages trigger the script). Each entry is either a leaf `{key, value}` (top-level field) or `{path, value}` (dotted path into nested objects), or a group `{any_of: [...]}` / `{all_of: [...]}` / `{none_of: [...]}` nesting more entries. Entries at the top level are combined with `filter_logic`." items: - type: object - properties: - key: - type: string - value: {} - required: - - key - - value + $ref: "#/components/schemas/TriggerFilter" filter_logic: type: string enum: - and - or default: and - description: "Logic to apply when evaluating filters. 'and' requires all filters to match, 'or' requires any filter to match." + description: "Logic to apply when evaluating the top-level filters. 'and' requires all of them to match, 'or' requires any of them to match. Nested `any_of`/`all_of`/`none_of` groups carry their own logic." initial_messages: type: array nullable: true @@ -30559,22 +30588,16 @@ components: description: Array of Kafka topic names to subscribe to filters: type: array + description: "Filters to match incoming messages (only matching messages trigger the script). Each entry is either a leaf `{key, value}` (top-level field) or `{path, value}` (dotted path into nested objects), or a group `{any_of: [...]}` / `{all_of: [...]}` / `{none_of: [...]}` nesting more entries. Entries at the top level are combined with `filter_logic`." items: - type: object - properties: - key: - type: string - value: {} - required: - - key - - value + $ref: "#/components/schemas/TriggerFilter" filter_logic: type: string enum: - and - or default: and - description: "Logic to apply when evaluating filters. 'and' requires all filters to match, 'or' requires any filter to match." + description: "Logic to apply when evaluating the top-level filters. 'and' requires all of them to match, 'or' requires any of them to match. Nested `any_of`/`all_of`/`none_of` groups carry their own logic." auto_offset_reset: type: string enum: @@ -30637,22 +30660,16 @@ components: description: Array of Kafka topic names to subscribe to filters: type: array + description: "Filters to match incoming messages (only matching messages trigger the script). Each entry is either a leaf `{key, value}` (top-level field) or `{path, value}` (dotted path into nested objects), or a group `{any_of: [...]}` / `{all_of: [...]}` / `{none_of: [...]}` nesting more entries. Entries at the top level are combined with `filter_logic`." items: - type: object - properties: - key: - type: string - value: {} - required: - - key - - value + $ref: "#/components/schemas/TriggerFilter" filter_logic: type: string enum: - and - or default: and - description: "Logic to apply when evaluating filters. 'and' requires all filters to match, 'or' requires any filter to match." + description: "Logic to apply when evaluating the top-level filters. 'and' requires all of them to match, 'or' requires any of them to match. Nested `any_of`/`all_of`/`none_of` groups carry their own logic." auto_offset_reset: type: string enum: @@ -30711,22 +30728,16 @@ components: description: Array of Kafka topic names to subscribe to filters: type: array + description: "Filters to match incoming messages (only matching messages trigger the script). Each entry is either a leaf `{key, value}` (top-level field) or `{path, value}` (dotted path into nested objects), or a group `{any_of: [...]}` / `{all_of: [...]}` / `{none_of: [...]}` nesting more entries. Entries at the top level are combined with `filter_logic`." items: - type: object - properties: - key: - type: string - value: {} - required: - - key - - value + $ref: "#/components/schemas/TriggerFilter" filter_logic: type: string enum: - and - or default: and - description: "Logic to apply when evaluating filters. 'and' requires all filters to match, 'or' requires any filter to match." + description: "Logic to apply when evaluating the top-level filters. 'and' requires all of them to match, 'or' requires any of them to match. Nested `any_of`/`all_of`/`none_of` groups carry their own logic." auto_offset_reset: type: string enum: diff --git a/backend/windmill-trigger-websocket/src/handler.rs b/backend/windmill-trigger-websocket/src/handler.rs index c2ca37c500..a428bc7f87 100644 --- a/backend/windmill-trigger-websocket/src/handler.rs +++ b/backend/windmill-trigger-websocket/src/handler.rs @@ -12,7 +12,7 @@ use windmill_common::{ worker::to_raw_value, }; use windmill_git_sync::DeployedObject; -use windmill_trigger::{Trigger, TriggerCrud, TriggerData}; +use windmill_trigger::{filter::CompiledFilters, Trigger, TriggerCrud, TriggerData}; use super::{ get_url_from_runnable_value, listener::InitialMessage, proxy::connect_async_with_proxy, @@ -106,6 +106,8 @@ impl TriggerCrud for WebsocketTrigger { } } + CompiledFilters::validate(&config.filters)?; + if let Some(ref hb) = config.heartbeat { if hb.interval_secs < 1 { return Err(Error::BadRequest( diff --git a/backend/windmill-trigger-websocket/src/listener.rs b/backend/windmill-trigger-websocket/src/listener.rs index a38f2434a4..91bbc4bf23 100644 --- a/backend/windmill-trigger-websocket/src/listener.rs +++ b/backend/windmill-trigger-websocket/src/listener.rs @@ -21,7 +21,7 @@ use windmill_common::{ DB, }; use windmill_queue::PushArgsOwned; -use windmill_trigger::filter::{check_filters, Filter}; +use windmill_trigger::filter::CompiledFilters; use windmill_trigger::listener::{update_rw_lock, ListeningTrigger}; use windmill_trigger::trigger_helpers::{ trigger_runnable, trigger_runnable_and_wait_for_raw_result, @@ -362,15 +362,14 @@ impl Listener for WebsocketTrigger { } => {}, // Message reader _ = async { - let filters: Vec = if listening_trigger.trigger_mode { - listening_trigger - .trigger_config - .filters - .iter() - .filter_map(|m| serde_json::from_str(m.get()).ok()) - .collect_vec() + let filters = if listening_trigger.trigger_mode { + CompiledFilters::parse( + listening_trigger.trigger_config.filters.iter().map(|m| m.get()), + listening_trigger.trigger_config.filter_logic == "or", + &listening_trigger.path, + ) } else { - vec![] + CompiledFilters::default() }; loop { if let Some(msg) = reader.next().await { @@ -391,9 +390,7 @@ impl Listener for WebsocketTrigger { } } - let use_or = listening_trigger.trigger_config.filter_logic == "or"; - let should_handle = check_filters(&text, &filters, use_or); - if should_handle { + if filters.matches(&text) { let trigger_info = HashMap::from([ ("url".to_string(), to_raw_value(&listening_trigger.trigger_config.url)), ]); diff --git a/backend/windmill-trigger/src/filter.rs b/backend/windmill-trigger/src/filter.rs index 3c9c857058..3c1ab7db76 100644 --- a/backend/windmill-trigger/src/filter.rs +++ b/backend/windmill-trigger/src/filter.rs @@ -2,52 +2,301 @@ use serde::{ de::{self, MapAccess, Visitor}, Deserialize, Deserializer, }; -use serde_json::Value; -use std::fmt; +use serde_json::{value::RawValue, Value}; +use std::{collections::HashMap, fmt}; -#[derive(Deserialize)] +#[derive(Debug, Deserialize)] pub struct JsonFilter { pub key: String, pub value: Value, } -#[derive(Deserialize)] +/// Same comparison as [`JsonFilter`], but the field is addressed by a dotted path into +/// nested objects. A separate field rather than dots in `key`, because a `key` containing +/// a dot already means the top-level field spelled that way. +#[derive(Debug, Deserialize)] +pub struct PathFilter { + pub path: String, + pub value: Value, +} + +/// Boolean group of nested filters, externally tagged (`{"any_of": [...]}`) so it is +/// unambiguous against a leaf filter, which is `{"key": ..., "value": ...}`. +#[derive(Debug, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum FilterGroup { + AnyOf(Vec), + AllOf(Vec), + NoneOf(Vec), +} + +#[derive(Debug, Deserialize)] #[serde(untagged)] pub enum Filter { JsonFilter(JsonFilter), + PathFilter(PathFilter), + Group(FilterGroup), } -struct SupersetVisitor<'a> { - key: &'a str, - value_to_check: &'a Value, +/// The scanned top-level key, and whatever is left to walk inside its value. +fn split_path(path: &str) -> (&str, &str) { + path.split_once('.').unwrap_or((path, "")) } -impl<'de, 'a> Visitor<'de> for SupersetVisitor<'a> { - type Value = bool; +/// Objects only: `Value::get` yields nothing for a string index into an array, so a path +/// through one simply does not match rather than guessing an element. +fn resolve<'v>(root: &'v Value, rest: &str) -> Option<&'v Value> { + let mut current = root; + if !rest.is_empty() { + for segment in rest.split('.') { + current = current.get(segment)?; + } + } + Some(current) +} + +/// Filters prepared for repeated evaluation against a stream of messages. The set of +/// top-level keys the whole tree references is computed once, so each message is scanned +/// in a single pass instead of once per leaf filter. +#[derive(Debug, Default)] +pub struct CompiledFilters { + filters: Vec, + use_or_logic: bool, + keys: Vec, +} + +impl CompiledFilters { + pub fn new(filters: Vec, use_or_logic: bool) -> Self { + let filters = drop_empty_groups(filters); + let mut keys = Vec::new(); + collect_keys(&filters, &mut keys); + Self { filters, use_or_logic, keys } + } + + /// Build from the raw JSON of each filter as stored in the trigger config. An entry + /// that fails to parse is skipped rather than dropping the other filters, but it + /// widens what the trigger accepts, so it is reported. + pub fn parse<'a>( + raw_filters: impl IntoIterator, + use_or_logic: bool, + trigger_path: &str, + ) -> Self { + let filters = raw_filters + .into_iter() + .filter_map(|raw| match serde_json::from_str::(raw) { + Ok(filter) => Some(filter), + Err(err) => { + tracing::error!( + "Ignoring unparseable filter of trigger {}: {} ({})", + trigger_path, + raw, + err + ); + None + } + }) + .collect(); + Self::new(filters, use_or_logic) + } + + pub fn is_empty(&self) -> bool { + self.filters.is_empty() + } + + /// Reject at save time what [`Self::parse`] would drop at listen time. A group nests + /// arbitrarily many criteria, so one mistyped entry silently widens the trigger by the + /// whole subtree it belongs to. + pub fn validate(filters: &[Value]) -> windmill_common::error::Result<()> { + for (index, filter) in filters.iter().enumerate() { + validate_filter(filter, &format!("filter #{}", index + 1))?; + } + Ok(()) + } + + /// Whether `text`, parsed as a JSON object, satisfies the filters. + pub fn matches(&self, text: &str) -> bool { + if self.filters.is_empty() { + return true; + } + + let mut deserializer = serde_json::Deserializer::from_str(text); + let values = + Deserializer::deserialize_map(&mut deserializer, KeysVisitor { keys: &self.keys }) + .unwrap_or_default(); + + eval_all(&self.filters, self.use_or_logic, &values) + } +} + +/// Groups are descended into by hand so a bad entry is named on its own: serde's untagged +/// error only reports that the outermost entry matched no variant, whatever depth is wrong. +fn validate_filter(filter: &Value, path: &str) -> windmill_common::error::Result<()> { + const GROUP_KEYS: [&str; 3] = ["any_of", "all_of", "none_of"]; + + let group = filter + .as_object() + .filter(|object| object.len() == 1) + .and_then(|object| { + GROUP_KEYS + .into_iter() + .find_map(|key| object.get(key).map(|nested| (key, nested))) + }); + + if let Some((key, nested)) = group { + let nested = nested.as_array().ok_or_else(|| { + windmill_common::error::Error::BadRequest(format!( + "{}: {} must be an array of filters", + path, key + )) + })?; + for (index, child) in nested.iter().enumerate() { + validate_filter(child, &format!("{} -> {}[{}]", path, key, index))?; + } + return Ok(()); + } + + // Everything below is meant to be a leaf. The untagged enum resolves a half-and-half + // entry by taking the first variant that fits and ignoring the rest of it, so a + // criterion carrying a group key would silently lose the whole subtree. + if let Some(group_key) = GROUP_KEYS.iter().find(|key| filter.get(*key).is_some()) { + return Err(windmill_common::error::Error::BadRequest(format!( + "{} combines a criterion with a {} group; an entry is one or the other", + path, group_key + ))); + } + + if filter.get("key").is_some() && filter.get("path").is_some() { + return Err(windmill_common::error::Error::BadRequest(format!( + "{} names its field with both key and path; use one or the other", + path + ))); + } + + let parsed = serde_json::from_value::(filter.clone()).map_err(|err| { + windmill_common::error::Error::BadRequest(format!( + "{} is neither a {{key, value}} / {{path, value}} criterion nor an any_of/all_of/none_of group: {}", + path, err + )) + })?; + + // An empty segment addresses no field, so the filter could only ever reject everything + if let Filter::PathFilter(PathFilter { path: dotted, .. }) = &parsed { + if dotted.split('.').any(|segment| segment.is_empty()) { + return Err(windmill_common::error::Error::BadRequest(format!( + "{}: path {:?} has an empty segment", + path, dotted + ))); + } + } + + Ok(()) +} + +/// A group with no criterion cannot evaluate to a constant: `true` makes an `or` list +/// accept every message, `false` mutes an `and` list. Dropping it instead leaves its +/// siblings in force, which is what a group left empty in the editor should mean. +fn drop_empty_groups(filters: Vec) -> Vec { + filters + .into_iter() + .filter_map(|filter| { + let (rebuild, nested): (fn(Vec) -> FilterGroup, _) = match filter { + Filter::Group(FilterGroup::AnyOf(nested)) => (FilterGroup::AnyOf, nested), + Filter::Group(FilterGroup::AllOf(nested)) => (FilterGroup::AllOf, nested), + Filter::Group(FilterGroup::NoneOf(nested)) => (FilterGroup::NoneOf, nested), + leaf => return Some(leaf), + }; + let nested = drop_empty_groups(nested); + (!nested.is_empty()).then(|| Filter::Group(rebuild(nested))) + }) + .collect() +} + +fn push_key(keys: &mut Vec, key: &str) { + if !keys.iter().any(|k| k == key) { + keys.push(key.to_string()); + } +} + +fn collect_keys(filters: &[Filter], keys: &mut Vec) { + for filter in filters { + match filter { + Filter::JsonFilter(JsonFilter { key, .. }) => push_key(keys, key), + Filter::PathFilter(PathFilter { path, .. }) => push_key(keys, split_path(path).0), + Filter::Group( + FilterGroup::AnyOf(nested) + | FilterGroup::AllOf(nested) + | FilterGroup::NoneOf(nested), + ) => collect_keys(nested, keys), + } + } +} + +/// `filters` is never empty: the top level is short-circuited by [`CompiledFilters::matches`], +/// and [`drop_empty_groups`] removes empty groups. +fn eval_all(filters: &[Filter], use_or_logic: bool, values: &HashMap<&str, &RawValue>) -> bool { + let eval = |filter: &Filter| match filter { + // Parsed here rather than during the scan so that `any`/`all` short-circuiting keeps + // a large field the verdict never depends on from being materialized at all. + Filter::JsonFilter(JsonFilter { key, value }) => values + .get(key.as_str()) + .and_then(|raw| serde_json::from_str::(raw.get()).ok()) + .map_or(false, |found| is_superset(&found, value)), + Filter::PathFilter(PathFilter { path, value }) => { + let (root, rest) = split_path(path); + values + .get(root) + .and_then(|raw| serde_json::from_str::(raw.get()).ok()) + .and_then(|found| resolve(&found, rest).map(|at| is_superset(at, value))) + .unwrap_or(false) + } + Filter::Group(FilterGroup::AnyOf(nested)) => eval_all(nested, true, values), + Filter::Group(FilterGroup::AllOf(nested)) => eval_all(nested, false, values), + // A key the message does not carry satisfies a negation: nothing there can match. + Filter::Group(FilterGroup::NoneOf(nested)) => !eval_all(nested, true, values), + }; + + if use_or_logic { + filters.iter().any(eval) + } else { + filters.iter().all(eval) + } +} + +/// Locates the requested top-level keys in a single pass, skipping every other value. The +/// ones it wants are borrowed as raw slices of the message rather than deserialized, so a +/// key the boolean evaluation never reaches costs nothing beyond the scan. +struct KeysVisitor<'k> { + keys: &'k [String], +} + +impl<'de, 'k> Visitor<'de> for KeysVisitor<'k> { + type Value = HashMap<&'k str, &'de RawValue>; fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result { - formatter.write_str("a JSON object with a specific key at the top level") + formatter.write_str("a JSON object") } fn visit_map(self, mut map: V) -> std::result::Result where V: MapAccess<'de>, { - let mut result = false; - let mut found = false; + let mut found = HashMap::with_capacity(self.keys.len()); // Must consume entire map to satisfy deserializer contract while let Some(key) = map.next_key::()? { - if !found && key == self.key { - let json_value: Value = map.next_value()?; - result = is_superset(&json_value, self.value_to_check); - found = true; - } else { - // Skip values we don't need (cheaper than full deserialization) - let _ = map.next_value::()?; + match self.keys.iter().find(|k| k.as_str() == key) { + // On a duplicated key the first occurrence wins + Some(k) if !found.contains_key(k.as_str()) => { + found.insert(k.as_str(), map.next_value::<&'de RawValue>()?); + } + _ => { + // Skip values we don't need (cheaper than full deserialization) + let _ = map.next_value::()?; + } } } - Ok(result) + + Ok(found) } } @@ -69,96 +318,359 @@ pub fn is_superset(json_value: &Value, value_to_check: &Value) -> bool { } } -pub fn is_value_superset<'a, 'de, D>( - deserializer: D, - key: &'a str, - value_to_check: &'a Value, -) -> std::result::Result -where - D: Deserializer<'de>, -{ - deserializer.deserialize_map(SupersetVisitor { key, value_to_check }) -} - -pub fn check_filters(text: &str, filters: &[Filter], use_or_logic: bool) -> bool { - if filters.is_empty() { - return true; - } - - let check = |filter: &Filter| -> bool { - match filter { - Filter::JsonFilter(JsonFilter { key, value }) => { - let mut deserializer = serde_json::Deserializer::from_str(text); - is_value_superset(&mut deserializer, key, value).unwrap_or(false) - } - } - }; - - if use_or_logic { - filters.iter().any(check) - } else { - filters.iter().all(check) - } -} - #[cfg(test)] mod tests { use super::*; use serde_json::json; + fn matches(payload: &str, filters: serde_json::Value, use_or_logic: bool) -> bool { + let filters: Vec = serde_json::from_value(filters).unwrap(); + CompiledFilters::new(filters, use_or_logic).matches(payload) + } + #[test] fn test_filter_with_other_top_level_keys() { let payload = r#"{"event_type": "test", "other": "data"}"#; - let key = "event_type"; - let value = json!("test"); - - let mut deserializer = serde_json::Deserializer::from_str(payload); - let result = is_value_superset(&mut deserializer, key, &value).unwrap(); - assert!(result, "Should match when key exists with correct value"); + let filters = json!([{"key": "event_type", "value": "test"}]); + assert!( + matches(payload, filters, false), + "Should match when key exists with correct value" + ); } #[test] fn test_filter_with_key_not_first() { let payload = r#"{"other": "data", "event_type": "test"}"#; - let key = "event_type"; - let value = json!("test"); - - let mut deserializer = serde_json::Deserializer::from_str(payload); - let result = is_value_superset(&mut deserializer, key, &value).unwrap(); - assert!(result, "Should match even when key is not first"); + let filters = json!([{"key": "event_type", "value": "test"}]); + assert!( + matches(payload, filters, false), + "Should match even when key is not first" + ); } #[test] fn test_filter_with_nested_object() { let payload = r#"{"data": {"status": "active", "count": 5}, "other": "value"}"#; - let key = "data"; - let value = json!({"status": "active"}); - - let mut deserializer = serde_json::Deserializer::from_str(payload); - let result = is_value_superset(&mut deserializer, key, &value).unwrap(); - assert!(result, "Should match when nested object is superset"); + let filters = json!([{"key": "data", "value": {"status": "active"}}]); + assert!( + matches(payload, filters, false), + "Should match when nested object is superset" + ); } #[test] fn test_filter_no_match() { let payload = r#"{"event_type": "other", "data": "value"}"#; - let key = "event_type"; - let value = json!("test"); - - let mut deserializer = serde_json::Deserializer::from_str(payload); - let result = is_value_superset(&mut deserializer, key, &value).unwrap(); - assert!(!result, "Should not match when value differs"); + let filters = json!([{"key": "event_type", "value": "test"}]); + assert!( + !matches(payload, filters, false), + "Should not match when value differs" + ); } #[test] fn test_filter_key_not_found() { let payload = r#"{"other": "data"}"#; - let key = "event_type"; - let value = json!("test"); + let filters = json!([{"key": "event_type", "value": "test"}]); + assert!( + !matches(payload, filters, false), + "Should not match when key doesn't exist" + ); + } - let mut deserializer = serde_json::Deserializer::from_str(payload); - let result = is_value_superset(&mut deserializer, key, &value).unwrap(); - assert!(!result, "Should not match when key doesn't exist"); + #[test] + fn test_no_filters_matches_everything() { + assert!(matches(r#"{"a": 1}"#, json!([]), false)); + assert!(matches("not even json", json!([]), true)); + } + + #[test] + fn test_non_object_payload_never_matches() { + let filters = json!([{"key": "a", "value": 1}]); + assert!(!matches("[1, 2]", filters.clone(), false)); + assert!(!matches("nope", filters, true)); + } + + #[test] + fn test_top_level_and_or_logic() { + let payload = r#"{"a": 1, "b": 2}"#; + let filters = json!([{"key": "a", "value": 1}, {"key": "b", "value": 99}]); + assert!(!matches(payload, filters.clone(), false)); + assert!(matches(payload, filters, true)); + } + + // --- nested groups --- + + #[test] + fn test_any_of_group_nested_in_and() { + let payload = + r#"{"event": "message_created", "previous_message": {"sent_by": "reminder"}}"#; + let filters = json!([ + {"key": "event", "value": "message_created"}, + {"any_of": [ + {"key": "in_reply_to", "value": {"sent_by": "reminder"}}, + {"key": "previous_message", "value": {"sent_by": "reminder"}} + ]} + ]); + assert!(matches(payload, filters, false)); + } + + #[test] + fn test_any_of_group_all_branches_fail() { + let payload = r#"{"event": "message_created", "previous_message": {"sent_by": "someone"}}"#; + let filters = json!([ + {"key": "event", "value": "message_created"}, + {"any_of": [ + {"key": "in_reply_to", "value": {"sent_by": "reminder"}}, + {"key": "previous_message", "value": {"sent_by": "reminder"}} + ]} + ]); + assert!(!matches(payload, filters, false)); + } + + #[test] + fn test_all_of_group_nested_in_or() { + let payload = r#"{"a": 1, "b": 2}"#; + let filters = json!([ + {"key": "missing", "value": true}, + {"all_of": [{"key": "a", "value": 1}, {"key": "b", "value": 2}]} + ]); + assert!(matches(payload, filters.clone(), true)); + assert!(!matches(payload, filters, false)); + } + + /// `1e400` overflows `Value`'s f64 and only fails to parse if something reads it, so a + /// match here means the short-circuit really did skip that field rather than + /// materializing every referenced key up front. + #[test] + fn test_unreached_branch_is_never_materialized() { + let payload = r#"{"gate": "match", "huge": 1e400}"#; + let filters = json!([{"key": "gate", "value": "match"}, {"key": "huge", "value": 1}]); + assert!(matches(payload, filters.clone(), true)); + assert!(!matches(payload, filters, false)); + } + + #[test] + fn test_deeply_nested_groups() { + let payload = r#"{"a": 1, "b": 2, "c": 3}"#; + let filters = json!([ + {"any_of": [ + {"key": "a", "value": 99}, + {"all_of": [ + {"key": "b", "value": 2}, + {"any_of": [{"key": "c", "value": 3}, {"key": "c", "value": 4}]} + ]} + ]} + ]); + assert!(matches(payload, filters, false)); + } + + #[test] + fn test_path_reaches_a_nested_field() { + let payload = r#"{"in_reply_to_message": {"content_attributes": {"sent_by": "reminder"}}}"#; + let path = "in_reply_to_message.content_attributes.sent_by"; + assert!(matches( + payload, + json!([{"path": path, "value": "reminder"}]), + false + )); + assert!(!matches( + payload, + json!([{"path": path, "value": "other"}]), + false + )); + // a missing intermediate segment is a miss, not an error + assert!(!matches( + payload, + json!([{"path": "in_reply_to_message.nope.sent_by", "value": "reminder"}]), + false + )); + } + + #[test] + fn test_path_and_key_keep_their_own_meaning() { + // `key` still addresses the top-level field spelled with dots, `path` traverses + assert!(matches( + r#"{"a.b": 1}"#, + json!([{"key": "a.b", "value": 1}]), + false + )); + assert!(!matches( + r#"{"a.b": 1}"#, + json!([{"path": "a.b", "value": 1}]), + false + )); + assert!(matches( + r#"{"a": {"b": 1}}"#, + json!([{"path": "a.b", "value": 1}]), + false + )); + assert!(!matches( + r#"{"a": {"b": 1}}"#, + json!([{"key": "a.b", "value": 1}]), + false + )); + } + + #[test] + fn test_path_does_not_traverse_arrays() { + // Deliberately unsupported for now: an element index is not implied + assert!(!matches( + r#"{"items": [{"id": 1}]}"#, + json!([{"path": "items.id", "value": 1}]), + false + )); + } + + #[test] + fn test_path_without_dots_is_a_top_level_field() { + assert!(matches( + r#"{"a": 1}"#, + json!([{"path": "a", "value": 1}]), + false + )); + } + + #[test] + fn test_none_of_excludes_matching_messages() { + let filters = json!([ + {"key": "event", "value": "message_created"}, + {"none_of": [{"key": "sender", "value": "bot"}, {"key": "kind", "value": "draft"}]} + ]); + assert!(matches( + r#"{"event": "message_created", "sender": "human"}"#, + filters.clone(), + false + )); + assert!(!matches( + r#"{"event": "message_created", "sender": "bot"}"#, + filters.clone(), + false + )); + // any one branch matching is enough to exclude + assert!(!matches( + r#"{"event": "message_created", "sender": "human", "kind": "draft"}"#, + filters, + false + )); + } + + #[test] + fn test_none_of_is_satisfied_by_a_missing_key() { + // Nothing is there to match, so the negation holds — the alternative would make + // every negative filter also require the field to be present. + assert!(matches( + r#"{"event": "message_created"}"#, + json!([{"none_of": [{"key": "sender", "value": "bot"}]}]), + false + )); + } + + #[test] + fn test_empty_group_is_dropped_not_constant() { + let payload = r#"{"a": 1}"#; + // On its own it leaves the trigger unfiltered, like an empty filter list + assert!(matches(payload, json!([{"any_of": []}]), false)); + assert!(matches( + payload, + json!([{"all_of": []}, {"any_of": [{"all_of": []}]}]), + true + )); + // Alongside a real criterion it must not decide the outcome either way + let with_failing_leaf = json!([{"any_of": []}, {"key": "a", "value": 99}]); + assert!(!matches(payload, with_failing_leaf.clone(), true)); + assert!(!matches(payload, with_failing_leaf, false)); + } + + #[test] + fn test_parses_legacy_and_group_entries_side_by_side() { + let filters = CompiledFilters::parse( + [ + r#"{"key": "event", "value": "created"}"#, + r#"{"any_of": [{"key": "a", "value": 1}, {"key": "b", "value": 2}]}"#, + ], + false, + "u/admin/trigger", + ); + assert!(filters.matches(r#"{"event": "created", "b": 2}"#)); + assert!(!filters.matches(r#"{"event": "created", "b": 3}"#)); + assert!(!filters.matches(r#"{"event": "other", "a": 1}"#)); + } + + #[test] + fn test_duplicated_payload_key_resolves_to_first_occurrence() { + let filters = json!([{"key": "a", "value": 1}]); + assert!(matches(r#"{"a": 1, "a": 2}"#, filters.clone(), false)); + assert!(!matches(r#"{"a": 2, "a": 1}"#, filters, false)); + } + + #[test] + fn test_validate_rejects_entries_the_listener_would_drop() { + assert!(CompiledFilters::validate(&[ + json!({"key": "a", "value": 1}), + json!({"all_of": []}) + ]) + .is_ok()); + assert!( + CompiledFilters::validate(&[json!({"anyOf": [{"key": "a", "value": 1}]})]).is_err() + ); + assert!(CompiledFilters::validate(&[json!({"key": "a"})]).is_err()); + } + + #[test] + fn test_validate_rejects_a_leaf_that_also_carries_a_group() { + // Untagged would settle each of these on one variant and drop the rest of the entry + for mixed in [ + json!({"key": "a", "value": 1, "none_of": [{"key": "b", "value": 2}]}), + json!({"path": "a.b", "value": 1, "any_of": [{"key": "b", "value": 2}]}), + json!({"any_of": [{"key": "a", "value": 1}], "all_of": [{"key": "b", "value": 2}]}), + ] { + let err = CompiledFilters::validate(&[mixed.clone()]) + .unwrap_err() + .to_string(); + assert!( + err.contains("combines a criterion with"), + "{} should be rejected, got: {}", + mixed, + err + ); + } + } + + #[test] + fn test_validate_rejects_a_leaf_naming_both_key_and_path() { + // Untagged would take it as a `key` criterion and drop the `path` without a word + let err = CompiledFilters::validate(&[json!({"key": "a", "path": "b.c", "value": 1})]) + .unwrap_err() + .to_string(); + assert!(err.contains("both key and path"), "got: {}", err); + } + + #[test] + fn test_validate_rejects_an_empty_path_segment() { + assert!(CompiledFilters::validate(&[json!({"path": "a.b", "value": 1})]).is_ok()); + for dead in ["", "a.", ".a", "a..b"] { + assert!( + CompiledFilters::validate(&[json!({"path": dead, "value": 1})]).is_err(), + "path {:?} addresses no field and should be rejected", + dead + ); + } + } + + #[test] + fn test_validate_names_the_offending_nested_entry() { + let err = CompiledFilters::validate(&[ + json!({"key": "a", "value": 1}), + json!({"all_of": [{"key": "b", "value": 2}, {"any_of": [{"key": "c"}]}]}), + ]) + .unwrap_err() + .to_string(); + assert!( + err.contains("filter #2 -> all_of[1] -> any_of[0]"), + "error should point at the entry that is wrong, got: {}", + err + ); } // --- is_superset unit tests --- diff --git a/cli/src/guidance/skills.gen.ts b/cli/src/guidance/skills.gen.ts index 56cb9730cb..eb386b3ec8 100644 --- a/cli/src/guidance/skills.gen.ts +++ b/cli/src/guidance/skills.gen.ts @@ -8287,18 +8287,62 @@ properties: filters: type: array items: - type: object - properties: - key: - type: string - value: {} + oneOf: + - type: object + properties: + key: + type: string + value: {} + required: + - key + - value + - type: object + properties: + path: + type: string + description: Dotted path into nested objects, e.g. \`a.b.c\`. Does not traverse + arrays. + value: {} + required: + - path + - value + - type: object + properties: + any_of: + type: array + items: + type: object + required: + - any_of + - type: object + properties: + all_of: + type: array + items: + type: object + required: + - all_of + - type: object + properties: + none_of: + type: array + items: + type: object + required: + - none_of + description: 'Filters to match incoming messages (only matching messages trigger + the script). Each entry is either a leaf \`{key, value}\` (top-level field) or + \`{path, value}\` (dotted path into nested objects), or a group \`{any_of: [...]}\` + / \`{all_of: [...]}\` / \`{none_of: [...]}\` nesting more entries. Entries at the + top level are combined with \`filter_logic\`.' filter_logic: type: string enum: - and - or - description: Logic to apply when evaluating filters. 'and' requires all filters - to match, 'or' requires any filter to match. + description: Logic to apply when evaluating the top-level filters. 'and' requires + all of them to match, 'or' requires any of them to match. Nested \`any_of\`/\`all_of\`/\`none_of\` + groups carry their own logic. auto_offset_reset: type: string enum: @@ -8399,6 +8443,18 @@ properties: type: array items: type: object + properties: + qos: + type: string + enum: + - qos0 + - qos1 + - qos2 + topic: + type: string + required: + - qos + - topic description: Array of MQTT topics to subscribe to, each with topic name and QoS level v3_config: @@ -8943,24 +8999,91 @@ properties: filters: type: array items: - type: object - properties: - key: - type: string - value: {} - description: Array of key-value filters to match incoming messages (only matching - messages trigger the script) + oneOf: + - type: object + properties: + key: + type: string + value: {} + required: + - key + - value + - type: object + properties: + path: + type: string + description: Dotted path into nested objects, e.g. \`a.b.c\`. Does not traverse + arrays. + value: {} + required: + - path + - value + - type: object + properties: + any_of: + type: array + items: + type: object + required: + - any_of + - type: object + properties: + all_of: + type: array + items: + type: object + required: + - all_of + - type: object + properties: + none_of: + type: array + items: + type: object + required: + - none_of + description: 'Filters to match incoming messages (only matching messages trigger + the script). Each entry is either a leaf \`{key, value}\` (top-level field) or + \`{path, value}\` (dotted path into nested objects), or a group \`{any_of: [...]}\` + / \`{all_of: [...]}\` / \`{none_of: [...]}\` nesting more entries. Entries at the + top level are combined with \`filter_logic\`.' filter_logic: type: string enum: - and - or - description: Logic to apply when evaluating filters. 'and' requires all filters - to match, 'or' requires any filter to match. + description: Logic to apply when evaluating the top-level filters. 'and' requires + all of them to match, 'or' requires any of them to match. Nested \`any_of\`/\`all_of\`/\`none_of\` + groups carry their own logic. initial_messages: type: array items: - type: object + oneOf: + - type: object + properties: + raw_message: + type: string + required: + - raw_message + - type: object + properties: + runnable_result: + type: object + properties: + path: + type: string + args: + type: object + description: The arguments to pass to the script or flow + additionalProperties: true + is_flow: + type: boolean + required: + - path + - args + - is_flow + required: + - runnable_result description: Messages to send immediately after connecting (can be raw strings or computed by runnables) url_runnable_args: diff --git a/frontend/src/lib/components/copilot/chat/workspaceToolsZod.gen.ts b/frontend/src/lib/components/copilot/chat/workspaceToolsZod.gen.ts index 0518fde90b..d22dceffa6 100644 --- a/frontend/src/lib/components/copilot/chat/workspaceToolsZod.gen.ts +++ b/frontend/src/lib/components/copilot/chat/workspaceToolsZod.gen.ts @@ -97,11 +97,20 @@ export const websocketTriggerRequestSchema = z.object({ "is_flow": z.boolean().describe("True if script_path points to a flow, false if it points to a script"), "url": z.string().describe("The WebSocket URL to connect to (can be a static URL or computed by a runnable)"), "mode": z.enum(["enabled", "disabled", "suspended"]).describe("job trigger mode").optional(), - "filters": z.array(z.object({ + "filters": z.array(z.union([z.object({ "key": z.string(), "value": z.any() - })).describe("Array of key-value filters to match incoming messages (only matching messages trigger the script)"), - "filter_logic": z.enum(["and", "or"]).describe("Logic to apply when evaluating filters. 'and' requires all filters to match, 'or' requires any filter to match.").default("and").optional(), + }), z.object({ + "path": z.string().describe("Dotted path into nested objects, e.g. `a.b.c`. Does not traverse arrays."), + "value": z.any() + }), z.object({ + "any_of": z.array(z.record(z.string(), z.any())) + }), z.object({ + "all_of": z.array(z.record(z.string(), z.any())) + }), z.object({ + "none_of": z.array(z.record(z.string(), z.any())) + })]).describe("Either a leaf filter, matching a field of the message (parsed as JSON) against a value by equality (or superset, when the value is an object or array) \u2014 addressed by `key` for a top-level field or `path` for a dotted path into nested objects \u2014 or a group nesting sub-filters under a boolean operator (`none_of` matches when none of its sub-filters do).\n")).describe("Filters to match incoming messages (only matching messages trigger the script). Each entry is either a leaf `{key, value}` (top-level field) or `{path, value}` (dotted path into nested objects), or a group `{any_of: [...]}` / `{all_of: [...]}` / `{none_of: [...]}` nesting more entries. Entries at the top level are combined with `filter_logic`."), + "filter_logic": z.enum(["and", "or"]).describe("Logic to apply when evaluating the top-level filters. 'and' requires all of them to match, 'or' requires any of them to match. Nested `any_of`/`all_of`/`none_of` groups carry their own logic.").default("and").optional(), "initial_messages": z.array(z.union([z.object({ "raw_message": z.string() }), z.object({ @@ -148,11 +157,20 @@ export const kafkaTriggerRequestSchema = z.object({ "kafka_resource_path": z.string().describe("Path to the Kafka resource containing connection configuration"), "group_id": z.string().describe("Kafka consumer group ID for this trigger"), "topics": z.array(z.string()).describe("Array of Kafka topic names to subscribe to"), - "filters": z.array(z.object({ + "filters": z.array(z.union([z.object({ "key": z.string(), "value": z.any() - })), - "filter_logic": z.enum(["and", "or"]).describe("Logic to apply when evaluating filters. 'and' requires all filters to match, 'or' requires any filter to match.").default("and").optional(), + }), z.object({ + "path": z.string().describe("Dotted path into nested objects, e.g. `a.b.c`. Does not traverse arrays."), + "value": z.any() + }), z.object({ + "any_of": z.array(z.record(z.string(), z.any())) + }), z.object({ + "all_of": z.array(z.record(z.string(), z.any())) + }), z.object({ + "none_of": z.array(z.record(z.string(), z.any())) + })]).describe("Either a leaf filter, matching a field of the message (parsed as JSON) against a value by equality (or superset, when the value is an object or array) \u2014 addressed by `key` for a top-level field or `path` for a dotted path into nested objects \u2014 or a group nesting sub-filters under a boolean operator (`none_of` matches when none of its sub-filters do).\n")).describe("Filters to match incoming messages (only matching messages trigger the script). Each entry is either a leaf `{key, value}` (top-level field) or `{path, value}` (dotted path into nested objects), or a group `{any_of: [...]}` / `{all_of: [...]}` / `{none_of: [...]}` nesting more entries. Entries at the top level are combined with `filter_logic`."), + "filter_logic": z.enum(["and", "or"]).describe("Logic to apply when evaluating the top-level filters. 'and' requires all of them to match, 'or' requires any of them to match. Nested `any_of`/`all_of`/`none_of` groups carry their own logic.").default("and").optional(), "auto_offset_reset": z.enum(["latest", "earliest"]).describe("Initial offset behavior when consumer group has no committed offset.").default("latest").optional(), "auto_commit": z.boolean().describe("When true (default), offsets are committed automatically after receiving each message. When false, you must manually commit offsets using the commit_offsets endpoint.").default(true).optional(), "mode": z.enum(["enabled", "disabled", "suspended"]).describe("job trigger mode").optional(), diff --git a/frontend/src/lib/components/triggers/TriggerFilterList.svelte b/frontend/src/lib/components/triggers/TriggerFilterList.svelte new file mode 100644 index 0000000000..d1d0522a65 --- /dev/null +++ b/frontend/src/lib/components/triggers/TriggerFilterList.svelte @@ -0,0 +1,161 @@ + + +
+ {#if depth > 0 || filters.length > 0} +
+