From 3876902a7be798fd5ef208bc5756b28fb55e569e Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Mon, 30 Mar 2026 19:32:24 +0000 Subject: [PATCH] feat: add OR logic support to kafka/websocket trigger filters (#8580) * feat: add OR logic support to kafka/websocket trigger filters Co-Authored-By: Claude Opus 4.6 (1M context) * chore: update ee-repo-ref for OR logic filter support Co-Authored-By: Claude Opus 4.6 (1M context) * fix: add filter_logic to OpenAPI spec/save utils, fix websocket derive, show capture group ID - Add filter_logic field to all 6 Kafka/WebSocket OpenAPI schemas so it is included in the generated frontend client types - Include filter_logic in save request bodies (kafka/utils.ts, websocket/utils.ts) - Fix misplaced #[derive(FromRow)] on WebsocketConfig (was on the default fn) - Show copyable "Test group ID" in Kafka capture UI - Remove capture event-loss warning for Kafka (uses separate consumer group) Co-Authored-By: Claude Opus 4.6 (1M context) * update sqlx * update ee ref * chore: regenerate system prompts for filter_logic schema changes Co-Authored-By: Claude Opus 4.6 (1M context) * fix: remove banned $bindable(default_value) pattern in TriggerFilters Use $bindable() without default and $derived with ?? for the effective value, per CLAUDE.md rules. Co-Authored-By: Claude Opus 4.6 (1M context) * fix: make filterLogic prop required in TriggerFilters All callers always pass it, no need for optional + derived fallback. Co-Authored-By: Claude Opus 4.6 (1M context) * chore: update ee-repo-ref to 5ee1382dfb23b6a1516e3c7586058cec8240fdf2 This commit updates the EE repository reference after PR #498 was merged in windmill-ee-private. Previous ee-repo-ref: bbd674991c07bff1cb2f3744e71fda10df53f09d New ee-repo-ref: 5ee1382dfb23b6a1516e3c7586058cec8240fdf2 Automated by sync-ee-ref workflow. * fix: reset filterLogic to 'and' in openNew for kafka/websocket editors Prevents stale OR logic from carrying over when creating a new trigger after editing one with OR filters. Co-Authored-By: Claude Opus 4.6 (1M context) --------- Co-authored-by: Claude Opus 4.6 (1M context) Co-authored-by: hugocasa Co-authored-by: windmill-internal-app[bot] --- ...92ba37da67cd2a56834c9d1378eab7551284d.json | 30 +++++++++++++ ...472080ecab3ed652912397245e5216ae0389.json} | 5 ++- ...70aa2e83fcb5aa0c953febdd9bac2d95bbec.json} | 5 ++- ...5ebc784e5da5ee25d47af187a75220d8fded7.json | 29 ------------- ...cbce7cdbdd1feb553718cbde60bb8ccff4733.json | 29 ------------- ...aca47771b480ff194884f9c947dcaf71d6cf9.json | 30 +++++++++++++ backend/ee-repo-ref.txt | 2 +- ...260328000000_trigger_filter_logic.down.sql | 2 + ...20260328000000_trigger_filter_logic.up.sql | 2 + backend/windmill-api/openapi.yaml | 42 +++++++++++++++++++ .../windmill-trigger-websocket/src/handler.rs | 27 +++++++----- backend/windmill-trigger-websocket/src/lib.rs | 16 +++++-- .../src/listener.rs | 24 ++--------- backend/windmill-trigger/src/filter.rs | 21 ++++++++++ cli/src/guidance/skills.ts | 14 +++++++ .../components/triggers/CaptureWrapper.svelte | 2 +- .../components/triggers/TriggerFilters.svelte | 26 +++++++++++- .../triggers/kafka/KafkaCapture.svelte | 11 +++-- .../kafka/KafkaTriggerEditorInner.svelte | 6 ++- .../lib/components/triggers/kafka/utils.ts | 1 + .../WebsocketTriggerEditorInner.svelte | 6 ++- .../components/triggers/websocket/utils.ts | 1 + .../schemas/kafka_trigger.schema.yaml | 7 ++++ .../schemas/websocket_trigger.schema.yaml | 7 ++++ 24 files changed, 238 insertions(+), 107 deletions(-) create mode 100644 backend/.sqlx/query-68c19cb0e18b94870bbe81f9aab92ba37da67cd2a56834c9d1378eab7551284d.json rename backend/.sqlx/{query-942c0abb55c910862fd45d3fa56a4eb6729f1a658101bda2d0b0fca96b3cfee5.json => query-6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389.json} (60%) rename backend/.sqlx/{query-a0a545fda5f3ebea0113d5daaf13358c964d9fb0f41bf2a1c834305b4d2398f2.json => query-6a8f4ed9946bb2a3c5e90695c90b70aa2e83fcb5aa0c953febdd9bac2d95bbec.json} (57%) delete mode 100644 backend/.sqlx/query-a37cfc632dd37cf37c06743239b5ebc784e5da5ee25d47af187a75220d8fded7.json delete mode 100644 backend/.sqlx/query-c7aed7fe3b6774477d403bc3e7fcbce7cdbdd1feb553718cbde60bb8ccff4733.json create mode 100644 backend/.sqlx/query-e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9.json create mode 100644 backend/migrations/20260328000000_trigger_filter_logic.down.sql create mode 100644 backend/migrations/20260328000000_trigger_filter_logic.up.sql diff --git a/backend/.sqlx/query-68c19cb0e18b94870bbe81f9aab92ba37da67cd2a56834c9d1378eab7551284d.json b/backend/.sqlx/query-68c19cb0e18b94870bbe81f9aab92ba37da67cd2a56834c9d1378eab7551284d.json new file mode 100644 index 0000000000..de69c0fc8c --- /dev/null +++ b/backend/.sqlx/query-68c19cb0e18b94870bbe81f9aab92ba37da67cd2a56834c9d1378eab7551284d.json @@ -0,0 +1,30 @@ +{ + "db_name": "PostgreSQL", + "query": "\n UPDATE kafka_trigger\n SET\n kafka_resource_path = $1,\n group_id = $2,\n topics = $3,\n filters = $4,\n filter_logic = $5,\n auto_offset_reset = $6,\n auto_commit = $7,\n script_path = $8,\n path = $9,\n is_flow = $10,\n edited_by = $11,\n permissioned_as = $12,\n edited_at = now(),\n server_id = NULL,\n error = NULL,\n error_handler_path = $15,\n error_handler_args = $16,\n retry = $17\n WHERE\n workspace_id = $13 AND path = $14\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "VarcharArray", + "JsonbArray", + "Varchar", + "Varchar", + "Bool", + "Varchar", + "Varchar", + "Bool", + "Varchar", + "Varchar", + "Text", + "Text", + "Varchar", + "Jsonb", + "Jsonb" + ] + }, + "nullable": [] + }, + "hash": "68c19cb0e18b94870bbe81f9aab92ba37da67cd2a56834c9d1378eab7551284d" +} diff --git a/backend/.sqlx/query-942c0abb55c910862fd45d3fa56a4eb6729f1a658101bda2d0b0fca96b3cfee5.json b/backend/.sqlx/query-6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389.json similarity index 60% rename from backend/.sqlx/query-942c0abb55c910862fd45d3fa56a4eb6729f1a658101bda2d0b0fca96b3cfee5.json rename to backend/.sqlx/query-6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389.json index e3ca43fbd0..a6f1f5f7cc 100644 --- a/backend/.sqlx/query-942c0abb55c910862fd45d3fa56a4eb6729f1a658101bda2d0b0fca96b3cfee5.json +++ b/backend/.sqlx/query-6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n INSERT INTO websocket_trigger (\n workspace_id,\n path,\n url,\n script_path,\n is_flow,\n mode,\n filters,\n initial_messages,\n url_runnable_args,\n edited_by,\n can_return_message,\n can_return_error_result,\n permissioned_as,\n edited_at,\n error_handler_path,\n error_handler_args,\n retry\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, now(), $14, $15, $16\n )\n ", + "query": "\n INSERT INTO websocket_trigger (\n workspace_id,\n path,\n url,\n script_path,\n is_flow,\n mode,\n filters,\n filter_logic,\n initial_messages,\n url_runnable_args,\n edited_by,\n can_return_message,\n can_return_error_result,\n permissioned_as,\n edited_at,\n error_handler_path,\n error_handler_args,\n retry\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, now(), $15, $16, $17\n )\n ", "describe": { "columns": [], "parameters": { @@ -23,6 +23,7 @@ } }, "JsonbArray", + "Varchar", "JsonbArray", "Jsonb", "Varchar", @@ -36,5 +37,5 @@ }, "nullable": [] }, - "hash": "942c0abb55c910862fd45d3fa56a4eb6729f1a658101bda2d0b0fca96b3cfee5" + "hash": "6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389" } diff --git a/backend/.sqlx/query-a0a545fda5f3ebea0113d5daaf13358c964d9fb0f41bf2a1c834305b4d2398f2.json b/backend/.sqlx/query-6a8f4ed9946bb2a3c5e90695c90b70aa2e83fcb5aa0c953febdd9bac2d95bbec.json similarity index 57% rename from backend/.sqlx/query-a0a545fda5f3ebea0113d5daaf13358c964d9fb0f41bf2a1c834305b4d2398f2.json rename to backend/.sqlx/query-6a8f4ed9946bb2a3c5e90695c90b70aa2e83fcb5aa0c953febdd9bac2d95bbec.json index 99d8d2a181..c6d3374639 100644 --- a/backend/.sqlx/query-a0a545fda5f3ebea0113d5daaf13358c964d9fb0f41bf2a1c834305b4d2398f2.json +++ b/backend/.sqlx/query-6a8f4ed9946bb2a3c5e90695c90b70aa2e83fcb5aa0c953febdd9bac2d95bbec.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n INSERT INTO kafka_trigger (\n workspace_id,\n path,\n kafka_resource_path,\n group_id,\n topics,\n filters,\n auto_offset_reset,\n auto_commit,\n script_path,\n is_flow,\n mode,\n edited_by,\n permissioned_as,\n edited_at,\n error_handler_path,\n error_handler_args,\n retry\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, now(), $14, $15, $16\n )\n ", + "query": "\n INSERT INTO kafka_trigger (\n workspace_id,\n path,\n kafka_resource_path,\n group_id,\n topics,\n filters,\n filter_logic,\n auto_offset_reset,\n auto_commit,\n script_path,\n is_flow,\n mode,\n edited_by,\n permissioned_as,\n edited_at,\n error_handler_path,\n error_handler_args,\n retry\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, now(), $15, $16, $17\n )\n ", "describe": { "columns": [], "parameters": { @@ -12,6 +12,7 @@ "VarcharArray", "JsonbArray", "Varchar", + "Varchar", "Bool", "Varchar", "Bool", @@ -36,5 +37,5 @@ }, "nullable": [] }, - "hash": "a0a545fda5f3ebea0113d5daaf13358c964d9fb0f41bf2a1c834305b4d2398f2" + "hash": "6a8f4ed9946bb2a3c5e90695c90b70aa2e83fcb5aa0c953febdd9bac2d95bbec" } diff --git a/backend/.sqlx/query-a37cfc632dd37cf37c06743239b5ebc784e5da5ee25d47af187a75220d8fded7.json b/backend/.sqlx/query-a37cfc632dd37cf37c06743239b5ebc784e5da5ee25d47af187a75220d8fded7.json deleted file mode 100644 index c993120dae..0000000000 --- a/backend/.sqlx/query-a37cfc632dd37cf37c06743239b5ebc784e5da5ee25d47af187a75220d8fded7.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE kafka_trigger\n SET\n kafka_resource_path = $1,\n group_id = $2,\n topics = $3,\n filters = $4,\n auto_offset_reset = $5,\n auto_commit = $6,\n script_path = $7,\n path = $8,\n is_flow = $9,\n edited_by = $10,\n permissioned_as = $11,\n edited_at = now(),\n server_id = NULL,\n error = NULL,\n error_handler_path = $14,\n error_handler_args = $15,\n retry = $16\n WHERE\n workspace_id = $12 AND path = $13\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Varchar", - "Varchar", - "VarcharArray", - "JsonbArray", - "Varchar", - "Bool", - "Varchar", - "Varchar", - "Bool", - "Varchar", - "Varchar", - "Text", - "Text", - "Varchar", - "Jsonb", - "Jsonb" - ] - }, - "nullable": [] - }, - "hash": "a37cfc632dd37cf37c06743239b5ebc784e5da5ee25d47af187a75220d8fded7" -} diff --git a/backend/.sqlx/query-c7aed7fe3b6774477d403bc3e7fcbce7cdbdd1feb553718cbde60bb8ccff4733.json b/backend/.sqlx/query-c7aed7fe3b6774477d403bc3e7fcbce7cdbdd1feb553718cbde60bb8ccff4733.json deleted file mode 100644 index 772b86a11c..0000000000 --- a/backend/.sqlx/query-c7aed7fe3b6774477d403bc3e7fcbce7cdbdd1feb553718cbde60bb8ccff4733.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE\n websocket_trigger\n SET\n url = $1,\n script_path = $2,\n path = $3,\n is_flow = $4,\n filters = $5,\n initial_messages = $6,\n url_runnable_args = $7,\n edited_by = $8,\n permissioned_as = $9,\n can_return_message = $10,\n can_return_error_result = $11,\n edited_at = now(),\n server_id = NULL,\n error = NULL,\n error_handler_path = $14,\n error_handler_args = $15,\n retry = $16\n WHERE\n workspace_id = $12 AND path = $13\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Varchar", - "Varchar", - "Varchar", - "Bool", - "JsonbArray", - "JsonbArray", - "Jsonb", - "Varchar", - "Varchar", - "Bool", - "Bool", - "Text", - "Text", - "Varchar", - "Jsonb", - "Jsonb" - ] - }, - "nullable": [] - }, - "hash": "c7aed7fe3b6774477d403bc3e7fcbce7cdbdd1feb553718cbde60bb8ccff4733" -} diff --git a/backend/.sqlx/query-e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9.json b/backend/.sqlx/query-e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9.json new file mode 100644 index 0000000000..8ff6f2e89c --- /dev/null +++ b/backend/.sqlx/query-e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9.json @@ -0,0 +1,30 @@ +{ + "db_name": "PostgreSQL", + "query": "\n UPDATE\n websocket_trigger\n SET\n url = $1,\n script_path = $2,\n path = $3,\n is_flow = $4,\n filters = $5,\n filter_logic = $6,\n initial_messages = $7,\n url_runnable_args = $8,\n edited_by = $9,\n permissioned_as = $10,\n can_return_message = $11,\n can_return_error_result = $12,\n edited_at = now(),\n server_id = NULL,\n error = NULL,\n error_handler_path = $15,\n error_handler_args = $16,\n retry = $17\n WHERE\n workspace_id = $13 AND path = $14\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Varchar", + "Bool", + "JsonbArray", + "Varchar", + "JsonbArray", + "Jsonb", + "Varchar", + "Varchar", + "Bool", + "Bool", + "Text", + "Text", + "Varchar", + "Jsonb", + "Jsonb" + ] + }, + "nullable": [] + }, + "hash": "e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9" +} diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index aad0b0cd6c..3fa40f90dc 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -780857855e231c9d71f02fefd8253c254542ef32 +5ee1382dfb23b6a1516e3c7586058cec8240fdf2 diff --git a/backend/migrations/20260328000000_trigger_filter_logic.down.sql b/backend/migrations/20260328000000_trigger_filter_logic.down.sql new file mode 100644 index 0000000000..f6beac9bff --- /dev/null +++ b/backend/migrations/20260328000000_trigger_filter_logic.down.sql @@ -0,0 +1,2 @@ +ALTER TABLE kafka_trigger DROP COLUMN filter_logic; +ALTER TABLE websocket_trigger DROP COLUMN filter_logic; diff --git a/backend/migrations/20260328000000_trigger_filter_logic.up.sql b/backend/migrations/20260328000000_trigger_filter_logic.up.sql new file mode 100644 index 0000000000..ecb99bf396 --- /dev/null +++ b/backend/migrations/20260328000000_trigger_filter_logic.up.sql @@ -0,0 +1,2 @@ +ALTER TABLE kafka_trigger ADD COLUMN filter_logic VARCHAR(3) NOT NULL DEFAULT 'and'; +ALTER TABLE websocket_trigger ADD COLUMN filter_logic VARCHAR(3) NOT NULL DEFAULT 'and'; diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 1adc1e9997..23b2094676 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -21814,6 +21814,13 @@ components: required: - key - value + 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." initial_messages: type: array nullable: true @@ -21875,6 +21882,13 @@ components: required: - key - value + 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." initial_messages: type: array nullable: true @@ -21943,6 +21957,13 @@ components: required: - key - value + 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." initial_messages: type: array nullable: true @@ -22819,6 +22840,13 @@ components: required: - key - value + 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." auto_offset_reset: type: string enum: @@ -22890,6 +22918,13 @@ components: required: - key - value + 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." auto_offset_reset: type: string enum: @@ -22953,6 +22988,13 @@ components: required: - key - value + 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." 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 8411bb0246..223868e96e 100644 --- a/backend/windmill-trigger-websocket/src/handler.rs +++ b/backend/windmill-trigger-websocket/src/handler.rs @@ -36,6 +36,7 @@ impl TriggerCrud for WebsocketTrigger { const ADDITIONAL_SELECT_FIELDS: &[&'static str] = &[ "url", "filters", + "filter_logic", "initial_messages", "url_runnable_args", "can_return_message", @@ -103,6 +104,7 @@ impl TriggerCrud for WebsocketTrigger { is_flow, mode, filters, + filter_logic, initial_messages, url_runnable_args, edited_by, @@ -114,7 +116,7 @@ impl TriggerCrud for WebsocketTrigger { error_handler_args, retry ) VALUES ( - $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, now(), $14, $15, $16 + $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, now(), $15, $16, $17 ) "#, w_id, @@ -124,6 +126,7 @@ impl TriggerCrud for WebsocketTrigger { trigger.base.is_flow, trigger.base.mode() as _, &filters as _, + trigger.config.filter_logic, &initial_messages as _, trigger .config @@ -178,26 +181,28 @@ impl TriggerCrud for WebsocketTrigger { path = $3, is_flow = $4, filters = $5, - initial_messages = $6, - url_runnable_args = $7, - edited_by = $8, - permissioned_as = $9, - can_return_message = $10, - can_return_error_result = $11, + filter_logic = $6, + initial_messages = $7, + url_runnable_args = $8, + edited_by = $9, + permissioned_as = $10, + can_return_message = $11, + can_return_error_result = $12, edited_at = now(), server_id = NULL, error = NULL, - error_handler_path = $14, - error_handler_args = $15, - retry = $16 + error_handler_path = $15, + error_handler_args = $16, + retry = $17 WHERE - workspace_id = $12 AND path = $13 + workspace_id = $13 AND path = $14 ", trigger.config.url, trigger.base.script_path, trigger.base.path, trigger.base.is_flow, filters.as_slice() as &[SqlxJson>], + trigger.config.filter_logic, initial_messages.as_slice() as &[SqlxJson>], trigger .config diff --git a/backend/windmill-trigger-websocket/src/lib.rs b/backend/windmill-trigger-websocket/src/lib.rs index ca754d8349..4fe96bb7c2 100644 --- a/backend/windmill-trigger-websocket/src/lib.rs +++ b/backend/windmill-trigger-websocket/src/lib.rs @@ -1,12 +1,9 @@ use std::collections::HashMap; -use windmill_api_auth::ApiAuthed; -use windmill_trigger::trigger_helpers::{ - trigger_runnable_and_wait_for_raw_result_with_error_ctx, TriggerJobArgs, -}; use serde::{Deserialize, Serialize}; use serde_json::value::RawValue; use sqlx::{types::Json as SqlxJson, FromRow}; +use windmill_api_auth::ApiAuthed; use windmill_common::{ error::{Error, Result}, jobs::JobTriggerKind, @@ -15,6 +12,9 @@ use windmill_common::{ DB, }; use windmill_queue::PushArgsOwned; +use windmill_trigger::trigger_helpers::{ + trigger_runnable_and_wait_for_raw_result_with_error_ctx, TriggerJobArgs, +}; pub mod handler; pub mod listener; @@ -30,11 +30,17 @@ impl TriggerJobArgs for WebsocketTrigger { } } +fn default_filter_logic() -> String { + "and".to_string() +} + #[derive(Debug, Clone, FromRow, Serialize, Deserialize)] pub struct WebsocketConfig { pub url: String, #[serde(default)] pub filters: Vec>>, + #[serde(default = "default_filter_logic")] + pub filter_logic: String, #[serde(skip_serializing_if = "Option::is_none")] pub initial_messages: Option>>>, #[serde(skip_serializing_if = "Option::is_none")] @@ -49,6 +55,8 @@ pub struct WebsocketConfig { pub struct WebsocketConfigRequest { url: String, filters: Vec, + #[serde(default = "default_filter_logic")] + filter_logic: String, initial_messages: Option>, url_runnable_args: Option, can_return_message: bool, diff --git a/backend/windmill-trigger-websocket/src/listener.rs b/backend/windmill-trigger-websocket/src/listener.rs index 8aab61d9c6..205f2fe501 100644 --- a/backend/windmill-trigger-websocket/src/listener.rs +++ b/backend/windmill-trigger-websocket/src/listener.rs @@ -18,7 +18,7 @@ use windmill_common::{ DB, }; use windmill_queue::PushArgsOwned; -use windmill_trigger::filter::{is_value_superset, Filter, JsonFilter}; +use windmill_trigger::filter::{check_filters, Filter}; use windmill_trigger::listener::ListeningTrigger; use windmill_trigger::trigger_helpers::{ trigger_runnable, trigger_runnable_and_wait_for_raw_result, @@ -267,26 +267,8 @@ impl Listener for WebsocketTrigger { match msg { tokio_tungstenite::tungstenite::Message::Text(text) => { tracing::debug!("Received text message from WebSocket {}: {}", url, text); - let mut should_handle = true; - for filter in &filters { - match filter { - Filter::JsonFilter(JsonFilter { key, value }) => { - let mut deserializer = serde_json::Deserializer::from_str(text.as_str()); - should_handle = match is_value_superset(&mut deserializer, key, &value) { - Ok(filter_match) => { - filter_match - }, - Err(err) => { - tracing::warn!("Error deserializing filter for WebSocket {}: {:?}", url, err); - false - } - }; - } - } - if !should_handle { - break; - } - } + let use_or = listening_trigger.trigger_config.filter_logic == "or"; + let should_handle = check_filters(&text, &filters, use_or); if should_handle { 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 a1dd59f848..3c9c857058 100644 --- a/backend/windmill-trigger/src/filter.rs +++ b/backend/windmill-trigger/src/filter.rs @@ -80,6 +80,27 @@ where 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::*; diff --git a/cli/src/guidance/skills.ts b/cli/src/guidance/skills.ts index d1546a1c8d..b66b08f844 100644 --- a/cli/src/guidance/skills.ts +++ b/cli/src/guidance/skills.ts @@ -5820,6 +5820,13 @@ properties: key: type: string value: {} + 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. auto_offset_reset: type: string enum: @@ -6347,6 +6354,13 @@ properties: value: {} description: Array of key-value filters to match incoming messages (only matching messages trigger the script) + 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. initial_messages: type: array items: diff --git a/frontend/src/lib/components/triggers/CaptureWrapper.svelte b/frontend/src/lib/components/triggers/CaptureWrapper.svelte index a6908c6957..0c763987a9 100644 --- a/frontend/src/lib/components/triggers/CaptureWrapper.svelte +++ b/frontend/src/lib/components/triggers/CaptureWrapper.svelte @@ -260,7 +260,7 @@ {hasPreprocessor} {isFlow} {captureLoading} - {triggerDeployed} + groupId={args?.group_id} on:applyArgs on:updateSchema on:addPreprocessor diff --git a/frontend/src/lib/components/triggers/TriggerFilters.svelte b/frontend/src/lib/components/triggers/TriggerFilters.svelte index 1fb4041b48..cb7f198346 100644 --- a/frontend/src/lib/components/triggers/TriggerFilters.svelte +++ b/frontend/src/lib/components/triggers/TriggerFilters.svelte @@ -1,23 +1,45 @@

- Filters will limit the execution of the trigger to only messages that match all criteria.
+ {description}
The JSON filter checks if the value at the key is equal or a superset of the filter value.

+ {#if filters.length > 0} +
+