mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-11 08:07:15 +00:00
triggers improvements (#5074)
* triggers improvements * nits * nits * update ee ref
This commit is contained in:
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE kafka_trigger SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "0ca4365e7144584ef5723db7e133bb42525ff91a734caa87be8b802c6607e6be"
|
||||
}
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE kafka_trigger SET last_server_ping = NULL WHERE workspace_id = $1 AND path = $2 AND server_id IS NULL",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "1c37f91192aa4f535c7fff80fa809260b906dc988f44a0db60952a2bf8b1cdaf"
|
||||
}
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE http_trigger SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "1c94d4f3b90a47b40263c254f85d14cc8b551e9ab0208e3a10ca717a4e85888b"
|
||||
}
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE websocket_trigger SET last_server_ping = NULL WHERE workspace_id = $1 AND path = $2 AND server_id IS NULL",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "51f37d683d5a48b96f6224111639c364ba9c41572eb6799f7c822c66d00d2500"
|
||||
}
|
||||
+2
-2
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "INSERT INTO capture_config\n (workspace_id, path, is_flow, trigger_kind, trigger_config, owner, email)\n VALUES ($1, $2, $3, $4, $5, $6, $7)\n ON CONFLICT (workspace_id, path, is_flow, trigger_kind)\n DO UPDATE SET trigger_config = $5, owner = $6, email = $7, server_id = NULL, last_server_ping = NULL, error = NULL",
|
||||
"query": "INSERT INTO capture_config\n (workspace_id, path, is_flow, trigger_kind, trigger_config, owner, email)\n VALUES ($1, $2, $3, $4, $5, $6, $7)\n ON CONFLICT (workspace_id, path, is_flow, trigger_kind)\n DO UPDATE SET trigger_config = $5, owner = $6, email = $7, server_id = NULL, error = NULL",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
@@ -30,5 +30,5 @@
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "ef299490c4674c4c76e18d84620a74407b78378d66d8a089407998074059e79b"
|
||||
"hash": "62475252dcf54f32433b97ae011daf5d4205d160d2aedf463c7dfe944e93257a"
|
||||
}
|
||||
+14
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM websocket_trigger WHERE workspace_id = $1",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "71fe2f596242c1726af032e8b0aa4894458ac0a11c647a321232eccb6d1d25e7"
|
||||
}
|
||||
+2
-2
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM capture_config WHERE workspace_id = $1",
|
||||
"query": "DELETE FROM http_trigger WHERE workspace_id = $1",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
@@ -10,5 +10,5 @@
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "3c9fc4d8579767f3ce7c3633fca770e6341624e98117d21b9f01e68b4e0ce033"
|
||||
"hash": "829a7108c164fc48387d5103b773eaba0a4dbdc3cbf2fd5bf0ac10c86290b3bc"
|
||||
}
|
||||
+1
-1
@@ -5,7 +5,7 @@
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "bool",
|
||||
"name": "?column?",
|
||||
"type_info": "Bool"
|
||||
}
|
||||
],
|
||||
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE capture_config SET last_server_ping = NULL WHERE workspace_id = $1 AND path = $2 AND is_flow = $3 AND trigger_kind = 'nats' AND server_id IS NULL",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text",
|
||||
"Bool"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "93d64930c74ccb1abdc9bda8287540a354f8eec49913aaf07fc1e737e7c93330"
|
||||
}
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE websocket_trigger SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "98880172e7d163e0db06eea2b9008c5082a3e34d8d294b341f01f549be6c7333"
|
||||
}
|
||||
+1
-1
@@ -5,7 +5,7 @@
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "bool",
|
||||
"name": "?column?",
|
||||
"type_info": "Bool"
|
||||
}
|
||||
],
|
||||
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE capture_config SET last_server_ping = NULL WHERE workspace_id = $1 AND path = $2 AND is_flow = $3 AND trigger_kind = 'kafka' AND server_id IS NULL",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text",
|
||||
"Bool"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "b223c56f55a138abef8c8cb2df40472a87b7e267a143709728b09ac063c33b96"
|
||||
}
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE capture_config SET last_server_ping = NULL WHERE workspace_id = $1 AND path = $2 AND is_flow = $3 AND trigger_kind = 'websocket' AND server_id IS NULL",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text",
|
||||
"Bool"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "cba3bfb174829ee3b08ea195831fdba4335ca221a5fd25c72519ac376c593e44"
|
||||
}
|
||||
+1
-1
@@ -5,7 +5,7 @@
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "bool",
|
||||
"name": "?column?",
|
||||
"type_info": "Bool"
|
||||
}
|
||||
],
|
||||
|
||||
+1
-1
@@ -5,7 +5,7 @@
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "bool",
|
||||
"name": "?column?",
|
||||
"type_info": "Bool"
|
||||
}
|
||||
],
|
||||
|
||||
+1
-1
@@ -5,7 +5,7 @@
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "bool",
|
||||
"name": "?column?",
|
||||
"type_info": "Bool"
|
||||
}
|
||||
],
|
||||
|
||||
+14
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM kafka_trigger WHERE workspace_id = $1",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "de816241885fee0f5f8f99bd34e79214f9ec492d458e72ecb032a6a92a908ca2"
|
||||
}
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE nats_trigger SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "e2b17a69501978f2617ff4d20198111e3dab821c25d3164db7983b3192d149a0"
|
||||
}
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE nats_trigger SET last_server_ping = NULL WHERE workspace_id = $1 AND path = $2 AND server_id IS NULL",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "e38a66d15382703b99054a9b609af9f3854bb65393d33b265476c0e68aec4f61"
|
||||
}
|
||||
+1
-1
@@ -5,7 +5,7 @@
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "bool",
|
||||
"name": "?column?",
|
||||
"type_info": "Bool"
|
||||
}
|
||||
],
|
||||
|
||||
@@ -1 +1 @@
|
||||
a515264a3c1c138e289da1ef3bb21d260aead486
|
||||
fb5e285c316f10fa2188318585573c7e85999784
|
||||
@@ -0,0 +1,16 @@
|
||||
-- Add down migration script here
|
||||
DROP POLICY admin_policy ON capture_config;
|
||||
DROP POLICY see_folder_extra_perms_user_select ON capture_config;
|
||||
DROP POLICY see_folder_extra_perms_user_insert ON capture_config;
|
||||
DROP POLICY see_folder_extra_perms_user_update ON capture_config;
|
||||
DROP POLICY see_folder_extra_perms_user_delete ON capture_config;
|
||||
DROP POLICY see_own ON capture_config;
|
||||
DROP POLICY see_member ON capture_config;
|
||||
|
||||
|
||||
DROP POLICY see_folder_extra_perms_user_select ON capture;
|
||||
DROP POLICY see_folder_extra_perms_user_insert ON capture;
|
||||
DROP POLICY see_folder_extra_perms_user_update ON capture;
|
||||
DROP POLICY see_folder_extra_perms_user_delete ON capture;
|
||||
DROP POLICY see_own ON capture;
|
||||
DROP POLICY see_member ON capture;
|
||||
@@ -0,0 +1,30 @@
|
||||
-- capture config
|
||||
CREATE POLICY admin_policy ON capture_config FOR ALL TO windmill_admin USING (true);
|
||||
CREATE POLICY see_folder_extra_perms_user_select ON capture_config FOR SELECT TO windmill_user
|
||||
USING (SPLIT_PART(capture_config.path, '/', 1) = 'f' AND SPLIT_PART(capture_config.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_read'), ',')::text[]));
|
||||
CREATE POLICY see_folder_extra_perms_user_insert ON capture_config FOR INSERT TO windmill_user
|
||||
WITH CHECK (SPLIT_PART(capture_config.path, '/', 1) = 'f' AND SPLIT_PART(capture_config.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
|
||||
CREATE POLICY see_folder_extra_perms_user_update ON capture_config FOR UPDATE TO windmill_user
|
||||
USING (SPLIT_PART(capture_config.path, '/', 1) = 'f' AND SPLIT_PART(capture_config.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
|
||||
CREATE POLICY see_folder_extra_perms_user_delete ON capture_config FOR DELETE TO windmill_user
|
||||
USING (SPLIT_PART(capture_config.path, '/', 1) = 'f' AND SPLIT_PART(capture_config.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
|
||||
CREATE POLICY see_own ON capture_config FOR ALL TO windmill_user
|
||||
USING (SPLIT_PART(capture_config.path, '/', 1) = 'u' AND SPLIT_PART(capture_config.path, '/', 2) = current_setting('session.user'));
|
||||
CREATE POLICY see_member ON capture_config FOR ALL TO windmill_user
|
||||
USING (SPLIT_PART(capture_config.path, '/', 1) = 'g' AND SPLIT_PART(capture_config.path, '/', 2) = any(regexp_split_to_array(current_setting('session.groups'), ',')::text[]));
|
||||
|
||||
|
||||
-- capture
|
||||
CREATE POLICY see_folder_extra_perms_user_select ON capture FOR SELECT TO windmill_user
|
||||
USING (SPLIT_PART(capture.path, '/', 1) = 'f' AND SPLIT_PART(capture.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_read'), ',')::text[]));
|
||||
CREATE POLICY see_folder_extra_perms_user_insert ON capture FOR INSERT TO windmill_user
|
||||
WITH CHECK (SPLIT_PART(capture.path, '/', 1) = 'f' AND SPLIT_PART(capture.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
|
||||
CREATE POLICY see_folder_extra_perms_user_update ON capture FOR UPDATE TO windmill_user
|
||||
USING (SPLIT_PART(capture.path, '/', 1) = 'f' AND SPLIT_PART(capture.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
|
||||
CREATE POLICY see_folder_extra_perms_user_delete ON capture FOR DELETE TO windmill_user
|
||||
USING (SPLIT_PART(capture.path, '/', 1) = 'f' AND SPLIT_PART(capture.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
|
||||
CREATE POLICY see_own ON capture FOR ALL TO windmill_user
|
||||
USING (SPLIT_PART(capture.path, '/', 1) = 'u' AND SPLIT_PART(capture.path, '/', 2) = current_setting('session.user'));
|
||||
CREATE POLICY see_member ON capture FOR ALL TO windmill_user
|
||||
USING (SPLIT_PART(capture.path, '/', 1) = 'g' AND SPLIT_PART(capture.path, '/', 2) = any(regexp_split_to_array(current_setting('session.groups'), ',')::text[]));
|
||||
|
||||
@@ -220,7 +220,7 @@ async fn set_config(
|
||||
(workspace_id, path, is_flow, trigger_kind, trigger_config, owner, email)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7)
|
||||
ON CONFLICT (workspace_id, path, is_flow, trigger_kind)
|
||||
DO UPDATE SET trigger_config = $5, owner = $6, email = $7, server_id = NULL, last_server_ping = NULL, error = NULL",
|
||||
DO UPDATE SET trigger_config = $5, owner = $6, email = $7, server_id = NULL, error = NULL",
|
||||
&w_id,
|
||||
&nc.path,
|
||||
nc.is_flow,
|
||||
|
||||
@@ -250,7 +250,7 @@ async fn update_websocket_trigger(
|
||||
.map(SqlxJson)
|
||||
.collect_vec();
|
||||
|
||||
// important to update server_id, last_server_ping and error to NULL to stop current websocket listener
|
||||
// important to update server_id to NULL to stop current websocket listener
|
||||
sqlx::query!(
|
||||
"UPDATE websocket_trigger SET url = $1, script_path = $2, path = $3, is_flow = $4, filters = $5, initial_messages = $6, url_runnable_args = $7, edited_by = $8, email = $9, edited_at = now(), server_id = NULL, error = NULL
|
||||
WHERE workspace_id = $10 AND path = $11",
|
||||
@@ -671,6 +671,12 @@ impl WebsocketTrigger {
|
||||
).fetch_optional(db).await {
|
||||
Ok(updated) => {
|
||||
if updated.flatten().is_none() {
|
||||
// allow faster restart of websocket trigger
|
||||
sqlx::query!(
|
||||
"UPDATE websocket_trigger SET last_server_ping = NULL WHERE workspace_id = $1 AND path = $2 AND server_id IS NULL",
|
||||
self.workspace_id,
|
||||
self.path,
|
||||
).execute(db).await.ok();
|
||||
tracing::info!("Websocket {} changed, disabled, or deleted, stopping...", self.url);
|
||||
return None;
|
||||
}
|
||||
@@ -875,6 +881,13 @@ impl CaptureConfigForWebsocket {
|
||||
).fetch_optional(db).await {
|
||||
Ok(updated) => {
|
||||
if updated.flatten().is_none() {
|
||||
// allow faster restart of websocket capture
|
||||
sqlx::query!(
|
||||
"UPDATE capture_config SET last_server_ping = NULL WHERE workspace_id = $1 AND path = $2 AND is_flow = $3 AND trigger_kind = 'websocket' AND server_id IS NULL",
|
||||
self.workspace_id,
|
||||
self.path,
|
||||
self.is_flow,
|
||||
).execute(db).await.ok();
|
||||
tracing::info!("Websocket capture {} changed, disabled, or deleted, stopping...", self.trigger_config.url);
|
||||
return None;
|
||||
}
|
||||
|
||||
@@ -111,6 +111,38 @@ pub(crate) async fn change_workspace_id(
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"UPDATE http_trigger SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
&rw.new_id,
|
||||
&old_id
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"UPDATE websocket_trigger SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
&rw.new_id,
|
||||
&old_id
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"UPDATE kafka_trigger SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
&rw.new_id,
|
||||
&old_id
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"UPDATE nats_trigger SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
&rw.new_id,
|
||||
&old_id
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"UPDATE completed_job SET workspace_id = $1 WHERE workspace_id = $2",
|
||||
&rw.new_id,
|
||||
@@ -404,9 +436,9 @@ pub(crate) async fn delete_workspace(
|
||||
sqlx::query!("DELETE FROM capture WHERE workspace_id = $1", &w_id)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
sqlx::query!("DELETE FROM capture_config WHERE workspace_id = $1", &w_id)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
// capture_config has on delete cascade
|
||||
|
||||
sqlx::query!("DELETE FROM draft WHERE workspace_id = $1", &w_id)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
@@ -497,6 +529,23 @@ pub(crate) async fn delete_workspace(
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!("DELETE FROM http_trigger WHERE workspace_id = $1", &w_id)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"DELETE FROM websocket_trigger WHERE workspace_id = $1",
|
||||
&w_id
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!("DELETE FROM kafka_trigger WHERE workspace_id = $1", &w_id)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
// NATS triggers have on delete cascade
|
||||
|
||||
sqlx::query!("DELETE FROM workspace WHERE id = $1", &w_id)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
|
||||
@@ -88,20 +88,15 @@
|
||||
}
|
||||
}
|
||||
|
||||
let ready = false
|
||||
function setDefaultArgs(captureConfigs: { [key: string]: CaptureConfig }) {
|
||||
if (captureType in captureConfigs) {
|
||||
const triggerConfig = captureConfigs[captureType].trigger_config
|
||||
args = isObject(triggerConfig) ? triggerConfig : {}
|
||||
} else if (captureType === 'kafka') {
|
||||
args = {
|
||||
...args,
|
||||
brokers: [''],
|
||||
topics: [''],
|
||||
group_id: `windmill_consumer-${$workspaceStore}-${path.replaceAll('/', '__')}`
|
||||
}
|
||||
} else {
|
||||
args = {}
|
||||
}
|
||||
ready = true
|
||||
}
|
||||
|
||||
onDestroy(() => {
|
||||
@@ -165,113 +160,118 @@
|
||||
$: args && (captureActive = false)
|
||||
</script>
|
||||
|
||||
<div class="flex flex-col gap-4 w-full">
|
||||
{#if cloudDisabled}
|
||||
<Alert title="Not compatible with multi-tenant cloud" type="warning" size="xs">
|
||||
{capitalize(captureType)} triggers are disabled in the multi-tenant cloud.
|
||||
</Alert>
|
||||
{:else if captureType === 'websocket'}
|
||||
<WebsocketEditorConfigSection
|
||||
can_write={true}
|
||||
headless={true}
|
||||
bind:url={args.url}
|
||||
bind:url_runnable_args={args.url_runnable_args}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'webhook'}
|
||||
<WebhooksConfigSection
|
||||
{isFlow}
|
||||
{path}
|
||||
hash={data?.hash}
|
||||
token={data?.token}
|
||||
{args}
|
||||
scopes={data?.scopes}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'http'}
|
||||
<RouteEditorConfigSection
|
||||
{showCapture}
|
||||
can_write={true}
|
||||
bind:args
|
||||
headless
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'email'}
|
||||
<EmailTriggerConfigSection
|
||||
hash={data?.hash}
|
||||
token={data?.token}
|
||||
{path}
|
||||
{isFlow}
|
||||
userSettings={data?.userSettings}
|
||||
emailDomain={data?.emailDomain}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'kafka'}
|
||||
<KafkaTriggersConfigSection
|
||||
headless={true}
|
||||
bind:args
|
||||
staticInputDisabled={false}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'nats'}
|
||||
<NatsTriggersConfigSection
|
||||
headless={true}
|
||||
bind:args
|
||||
{path}
|
||||
staticInputDisabled={false}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
/>
|
||||
{/if}
|
||||
</div>
|
||||
{#key ready}
|
||||
<div class="flex flex-col gap-4 w-full">
|
||||
{#if cloudDisabled}
|
||||
<Alert title="Not compatible with multi-tenant cloud" type="warning" size="xs">
|
||||
{capitalize(captureType)} triggers are disabled in the multi-tenant cloud.
|
||||
</Alert>
|
||||
{:else if captureType === 'websocket'}
|
||||
<WebsocketEditorConfigSection
|
||||
can_write={true}
|
||||
headless={true}
|
||||
bind:url={args.url}
|
||||
bind:url_runnable_args={args.url_runnable_args}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'webhook'}
|
||||
<WebhooksConfigSection
|
||||
{isFlow}
|
||||
{path}
|
||||
hash={data?.hash}
|
||||
token={data?.token}
|
||||
{args}
|
||||
scopes={data?.scopes}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'http'}
|
||||
<RouteEditorConfigSection
|
||||
{isFlow}
|
||||
{path}
|
||||
{showCapture}
|
||||
can_write={true}
|
||||
bind:args
|
||||
headless
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'email'}
|
||||
<EmailTriggerConfigSection
|
||||
hash={data?.hash}
|
||||
token={data?.token}
|
||||
{path}
|
||||
{isFlow}
|
||||
userSettings={data?.userSettings}
|
||||
emailDomain={data?.emailDomain}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'kafka'}
|
||||
<KafkaTriggersConfigSection
|
||||
headless={true}
|
||||
{path}
|
||||
bind:args
|
||||
staticInputDisabled={false}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
on:testWithArgs
|
||||
/>
|
||||
{:else if captureType === 'nats'}
|
||||
<NatsTriggersConfigSection
|
||||
headless={true}
|
||||
bind:args
|
||||
{path}
|
||||
staticInputDisabled={false}
|
||||
{showCapture}
|
||||
{captureInfo}
|
||||
bind:captureTable
|
||||
on:applyArgs
|
||||
on:updateSchema
|
||||
on:addPreprocessor
|
||||
on:captureToggle={() => {
|
||||
handleCapture()
|
||||
}}
|
||||
/>
|
||||
{/if}
|
||||
</div>
|
||||
{/key}
|
||||
|
||||
@@ -218,6 +218,7 @@
|
||||
<KafkaTriggersConfigSection
|
||||
bind:args
|
||||
bind:isValid
|
||||
{path}
|
||||
defaultValues={useDefaultValues() ? defaultValues : undefined}
|
||||
/>
|
||||
|
||||
|
||||
@@ -8,7 +8,9 @@
|
||||
import SchemaForm from '../SchemaForm.svelte'
|
||||
import CaptureSection, { type CaptureInfo } from './CaptureSection.svelte'
|
||||
import CaptureTable from './CaptureTable.svelte'
|
||||
import { workspaceStore } from '$lib/stores'
|
||||
|
||||
export let path: string
|
||||
export let defaultValues: Record<string, any> | undefined = undefined
|
||||
export let headless: boolean = false
|
||||
export let args: Record<string, any> = {}
|
||||
@@ -30,7 +32,8 @@
|
||||
type: 'string'
|
||||
},
|
||||
nullable: false,
|
||||
title: 'Brokers'
|
||||
title: 'Brokers',
|
||||
default: ['']
|
||||
},
|
||||
security: {
|
||||
type: 'object',
|
||||
@@ -120,12 +123,13 @@
|
||||
type: 'string'
|
||||
},
|
||||
nullable: false,
|
||||
title: 'Topics'
|
||||
title: 'Topics',
|
||||
default: ['']
|
||||
},
|
||||
group_id: {
|
||||
type: 'string',
|
||||
title: 'Group ID',
|
||||
pattern: '^[a-zA-Z0-9-_.]+$',
|
||||
pattern: '^((\\$var:[\\w-\\/]+)|[a-zA-Z0-9-_.]+)$',
|
||||
customErrorMessage: 'Invalid group ID'
|
||||
}
|
||||
},
|
||||
@@ -147,7 +151,15 @@
|
||||
args.topics.length > 0 &&
|
||||
args.topics.every((b) => /^[a-zA-Z0-9-_.]+$/.test(b))
|
||||
|
||||
$: args.kafka_resource_path && (selected = 'resource')
|
||||
$: usingResource = !!args.kafka_resource_path
|
||||
$: usingResource && (selected = 'resource')
|
||||
|
||||
function setGroupId() {
|
||||
if (!args.group_id) {
|
||||
args.group_id = `windmill_consumer-${$workspaceStore}-${path.replaceAll('/', '__')}`
|
||||
}
|
||||
}
|
||||
$: path && setGroupId()
|
||||
</script>
|
||||
|
||||
<div>
|
||||
|
||||
@@ -153,7 +153,7 @@
|
||||
stream_name: {
|
||||
type: 'string',
|
||||
title: 'Stream name',
|
||||
pattern: '^[a-zA-Z0-9-_]+$',
|
||||
pattern: '^((\\$var:[\\w-\\/]+)|[a-zA-Z0-9-_]+)$',
|
||||
customErrorMessage: 'Invalid stream name',
|
||||
showExpr: 'fields.use_jetstream',
|
||||
description:
|
||||
@@ -162,7 +162,7 @@
|
||||
consumer_name: {
|
||||
type: 'string',
|
||||
title: 'Consumer name',
|
||||
pattern: '^[a-zA-Z0-9-_]+$',
|
||||
pattern: '^((\\$var:[\\w-\\/]+)|[a-zA-Z0-9-_]+)$',
|
||||
customErrorMessage: 'Invalid consumer name',
|
||||
showExpr: 'fields.use_jetstream',
|
||||
description:
|
||||
|
||||
@@ -17,6 +17,8 @@
|
||||
import CaptureTable from './CaptureTable.svelte'
|
||||
import ClipboardPanel from '../details/ClipboardPanel.svelte'
|
||||
|
||||
export let isFlow: boolean
|
||||
export let path: string
|
||||
export let args: Record<string, any> = { route_path: '', http_method: 'get' }
|
||||
export let dirtyRoutePath: boolean = false
|
||||
export let route_path = ''
|
||||
@@ -59,6 +61,10 @@
|
||||
})
|
||||
}
|
||||
|
||||
$: captureURL = `${location.origin}${base}/api/w/${$workspaceStore}/capture_u/http/${
|
||||
isFlow ? 'flow' : 'script'
|
||||
}/${path.replaceAll('/', '.')}/${route_path}`
|
||||
|
||||
function getHttpRoute(route_path: string | undefined) {
|
||||
return `${location.origin}${base}/api/r/${
|
||||
isCloudHosted() ? $workspaceStore + '/' : ''
|
||||
@@ -94,14 +100,14 @@
|
||||
bind:captureTable
|
||||
>
|
||||
<Label label="URL">
|
||||
<ClipboardPanel content={fullRoute} disabled={!captureInfo.active} />
|
||||
<ClipboardPanel content={captureURL} disabled={!captureInfo.active} />
|
||||
</Label>
|
||||
|
||||
<Label label="Example cUrl">
|
||||
<CopyableCodeBlock
|
||||
disabled={!captureInfo.active}
|
||||
code={`curl \\
|
||||
-X POST ${fullRoute} \\
|
||||
-X POST ${captureURL} \\
|
||||
-H 'Content-Type: application/json' \\
|
||||
-d '{"foo": 42}'`}
|
||||
language={bash}
|
||||
|
||||
@@ -224,6 +224,8 @@
|
||||
</div>
|
||||
|
||||
<RouteEditorConfigSection
|
||||
isFlow={is_flow}
|
||||
{path}
|
||||
bind:route_path
|
||||
bind:args
|
||||
bind:isValid
|
||||
|
||||
@@ -41,7 +41,14 @@
|
||||
showCapture={false}
|
||||
/>
|
||||
{:else if triggerType === 'http'}
|
||||
<RouteEditorConfigSection showCapture={false} can_write={true} bind:args headless />
|
||||
<RouteEditorConfigSection
|
||||
showCapture={false}
|
||||
can_write={true}
|
||||
bind:args
|
||||
headless
|
||||
{isFlow}
|
||||
{path}
|
||||
/>
|
||||
{:else if triggerType === 'email'}
|
||||
<EmailTriggerConfigSection
|
||||
hash={data?.hash}
|
||||
@@ -52,7 +59,7 @@
|
||||
emailDomain={data?.emailDomain}
|
||||
/>
|
||||
{:else if triggerType === 'kafka'}
|
||||
<KafkaTriggersConfigSection headless={true} bind:args staticInputDisabled={false} />
|
||||
<KafkaTriggersConfigSection headless={true} bind:args staticInputDisabled={false} {path} />
|
||||
{:else if triggerType === 'nats'}
|
||||
<NatsTriggersConfigSection headless={true} bind:args staticInputDisabled={false} {path} />
|
||||
{/if}
|
||||
|
||||
@@ -67,7 +67,7 @@
|
||||
urlError = ''
|
||||
}
|
||||
} else if (!url || /^(ws:|wss:)\/\/[^\s]+$/.test(url) === false) {
|
||||
urlError = 'Invalid websocket URL'
|
||||
urlError = 'Websocket URL must start with ws:// or wss://'
|
||||
} else {
|
||||
urlError = ''
|
||||
}
|
||||
@@ -171,6 +171,7 @@
|
||||
on:input={() => {
|
||||
dirtyUrl = true
|
||||
}}
|
||||
placeholder="ws://example.com"
|
||||
class={urlError === ''
|
||||
? ''
|
||||
: 'border border-red-700 bg-red-100 border-opacity-30 focus:border-red-700 focus:border-opacity-30 focus-visible:ring-red-700 focus-visible:ring-opacity-25 focus-visible:border-red-700'}
|
||||
|
||||
@@ -81,7 +81,7 @@
|
||||
})
|
||||
} catch (err) {
|
||||
sendUserToast(
|
||||
`Cannot ` + (enabled ? 'enable' : 'disable') + ` nats trigger: ${err.body}`,
|
||||
`Cannot ` + (enabled ? 'enable' : 'disable') + ` NATS trigger: ${err.body}`,
|
||||
true
|
||||
)
|
||||
} finally {
|
||||
@@ -211,7 +211,7 @@
|
||||
<CenteredPage>
|
||||
<PageHeader
|
||||
title="NATS triggers"
|
||||
tooltip="Windmill can consume nats events and trigger scripts or flows based on them."
|
||||
tooltip="Windmill can consume NATS events and trigger scripts or flows based on them."
|
||||
>
|
||||
<Button size="md" startIcon={{ icon: Plus }} on:click={() => natsTriggerEditor.openNew(false)}>
|
||||
New NATS trigger
|
||||
@@ -258,7 +258,7 @@
|
||||
<Skeleton layout={[[6], 0.4]} />
|
||||
{/each}
|
||||
{:else if !triggers?.length}
|
||||
<div class="text-center text-sm text-tertiary mt-2"> No nats triggers </div>
|
||||
<div class="text-center text-sm text-tertiary mt-2"> No NATS triggers </div>
|
||||
{:else if items?.length}
|
||||
<div class="border rounded-md divide-y">
|
||||
{#each items.slice(0, nbDisplayed) as { path, edited_by, edited_at, script_path, is_flow, nats_resource_path, subjects, extra_perms, canWrite, marked, server_id, error, last_server_ping, enabled } (path)}
|
||||
|
||||
Reference in New Issue
Block a user