diff --git a/backend/.sqlx/query-6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389.json b/backend/.sqlx/query-11fd92de8688ef6b4d524aade507850a9cb3e097f2d957219d04ad134a5e0399.json similarity index 80% rename from backend/.sqlx/query-6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389.json rename to backend/.sqlx/query-11fd92de8688ef6b4d524aade507850a9cb3e097f2d957219d04ad134a5e0399.json index a6f1f5f7cc..0ce46a41e1 100644 --- a/backend/.sqlx/query-6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389.json +++ b/backend/.sqlx/query-11fd92de8688ef6b4d524aade507850a9cb3e097f2d957219d04ad134a5e0399.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 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 ", + "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 heartbeat\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, now(), $15, $16, $17, $18\n )\n ", "describe": { "columns": [], "parameters": { @@ -32,10 +32,11 @@ "Varchar", "Varchar", "Jsonb", + "Jsonb", "Jsonb" ] }, "nullable": [] }, - "hash": "6948eb5aabf82f2f4a08dd4410eb472080ecab3ed652912397245e5216ae0389" + "hash": "11fd92de8688ef6b4d524aade507850a9cb3e097f2d957219d04ad134a5e0399" } diff --git a/backend/.sqlx/query-e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9.json b/backend/.sqlx/query-492edd53e497a45c314d41044e85a5e3492227b5049dc6870d60817ad929e7af.json similarity index 84% rename from backend/.sqlx/query-e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9.json rename to backend/.sqlx/query-492edd53e497a45c314d41044e85a5e3492227b5049dc6870d60817ad929e7af.json index 8ff6f2e89c..b5f50c99b0 100644 --- a/backend/.sqlx/query-e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9.json +++ b/backend/.sqlx/query-492edd53e497a45c314d41044e85a5e3492227b5049dc6870d60817ad929e7af.json @@ -1,6 +1,6 @@ { "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 ", + "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 heartbeat = $18\n WHERE\n workspace_id = $13 AND path = $14\n ", "describe": { "columns": [], "parameters": { @@ -21,10 +21,11 @@ "Text", "Varchar", "Jsonb", + "Jsonb", "Jsonb" ] }, "nullable": [] }, - "hash": "e3d4f89ce36337af15d237b543eaca47771b480ff194884f9c947dcaf71d6cf9" + "hash": "492edd53e497a45c314d41044e85a5e3492227b5049dc6870d60817ad929e7af" } diff --git a/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json b/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json index 36ddb8ab9f..713ccb9dd3 100644 --- a/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json +++ b/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json @@ -15,7 +15,7 @@ ] }, "nullable": [ - true + null ] }, "hash": "5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55" diff --git a/backend/migrations/20260402213045_websocket_heartbeat.down.sql b/backend/migrations/20260402213045_websocket_heartbeat.down.sql new file mode 100644 index 0000000000..bf71be0f5c --- /dev/null +++ b/backend/migrations/20260402213045_websocket_heartbeat.down.sql @@ -0,0 +1 @@ +ALTER TABLE websocket_trigger DROP COLUMN heartbeat; diff --git a/backend/migrations/20260402213045_websocket_heartbeat.up.sql b/backend/migrations/20260402213045_websocket_heartbeat.up.sql new file mode 100644 index 0000000000..041b56deed --- /dev/null +++ b/backend/migrations/20260402213045_websocket_heartbeat.up.sql @@ -0,0 +1 @@ +ALTER TABLE websocket_trigger ADD COLUMN heartbeat JSONB NULL; diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 1e3a3a2fba..8f3953f014 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -21809,6 +21809,27 @@ components: google_count: type: number + WebsocketHeartbeat: + type: object + properties: + interval_secs: + type: integer + minimum: 1 + description: Interval in seconds between heartbeat messages + message: + type: string + description: >- + Message to send as heartbeat. Use {{state}} as a placeholder + for a value extracted from incoming messages (see state_field). + state_field: + type: string + description: >- + Optional. Top-level JSON field to extract from incoming messages. + The extracted value replaces {{state}} in the heartbeat message. + required: + - interval_secs + - message + WebsocketTrigger: allOf: - $ref: "#/components/schemas/TriggerExtraProperty" @@ -21862,6 +21883,10 @@ components: can_return_error_result: type: boolean description: If true, error results are sent back through the WebSocket + heartbeat: + $ref: "#/components/schemas/WebsocketHeartbeat" + nullable: true + description: Optional periodic heartbeat message configuration error_handler_path: type: string description: Path to a script or flow to run when the triggered job fails @@ -21930,6 +21955,10 @@ components: can_return_error_result: type: boolean description: If true, error results are sent back through the WebSocket + heartbeat: + $ref: "#/components/schemas/WebsocketHeartbeat" + nullable: true + description: Optional periodic heartbeat message configuration error_handler_path: type: string description: Path to a script or flow to run when the triggered job fails @@ -22005,6 +22034,10 @@ components: can_return_error_result: type: boolean description: If true, error results are sent back through the WebSocket + heartbeat: + $ref: "#/components/schemas/WebsocketHeartbeat" + nullable: true + description: Optional periodic heartbeat message configuration error_handler_path: type: string description: Path to a script or flow to run when the triggered job fails diff --git a/backend/windmill-trigger-websocket/src/handler.rs b/backend/windmill-trigger-websocket/src/handler.rs index 223868e96e..5adfc5a1e2 100644 --- a/backend/windmill-trigger-websocket/src/handler.rs +++ b/backend/windmill-trigger-websocket/src/handler.rs @@ -41,6 +41,7 @@ impl TriggerCrud for WebsocketTrigger { "url_runnable_args", "can_return_message", "can_return_error_result", + "heartbeat", ]; const IS_ALLOWED_ON_CLOUD: bool = false; @@ -68,6 +69,19 @@ impl TriggerCrud for WebsocketTrigger { } } + if let Some(ref hb) = config.heartbeat { + if hb.interval_secs < 1 { + return Err(Error::BadRequest( + "heartbeat interval_secs must be at least 1".to_string(), + )); + } + if hb.message.is_empty() { + return Err(Error::BadRequest( + "heartbeat message cannot be empty".to_string(), + )); + } + } + Ok(()) } @@ -114,9 +128,10 @@ impl TriggerCrud for WebsocketTrigger { edited_at, error_handler_path, error_handler_args, - retry + retry, + heartbeat ) VALUES ( - $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, now(), $15, $16, $17 + $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, now(), $15, $16, $17, $18 ) "#, w_id, @@ -138,7 +153,8 @@ impl TriggerCrud for WebsocketTrigger { resolved_permissioned_as, trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, - trigger.error_handling.retry as _ + trigger.error_handling.retry as _, + trigger.config.heartbeat.map(SqlxJson) as _ ) .execute(&mut *tx) .await?; @@ -193,7 +209,8 @@ impl TriggerCrud for WebsocketTrigger { error = NULL, error_handler_path = $15, error_handler_args = $16, - retry = $17 + retry = $17, + heartbeat = $18 WHERE workspace_id = $13 AND path = $14 ", @@ -217,7 +234,8 @@ impl TriggerCrud for WebsocketTrigger { path, trigger.error_handling.error_handler_path, trigger.error_handling.error_handler_args as _, - trigger.error_handling.retry as _ + trigger.error_handling.retry as _, + trigger.config.heartbeat.map(SqlxJson) as _ ) .execute(&mut *tx) .await?; diff --git a/backend/windmill-trigger-websocket/src/lib.rs b/backend/windmill-trigger-websocket/src/lib.rs index 4fe96bb7c2..bcf4a144db 100644 --- a/backend/windmill-trigger-websocket/src/lib.rs +++ b/backend/windmill-trigger-websocket/src/lib.rs @@ -34,6 +34,14 @@ fn default_filter_logic() -> String { "and".to_string() } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WebsocketHeartbeat { + pub interval_secs: u64, + pub message: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub state_field: Option, +} + #[derive(Debug, Clone, FromRow, Serialize, Deserialize)] pub struct WebsocketConfig { pub url: String, @@ -49,6 +57,8 @@ pub struct WebsocketConfig { pub can_return_message: bool, #[serde(default)] pub can_return_error_result: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub heartbeat: Option>, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -61,6 +71,7 @@ pub struct WebsocketConfigRequest { url_runnable_args: Option, can_return_message: bool, can_return_error_result: bool, + pub heartbeat: Option, } #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/backend/windmill-trigger-websocket/src/listener.rs b/backend/windmill-trigger-websocket/src/listener.rs index 205f2fe501..b729c49a28 100644 --- a/backend/windmill-trigger-websocket/src/listener.rs +++ b/backend/windmill-trigger-websocket/src/listener.rs @@ -216,8 +216,10 @@ impl Listener for WebsocketTrigger { } } + let needs_sender = listening_trigger.trigger_config.can_return_message + || listening_trigger.trigger_config.heartbeat.is_some(); let (return_message_channels, message_sender_handle) = if listening_trigger.trigger_mode - && listening_trigger.trigger_config.can_return_message + && needs_sender { let (send_message_tx, mut rx) = tokio::sync::mpsc::channel::(100); let w_id = listening_trigger.workspace_id.clone(); @@ -243,12 +245,62 @@ impl Listener for WebsocketTrigger { (None, None) }; + // Shared heartbeat state: tracks the latest value of the configured state_field + let heartbeat_state: Arc>> = Arc::new(RwLock::new(None)); + + // Clone send channel for heartbeat use (only when heartbeat is configured) + let heartbeat_tx = if listening_trigger.trigger_config.heartbeat.is_some() { + return_message_channels + .as_ref() + .map(|c| c.send_message_tx.clone()) + } else { + None + }; + tokio::select! { biased; _ = killpill_rx.recv() => { }, _ = self.loop_ping(db, listening_trigger, err_message.clone(), None) => { }, + // Heartbeat timer + _ = async { + if let Some(ref hb) = listening_trigger.trigger_config.heartbeat { + let interval_duration = std::time::Duration::from_secs(hb.interval_secs); + let mut interval = tokio::time::interval(interval_duration); + // Skip the first immediate tick + interval.tick().await; + + loop { + interval.tick().await; + + let msg = if hb.state_field.is_some() { + let state_val = heartbeat_state.read().await; + match &*state_val { + Some(val) => hb.message.replace("{{state}}", &val.to_string()), + None => hb.message.replace("{{state}}", "null"), + } + } else { + hb.message.clone() + }; + + tracing::debug!("Sending heartbeat to WebSocket {}: {}", url, msg); + + if let Some(tx) = heartbeat_tx.as_ref() { + if let Err(err) = tx.send(msg).await { + tracing::error!("Failed to send heartbeat to WebSocket {}: {}", url, err); + break; + } + } else { + tracing::warn!("Heartbeat configured but no send channel available for WebSocket {}", url); + break; + } + } + } else { + futures::future::pending::<()>().await; + } + } => {}, + // Message reader _ = async { let filters: Vec = if listening_trigger.trigger_mode { listening_trigger @@ -267,6 +319,18 @@ impl Listener for WebsocketTrigger { match msg { tokio_tungstenite::tungstenite::Message::Text(text) => { tracing::debug!("Received text message from WebSocket {}: {}", url, text); + + // Extract heartbeat state from every incoming message + if let Some(ref hb) = listening_trigger.trigger_config.heartbeat { + if let Some(ref field) = hb.state_field { + if let Ok(parsed) = serde_json::from_str::(&text) { + if let Some(val) = parsed.get(field) { + *heartbeat_state.write().await = Some(val.clone()); + } + } + } + } + let use_or = listening_trigger.trigger_config.filter_logic == "or"; let should_handle = check_filters(&text, &filters, use_or); if should_handle { @@ -349,7 +413,8 @@ impl Listener for WebsocketTrigger { None => (None, None, None), }; let trigger = TriggerMetadata::new(Some(path.to_owned()), Self::JOB_TRIGGER_KIND); - if *suspended_mode || extra.is_none() { + let can_return_message = trigger_config.can_return_message; + if *suspended_mode || extra.is_none() || !can_return_message { trigger_runnable( db, None, diff --git a/cli/src/guidance/skills.ts b/cli/src/guidance/skills.ts index 2417ed9c12..54e41069fa 100644 --- a/cli/src/guidance/skills.ts +++ b/cli/src/guidance/skills.ts @@ -6479,6 +6479,21 @@ properties: can_return_error_result: type: boolean description: If true, error results are sent back through the WebSocket + heartbeat: + type: object + properties: + interval_secs: + type: integer + minimum: 1 + description: Interval in seconds between heartbeat messages + message: + type: string + description: Message to send as heartbeat. Use {{state}} as a placeholder + for a value extracted from incoming messages (see state_field). + state_field: + type: string + description: Optional. Top-level JSON field to extract from incoming messages. + The extracted value replaces {{state}} in the heartbeat message. error_handler_path: type: string description: Path to a script or flow to run when the triggered job fails diff --git a/frontend/src/lib/components/triggers/websocket/WebsocketTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/websocket/WebsocketTriggerEditorInner.svelte index 3a50c71c20..bd2ea8426e 100644 --- a/frontend/src/lib/components/triggers/websocket/WebsocketTriggerEditorInner.svelte +++ b/frontend/src/lib/components/triggers/websocket/WebsocketTriggerEditorInner.svelte @@ -1,5 +1,6 @@