feat: add application-level heartbeat support for websocket triggers (#8686)

* feat: add application-level heartbeat support for websocket triggers

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

* Update SQLx metadata

* chore: regenerate auto-generated schema and skill files

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

* fix: handle missing heartbeat channel gracefully, fix TextInput props

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

* refactor: only clone heartbeat sender when heartbeat is configured

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
This commit is contained in:
Alexander Petric
2026-04-03 07:31:08 -04:00
committed by GitHub
parent cdf3c29664
commit 5b7fa63bf1
13 changed files with 257 additions and 12 deletions

View File

@@ -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"
}

View File

@@ -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"
}

View File

@@ -15,7 +15,7 @@
]
},
"nullable": [
true
null
]
},
"hash": "5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55"

View File

@@ -0,0 +1 @@
ALTER TABLE websocket_trigger DROP COLUMN heartbeat;

View File

@@ -0,0 +1 @@
ALTER TABLE websocket_trigger ADD COLUMN heartbeat JSONB NULL;

View File

@@ -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

View File

@@ -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?;

View File

@@ -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<String>,
}
#[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<SqlxJson<WebsocketHeartbeat>>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -61,6 +71,7 @@ pub struct WebsocketConfigRequest {
url_runnable_args: Option<serde_json::Value>,
can_return_message: bool,
can_return_error_result: bool,
pub heartbeat: Option<WebsocketHeartbeat>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]

View File

@@ -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::<String>(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<RwLock<Option<serde_json::Value>>> = 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<Filter> = 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::<serde_json::Value>(&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,

View File

@@ -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

View File

@@ -1,5 +1,6 @@
<script lang="ts">
import { Alert, Button } from '$lib/components/common'
import TextInput from '$lib/components/text_input/TextInput.svelte'
import Drawer from '$lib/components/common/drawer/Drawer.svelte'
import DrawerContent from '$lib/components/common/drawer/DrawerContent.svelte'
import Path from '$lib/components/Path.svelte'
@@ -13,6 +14,7 @@
type Script,
type ScriptArgs,
type WebsocketTriggerInitialMessage,
type WebsocketHeartbeat,
type Retry,
type ErrorHandler,
type TriggerMode
@@ -101,6 +103,10 @@
let url_runnable_args: Record<string, any> | undefined = $state({})
let can_return_message = $state(false)
let can_return_error_result = $state(false)
let heartbeat_enabled = $state(false)
let heartbeat_interval_secs = $state(41)
let heartbeat_message = $state('')
let heartbeat_state_field = $state('')
let dirtyPath = $state(false)
let can_write = $state(true)
let drawerLoading = $state(true)
@@ -235,6 +241,11 @@
url_runnable_args = cfg?.url_runnable_args
can_return_message = cfg?.can_return_message
can_return_error_result = cfg?.can_return_error_result
const hb = cfg?.heartbeat as WebsocketHeartbeat | undefined | null
heartbeat_enabled = !!hb
heartbeat_interval_secs = hb?.interval_secs ?? 41
heartbeat_message = hb?.message ?? ''
heartbeat_state_field = hb?.state_field ?? ''
can_write = canWrite(path, cfg?.extra_perms, $userStore)
error_handler_path = cfg?.error_handler_path
error_handler_args = cfg?.error_handler_args ?? {}
@@ -258,6 +269,13 @@
url_runnable_args,
can_return_message,
can_return_error_result,
heartbeat: heartbeat_enabled
? {
interval_secs: heartbeat_interval_secs,
message: heartbeat_message,
...(heartbeat_state_field ? { state_field: heartbeat_state_field } : {})
}
: undefined,
error_handler_path,
error_handler_args,
retry,
@@ -708,6 +726,71 @@
</div>
</Section>
<Section label="Heartbeat" collapsable collapsed={!heartbeat_enabled}>
{#snippet header()}
{#if heartbeat_enabled}
<span class="text-2xs text-tertiary ml-2">every {heartbeat_interval_secs}s</span>
{/if}
{/snippet}
<div class="flex flex-col gap-4">
<p class="text-xs text-tertiary">
Send periodic application-level heartbeat messages to keep the connection alive.
Required for protocols like Discord Gateway, STOMP, etc.
</p>
<Toggle
checked={heartbeat_enabled}
on:change={() => {
heartbeat_enabled = !heartbeat_enabled
}}
options={{
right: 'Enable heartbeat'
}}
disabled={!can_write}
/>
{#if heartbeat_enabled}
<Label label="Interval (seconds)">
<TextInput
bind:value={heartbeat_interval_secs}
inputProps={{ type: 'number', placeholder: '41', disabled: !can_write, min: 1 }}
/>
</Label>
<Label label="Message">
<svelte:boundary>
<textarea
class="textarea textarea-sm w-full font-mono text-xs"
rows={3}
bind:value={heartbeat_message}
placeholder={'{"op": 1, "d": {{state}}}'}
disabled={!can_write}
></textarea>
</svelte:boundary>
<p class="text-2xs text-tertiary mt-1">
Use <code>{'{{state}}'}</code> as a placeholder for a value extracted from incoming messages.
</p>
</Label>
<Label label="State field (optional)">
<TextInput
bind:value={heartbeat_state_field}
inputProps={{
placeholder: 'e.g. s (for Discord sequence number)',
disabled: !can_write
}}
/>
<p class="text-2xs text-tertiary mt-1">
Top-level JSON field to extract from incoming messages. Replaces <code
>{'{{state}}'}</code
> in the heartbeat message.
</p>
</Label>
{/if}
</div>
</Section>
<Section label="Advanced" collapsable>
{#snippet header()}
<TriggerAdvancedBadges

View File

@@ -30,6 +30,7 @@ export async function saveWebsocketTriggerFromCfg(
url_runnable_args: triggerCfg.url_runnable_args,
can_return_message: triggerCfg.can_return_message,
can_return_error_result: triggerCfg.can_return_error_result,
heartbeat: triggerCfg.heartbeat,
...errorHandlerAndRetries,
permissioned_as: triggerCfg.permissioned_as,
preserve_permissioned_as: triggerCfg.preserve_permissioned_as

View File

@@ -46,6 +46,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