feat: websocket triggers (#4505)

* feat: websocket triggers

* fix

* frontend nits

* nits

* new trigger UI

* sqlx
This commit is contained in:
HugoCasa
2024-10-18 00:21:52 +02:00
committed by GitHub
parent 2953154a25
commit b0fbcd8305
48 changed files with 2594 additions and 95 deletions
@@ -0,0 +1,24 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE websocket_trigger SET server_id = $1, last_server_ping = now() WHERE enabled IS TRUE AND workspace_id = $2 AND path = $3 AND (server_id IS NULL OR last_server_ping IS NULL OR last_server_ping < now() - interval '15 seconds') RETURNING true",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "?column?",
"type_info": "Bool"
}
],
"parameters": {
"Left": [
"Varchar",
"Text",
"Text"
]
},
"nullable": [
null
]
},
"hash": "02424907504848e983bfa89eec343061932dc5b4b17cf13d5cf8d833aedbe6d5"
}
@@ -0,0 +1,26 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE websocket_trigger SET enabled = $1, email = $2, edited_by = $3, edited_at = now(), server_id = NULL, last_server_ping = NULL, error = NULL\n WHERE path = $4 AND workspace_id = $5 RETURNING 1",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "?column?",
"type_info": "Int4"
}
],
"parameters": {
"Left": [
"Bool",
"Varchar",
"Varchar",
"Text",
"Text"
]
},
"nullable": [
null
]
},
"hash": "12a86755706ce030a0a9142da78d035d0c7d361240b60005cc864e4345eb0bc7"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT EXISTS(SELECT 1 FROM http_trigger WHERE workspace_id = $1)",
"query": "SELECT EXISTS(SELECT 1 FROM websocket_trigger WHERE path = $1 AND workspace_id = $2)",
"describe": {
"columns": [
{
@@ -11,6 +11,7 @@
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
@@ -18,5 +19,5 @@
null
]
},
"hash": "9ebb9c16948a695a053068d9a1df0152691271ac33d8c82db47be11746ccbcef"
"hash": "2238aaed46031f71b17c1cabcef935a0d3b2ca6d057f446ed9b02d4386e2ddd9"
}
@@ -0,0 +1,24 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COUNT(*) FROM websocket_trigger WHERE script_path = $1 AND is_flow = $2 AND workspace_id = $3",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "count",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Text",
"Bool",
"Text"
]
},
"nullable": [
null
]
},
"hash": "2e6165543e34216dfaedf6e10729f733cff812e62dea1c1bce8cefd3a3979b14"
}
@@ -0,0 +1,108 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO websocket_trigger (workspace_id, path, url, script_path, is_flow, enabled, filters, edited_by, email, edited_at) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, now()) RETURNING *",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "path",
"type_info": "Varchar"
},
{
"ordinal": 1,
"name": "url",
"type_info": "Varchar"
},
{
"ordinal": 2,
"name": "script_path",
"type_info": "Varchar"
},
{
"ordinal": 3,
"name": "is_flow",
"type_info": "Bool"
},
{
"ordinal": 4,
"name": "workspace_id",
"type_info": "Varchar"
},
{
"ordinal": 5,
"name": "edited_by",
"type_info": "Varchar"
},
{
"ordinal": 6,
"name": "email",
"type_info": "Varchar"
},
{
"ordinal": 7,
"name": "edited_at",
"type_info": "Timestamptz"
},
{
"ordinal": 8,
"name": "extra_perms",
"type_info": "Jsonb"
},
{
"ordinal": 9,
"name": "server_id",
"type_info": "Varchar"
},
{
"ordinal": 10,
"name": "last_server_ping",
"type_info": "Timestamptz"
},
{
"ordinal": 11,
"name": "error",
"type_info": "Text"
},
{
"ordinal": 12,
"name": "enabled",
"type_info": "Bool"
},
{
"ordinal": 13,
"name": "filters",
"type_info": "JsonbArray"
}
],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar",
"Varchar",
"Bool",
"Bool",
"JsonbArray",
"Varchar",
"Varchar"
]
},
"nullable": [
false,
false,
false,
false,
false,
false,
false,
false,
false,
true,
true,
true,
false,
false
]
},
"hash": "4e9668a46bad9e82baa51422946d373b18b6577198df7545c94bd19be3446775"
}
@@ -0,0 +1,22 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE websocket_trigger SET url = $1, script_path = $2, path = $3, is_flow = $4, filters = $5, edited_by = $6, email = $7, edited_at = now(), server_id = NULL, last_server_ping = NULL, error = NULL\n WHERE workspace_id = $8 AND path = $9",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar",
"Bool",
"JsonbArray",
"Varchar",
"Varchar",
"Text",
"Text"
]
},
"nullable": []
},
"hash": "acbf74cf3302bfcf7615285070d3f8958932bb8a2dda715f1b9152ab44442780"
}
@@ -0,0 +1,98 @@
{
"db_name": "PostgreSQL",
"query": "SELECT *\n FROM websocket_trigger\n WHERE enabled IS TRUE AND (server_id IS NULL OR last_server_ping IS NULL OR last_server_ping < now() - interval '15 seconds')",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "path",
"type_info": "Varchar"
},
{
"ordinal": 1,
"name": "url",
"type_info": "Varchar"
},
{
"ordinal": 2,
"name": "script_path",
"type_info": "Varchar"
},
{
"ordinal": 3,
"name": "is_flow",
"type_info": "Bool"
},
{
"ordinal": 4,
"name": "workspace_id",
"type_info": "Varchar"
},
{
"ordinal": 5,
"name": "edited_by",
"type_info": "Varchar"
},
{
"ordinal": 6,
"name": "email",
"type_info": "Varchar"
},
{
"ordinal": 7,
"name": "edited_at",
"type_info": "Timestamptz"
},
{
"ordinal": 8,
"name": "extra_perms",
"type_info": "Jsonb"
},
{
"ordinal": 9,
"name": "server_id",
"type_info": "Varchar"
},
{
"ordinal": 10,
"name": "last_server_ping",
"type_info": "Timestamptz"
},
{
"ordinal": 11,
"name": "error",
"type_info": "Text"
},
{
"ordinal": 12,
"name": "enabled",
"type_info": "Bool"
},
{
"ordinal": 13,
"name": "filters",
"type_info": "JsonbArray"
}
],
"parameters": {
"Left": []
},
"nullable": [
false,
false,
false,
false,
false,
false,
false,
false,
false,
true,
true,
true,
false,
false
]
},
"hash": "d490fef418e8567fa40aad60e5d46f233ce478086430d8070fe7b2b8a4f9580e"
}
@@ -0,0 +1,15 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM websocket_trigger WHERE workspace_id = $1 AND path = $2",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": []
},
"hash": "d54840373df5da9662ed11ed0a605cddac517bdc06dca0a8be071330338e948a"
}
@@ -0,0 +1,28 @@
{
"db_name": "PostgreSQL",
"query": "SELECT EXISTS(SELECT 1 FROM websocket_trigger WHERE workspace_id = $1) as \"websocket_used!\", EXISTS(SELECT 1 FROM http_trigger WHERE workspace_id = $1) as \"http_routes_used!\"",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "websocket_used!",
"type_info": "Bool"
},
{
"ordinal": 1,
"name": "http_routes_used!",
"type_info": "Bool"
}
],
"parameters": {
"Left": [
"Text"
]
},
"nullable": [
null,
null
]
},
"hash": "db24f1f6ef1e26b7de7926bd3f32d1fa673c4970c25ff9ad98f0fdf199801b53"
}
@@ -0,0 +1,101 @@
{
"db_name": "PostgreSQL",
"query": "SELECT *\n FROM websocket_trigger\n WHERE workspace_id = $1 AND path = $2",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "path",
"type_info": "Varchar"
},
{
"ordinal": 1,
"name": "url",
"type_info": "Varchar"
},
{
"ordinal": 2,
"name": "script_path",
"type_info": "Varchar"
},
{
"ordinal": 3,
"name": "is_flow",
"type_info": "Bool"
},
{
"ordinal": 4,
"name": "workspace_id",
"type_info": "Varchar"
},
{
"ordinal": 5,
"name": "edited_by",
"type_info": "Varchar"
},
{
"ordinal": 6,
"name": "email",
"type_info": "Varchar"
},
{
"ordinal": 7,
"name": "edited_at",
"type_info": "Timestamptz"
},
{
"ordinal": 8,
"name": "extra_perms",
"type_info": "Jsonb"
},
{
"ordinal": 9,
"name": "server_id",
"type_info": "Varchar"
},
{
"ordinal": 10,
"name": "last_server_ping",
"type_info": "Timestamptz"
},
{
"ordinal": 11,
"name": "error",
"type_info": "Text"
},
{
"ordinal": 12,
"name": "enabled",
"type_info": "Bool"
},
{
"ordinal": 13,
"name": "filters",
"type_info": "JsonbArray"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
false,
false,
false,
false,
false,
false,
false,
false,
false,
true,
true,
true,
false,
false
]
},
"hash": "f2baee15e6d1fecd6d2d7b39fda1b50a15ec8bd349f1081af54da5dc2f5e3021"
}
@@ -0,0 +1,25 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE websocket_trigger SET last_server_ping = now(), error = $1 WHERE workspace_id = $2 AND path = $3 AND server_id = $4 AND enabled IS TRUE RETURNING 1",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "?column?",
"type_info": "Int4"
}
],
"parameters": {
"Left": [
"Text",
"Text",
"Text",
"Text"
]
},
"nullable": [
null
]
},
"hash": "febd70e6d5304ad912c232d3d8fc3f8ce0133d8a1a9e5ca1124eff9a988dc2d7"
}
+40
View File
@@ -9315,6 +9315,20 @@ dependencies = [
"xattr",
]
[[package]]
name = "tokio-tungstenite"
version = "0.24.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "edc5f74e248dc973e0dbb7b74c7e0d6fcc301c694ff50049504004ef4d0cdcd9"
dependencies = [
"futures-util",
"log",
"native-tls",
"tokio",
"tokio-native-tls",
"tungstenite",
]
[[package]]
name = "tokio-util"
version = "0.7.12"
@@ -9702,6 +9716,25 @@ version = "0.2.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b"
[[package]]
name = "tungstenite"
version = "0.24.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "18e5b8366ee7a95b16d32197d0b2604b43a0be89dc5fac9f8e96ccafbaedda8a"
dependencies = [
"byteorder",
"bytes",
"data-encoding",
"http 1.1.0",
"httparse",
"log",
"native-tls",
"rand 0.8.5",
"sha1",
"thiserror",
"utf-8",
]
[[package]]
name = "twox-hash"
version = "1.6.3"
@@ -10015,6 +10048,12 @@ dependencies = [
"url",
]
[[package]]
name = "utf-8"
version = "0.7.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9"
[[package]]
name = "utf8-ranges"
version = "1.0.5"
@@ -10461,6 +10500,7 @@ dependencies = [
"tokio",
"tokio-native-tls",
"tokio-tar",
"tokio-tungstenite",
"tokio-util",
"tower 0.5.1",
"tower-cookies",
+1
View File
@@ -291,3 +291,4 @@ syn = { version = "2.0.74", features = ["full"] }
quote = "1.0.36"
regex-lite = "0.1.6"
yaml-rust = "0.4.5"
tokio-tungstenite = { version = "0.24.0", features = ["native-tls"] }
@@ -0,0 +1,2 @@
-- Add down migration script here
DROP TABLE websocket_trigger;
@@ -0,0 +1,67 @@
-- Add up migration script here
CREATE TABLE websocket_trigger (
path VARCHAR(255) NOT NULL,
url VARCHAR(255) NOT NULL,
script_path VARCHAR(255) NOT NULL,
is_flow BOOLEAN NOT NULL,
workspace_id VARCHAR(50) NOT NULL,
edited_by VARCHAR(50) NOT NULL,
email VARCHAR(255) NOT NULL,
edited_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
extra_perms JSONB NOT NULL DEFAULT '{}',
server_id VARCHAR(50) NULL,
last_server_ping TIMESTAMPTZ NULL,
error TEXT NULL,
enabled BOOLEAN NOT NULL,
filters JSONB[] NOT NULL DEFAULT '{}',
PRIMARY KEY (path, workspace_id)
);
GRANT ALL ON websocket_trigger TO windmill_user;
GRANT ALL ON websocket_trigger TO windmill_admin;
ALTER TABLE websocket_trigger ENABLE ROW LEVEL SECURITY;
CREATE POLICY admin_policy ON websocket_trigger FOR ALL TO windmill_admin USING (true);
CREATE POLICY see_folder_extra_perms_user_select ON websocket_trigger FOR SELECT TO windmill_user
USING (SPLIT_PART(websocket_trigger.path, '/', 1) = 'f' AND SPLIT_PART(websocket_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_read'), ',')::text[]));
CREATE POLICY see_folder_extra_perms_user_insert ON websocket_trigger FOR INSERT TO windmill_user
WITH CHECK (SPLIT_PART(websocket_trigger.path, '/', 1) = 'f' AND SPLIT_PART(websocket_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
CREATE POLICY see_folder_extra_perms_user_update ON websocket_trigger FOR UPDATE TO windmill_user
USING (SPLIT_PART(websocket_trigger.path, '/', 1) = 'f' AND SPLIT_PART(websocket_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
CREATE POLICY see_folder_extra_perms_user_delete ON websocket_trigger FOR DELETE TO windmill_user
USING (SPLIT_PART(websocket_trigger.path, '/', 1) = 'f' AND SPLIT_PART(websocket_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[]));
CREATE POLICY see_own ON websocket_trigger FOR ALL TO windmill_user
USING (SPLIT_PART(websocket_trigger.path, '/', 1) = 'u' AND SPLIT_PART(websocket_trigger.path, '/', 2) = current_setting('session.user'));
CREATE POLICY see_member ON websocket_trigger FOR ALL TO windmill_user
USING (SPLIT_PART(websocket_trigger.path, '/', 1) = 'g' AND SPLIT_PART(websocket_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.groups'), ',')::text[]));
CREATE POLICY see_extra_perms_user_select ON websocket_trigger FOR SELECT TO windmill_user
USING (extra_perms ? CONCAT('u/', current_setting('session.user')));
CREATE POLICY see_extra_perms_user_insert ON websocket_trigger FOR INSERT TO windmill_user
WITH CHECK ((extra_perms ->> CONCAT('u/', current_setting('session.user')))::boolean);
CREATE POLICY see_extra_perms_user_update ON websocket_trigger FOR UPDATE TO windmill_user
USING ((extra_perms ->> CONCAT('u/', current_setting('session.user')))::boolean);
CREATE POLICY see_extra_perms_user_delete ON websocket_trigger FOR DELETE TO windmill_user
USING ((extra_perms ->> CONCAT('u/', current_setting('session.user')))::boolean);
CREATE POLICY see_extra_perms_groups_select ON websocket_trigger FOR SELECT TO windmill_user
USING (extra_perms ?| regexp_split_to_array(current_setting('session.pgroups'), ',')::text[]);
CREATE POLICY see_extra_perms_groups_insert ON websocket_trigger FOR INSERT TO windmill_user
WITH CHECK (exists(
SELECT key, value FROM jsonb_each_text(extra_perms)
WHERE SPLIT_PART(key, '/', 1) = 'g' AND key = ANY(regexp_split_to_array(current_setting('session.pgroups'), ',')::text[])
AND value::boolean));
CREATE POLICY see_extra_perms_groups_update ON websocket_trigger FOR UPDATE TO windmill_user
USING (exists(
SELECT key, value FROM jsonb_each_text(extra_perms)
WHERE SPLIT_PART(key, '/', 1) = 'g' AND key = ANY(regexp_split_to_array(current_setting('session.pgroups'), ',')::text[])
AND value::boolean));
CREATE POLICY see_extra_perms_groups_delete ON websocket_trigger FOR DELETE TO windmill_user
USING (exists(
SELECT key, value FROM jsonb_each_text(extra_perms)
WHERE SPLIT_PART(key, '/', 1) = 'g' AND key = ANY(regexp_split_to_array(current_setting('session.pgroups'), ',')::text[])
AND value::boolean));
+1
View File
@@ -95,6 +95,7 @@ openidconnect = { workspace = true, optional = true}
url = { workspace = true, optional = true}
jsonwebtoken = { workspace = true }
matchit.workspace = true
tokio-tungstenite.workspace = true
pin-project.workspace = true
http.workspace = true
+303 -6
View File
@@ -2162,6 +2162,30 @@ paths:
schema:
type: number
/w/{workspace}/workspaces/used_triggers:
get:
summary: get used triggers
operationId: getUsedTriggers
tags:
- workspace
parameters:
- $ref: "#/components/parameters/WorkspaceId"
responses:
"200":
description: status
content:
application/json:
schema:
type: object
properties:
http_routes_used:
type: boolean
websocket_used:
type: boolean
required:
- http_routes_used
- websocket_used
/w/{workspace}/users/list:
get:
summary: list users
@@ -7209,22 +7233,170 @@ paths:
schema:
type: boolean
/w/{workspace}/http_triggers/used:
get:
summary: whether http triggers are used
operationId: used
/w/{workspace}/websocket_triggers/create:
post:
summary: create websocket trigger
operationId: createWebsocketTrigger
tags:
- http_trigger
- websocket_trigger
parameters:
- $ref: "#/components/parameters/WorkspaceId"
requestBody:
description: new websocket trigger
required: true
content:
application/json:
schema:
$ref: "#/components/schemas/NewWebsocketTrigger"
responses:
"201":
description: websocket trigger created
content:
text/plain:
schema:
type: string
/w/{workspace}/websocket_triggers/update/{path}:
post:
summary: update websocket trigger
operationId: updateWebsocketTrigger
tags:
- websocket_trigger
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/Path"
requestBody:
description: updated trigger
required: true
content:
application/json:
schema:
$ref: "#/components/schemas/EditWebsocketTrigger"
responses:
"200":
description: whether http triggers are used
description: websocket trigger updated
content:
text/plain:
schema:
type: string
/w/{workspace}/websocket_triggers/delete/{path}:
delete:
summary: delete websocket trigger
operationId: deleteWebsocketTrigger
tags:
- websocket_trigger
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/Path"
responses:
"200":
description: websocket trigger deleted
content:
text/plain:
schema:
type: string
/w/{workspace}/websocket_triggers/get/{path}:
get:
summary: get websocket trigger
operationId: getWebsocketTrigger
tags:
- websocket_trigger
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/Path"
responses:
"200":
description: websocket trigger deleted
content:
application/json:
schema:
$ref: "#/components/schemas/WebsocketTrigger"
/w/{workspace}/websocket_triggers/list:
get:
summary: list websocket triggers
operationId: listWebsocketTriggers
tags:
- websocket_trigger
parameters:
- $ref: "#/components/parameters/WorkspaceId"
required: true
- $ref: "#/components/parameters/Page"
- $ref: "#/components/parameters/PerPage"
- name: path
description: filter by path
in: query
schema:
type: string
- name: is_flow
in: query
schema:
type: boolean
- name: path_start
in: query
schema:
type: string
responses:
"200":
description: websocket trigger list
content:
application/json:
schema:
type: array
items:
$ref: "#/components/schemas/WebsocketTrigger"
/w/{workspace}/websocket_triggers/exists/{path}:
get:
summary: does websocket trigger exists
operationId: existsWebsocketTrigger
tags:
- websocket_trigger
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/Path"
responses:
"200":
description: websocket trigger exists
content:
application/json:
schema:
type: boolean
/w/{workspace}/websocket_triggers/setenabled/{path}:
post:
summary: set enabled websocket trigger
operationId: setWebsocketTriggerEnabled
tags:
- websocket_trigger
parameters:
- $ref: "#/components/parameters/WorkspaceId"
- $ref: "#/components/parameters/Path"
requestBody:
description: updated websocket trigger enable
required: true
content:
application/json:
schema:
type: object
properties:
enabled:
type: boolean
required:
- enabled
responses:
"200":
description: websocket trigger enabled set
content:
text/plain:
schema:
type: string
/groups/list:
get:
summary: list instance groups
@@ -8052,6 +8224,7 @@ paths:
app,
raw_app,
http_trigger,
websocket_trigger,
]
responses:
"200":
@@ -8089,6 +8262,7 @@ paths:
app,
raw_app,
http_trigger,
websocket_trigger,
]
requestBody:
description: acl to add
@@ -8137,6 +8311,7 @@ paths:
app,
raw_app,
http_trigger,
websocket_trigger,
]
requestBody:
description: acl to add
@@ -11270,6 +11445,128 @@ components:
type: number
email_count:
type: number
websocket_count:
type: number
WebsocketTrigger:
type: object
properties:
path:
type: string
edited_by:
type: string
edited_at:
type: string
format: date-time
script_path:
type: string
url:
type: string
is_flow:
type: boolean
extra_perms:
type: object
additionalProperties:
type: boolean
email:
type: string
workspace_id:
type: string
server_id:
type: string
last_server_ping:
type: string
format: date-time
error:
type: string
enabled:
type: boolean
filters:
type: array
items:
type: object
properties:
key:
type: string
value: {}
required:
- key
- value
required:
- path
- edited_by
- edited_at
- script_path
- url
- extra_perms
- is_flow
- email
- workspace_id
- enabled
- filters
NewWebsocketTrigger:
type: object
properties:
path:
type: string
script_path:
type: string
is_flow:
type: boolean
url:
type: string
enabled:
type: boolean
filters:
type: array
items:
type: object
properties:
key:
type: string
value: {}
required:
- key
- value
required:
- path
- script_path
- url
- is_flow
- filters
EditWebsocketTrigger:
type: object
properties:
url:
type: string
path:
type: string
script_path:
type: string
is_flow:
type: boolean
filters:
type: array
items:
type: object
properties:
key:
type: string
value: {}
required:
- key
- value
required:
- path
- script_path
- url
- is_flow
- filters
Group:
type: object
+1 -37
View File
@@ -12,10 +12,8 @@ use std::collections::HashMap;
use tower_http::cors::CorsLayer;
use windmill_audit::{audit_ee::audit_log, ActionKind};
use windmill_common::{
auth::fetch_authed_from_permissioned_as,
db::UserDB,
error::{self, JsonResult},
users::username_to_permissioned_as,
utils::{not_found_if_none, paginate, require_admin, Pagination, StripPath},
worker::{to_raw_value, CLOUD_HOSTED},
};
@@ -27,7 +25,7 @@ use crate::{
run_flow_by_path_inner, run_script_by_path_inner, run_wait_result_flow_by_path_internal,
run_wait_result_script_by_path_internal, RunJobQuery,
},
users::OptAuthed,
users::{fetch_api_authed, OptAuthed},
};
lazy_static::lazy_static! {
@@ -66,7 +64,6 @@ pub fn workspaced_service() -> Router {
.route("/update/*path", post(update_trigger))
.route("/delete/*path", delete(delete_trigger))
.route("/exists/*path", get(exists_trigger))
.route("/used", get(used))
.route("/route_exists", post(exists_route))
}
@@ -346,17 +343,6 @@ async fn delete_trigger(
Ok(format!("HTTP trigger {path} deleted"))
}
async fn used(Extension(db): Extension<DB>, Path(w_id): Path<String>) -> JsonResult<bool> {
let used = sqlx::query_scalar!(
r#"SELECT EXISTS(SELECT 1 FROM http_trigger WHERE workspace_id = $1)"#,
w_id,
)
.fetch_one(&db)
.await?
.unwrap_or(false);
Ok(Json(used))
}
async fn exists_trigger(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
@@ -421,28 +407,6 @@ struct TriggerRoute {
http_method: HttpMethod,
}
async fn fetch_api_authed(
username: String,
email: String,
w_id: &str,
db: &DB,
username_override: String,
) -> error::Result<ApiAuthed> {
let permissioned_as = username_to_permissioned_as(username.as_str());
let authed =
fetch_authed_from_permissioned_as(permissioned_as, email.clone(), w_id, db).await?;
Ok(ApiAuthed {
username: username,
email: email,
is_admin: authed.is_admin,
is_operator: authed.is_operator,
groups: authed.groups,
folders: authed.folders,
scopes: authed.scopes,
username_override: Some(username_override),
})
}
async fn get_http_route_trigger(
route_path: &str,
opt_authed: Option<ApiAuthed>,
+4 -2
View File
@@ -70,7 +70,10 @@ use windmill_common::{
oauth2::HmacSha256,
scripts::{ScriptHash, ScriptLang},
users::username_to_permissioned_as,
utils::{not_found_if_none, now_from_db, paginate, paginate_without_limits, require_admin, Pagination, StripPath},
utils::{
not_found_if_none, now_from_db, paginate, paginate_without_limits, require_admin,
Pagination, StripPath,
},
};
#[cfg(all(feature = "enterprise", feature = "parquet"))]
@@ -4751,7 +4754,6 @@ async fn get_job_update(
&w_id,
job_id,
"progress_perc"
)
.fetch_optional(&db)
.await?.and_then(|inner| inner)
+9 -2
View File
@@ -87,6 +87,7 @@ mod users;
mod utils;
mod variables;
mod webhook_util;
mod websocket_triggers;
mod workers;
mod workspaces;
@@ -225,7 +226,7 @@ pub async fn run_server(
db: db.clone(),
user_db: user_db,
auth_cache: auth_cache.clone(),
rsmq: rsmq,
rsmq: rsmq.clone(),
base_internal_url: base_internal_url.clone(),
});
if let Err(err) = smtp_server.start_listener_thread(addr).await {
@@ -245,6 +246,8 @@ pub async fn run_server(
}
};
websocket_triggers::start_websockets(db.clone(), rsmq).await;
// build our application with a route
let app = Router::new()
.nest(
@@ -285,7 +288,11 @@ pub async fn run_server(
.nest("/variables", variables::workspaced_service())
.nest("/workspaces", workspaces::workspaced_service())
.nest("/oidc", oidc_ee::workspaced_service())
.nest("/http_triggers", http_triggers::workspaced_service()),
.nest("/http_triggers", http_triggers::workspaced_service())
.nest(
"/websocket_triggers",
websocket_triggers::workspaced_service(),
),
)
.nest("/workspaces", workspaces::global_service())
.nest(
+12
View File
@@ -17,6 +17,7 @@ pub struct TriggersCount {
http_routes_count: i64,
webhook_count: i64,
email_count: i64,
websocket_count: i64,
}
pub(crate) async fn get_triggers_count_internal(
db: &DB,
@@ -53,6 +54,16 @@ pub(crate) async fn get_triggers_count_internal(
.await?
.unwrap_or(0);
let websocket_count = sqlx::query_scalar!(
"SELECT COUNT(*) FROM websocket_trigger WHERE script_path = $1 AND is_flow = $2 AND workspace_id = $3",
path,
is_flow,
w_id
)
.fetch_one(db)
.await?
.unwrap_or(0);
let webhook_count = (if is_flow {
sqlx::query_scalar!(
"SELECT COUNT(*) FROM token WHERE label LIKE 'webhook-%' AND workspace_id = $1 AND scopes @> ARRAY['run:flow/' || $2]::text[]",
@@ -93,6 +104,7 @@ pub(crate) async fn get_triggers_count_internal(
http_routes_count,
webhook_count,
email_count,
websocket_count,
}))
}
@@ -0,0 +1,632 @@
use axum::{
extract::{Path, Query},
routing::{delete, get, post},
Extension, Json, Router,
};
use futures::StreamExt;
use http::StatusCode;
use itertools::Itertools;
use rand::seq::SliceRandom;
use serde::{
de::{self, MapAccess, Visitor},
Deserialize, Deserializer, Serialize,
};
use serde_json::Value;
use sql_builder::{bind::Bind, SqlBuilder};
use sqlx::prelude::FromRow;
use std::{collections::HashMap, fmt};
use tokio_tungstenite::connect_async;
use windmill_audit::{audit_ee::audit_log, ActionKind};
use windmill_common::{
db::UserDB,
error::{self, JsonResult},
utils::{not_found_if_none, paginate, require_admin, Pagination, StripPath},
worker::to_raw_value,
INSTANCE_NAME,
};
use windmill_queue::PushArgsOwned;
use crate::{
db::{ApiAuthed, DB},
jobs::{
run_wait_result_flow_by_path_internal, run_wait_result_script_by_path_internal, RunJobQuery,
},
users::fetch_api_authed,
};
pub fn workspaced_service() -> Router {
Router::new()
.route("/create", post(create_websocket_trigger))
.route("/list", get(list_websocket_triggers))
.route("/get/*path", get(get_websocket_trigger))
.route("/update/*path", post(update_websocket_trigger))
.route("/delete/*path", delete(delete_websocket_trigger))
.route("/exists/*path", get(exists_websocket_trigger))
.route("/setenabled/*path", post(set_enabled))
}
#[derive(Deserialize)]
struct NewWebsocketTrigger {
path: String,
url: String,
script_path: String,
is_flow: bool,
enabled: Option<bool>,
filters: Vec<serde_json::Value>,
}
#[derive(FromRow, Serialize, Clone)]
pub struct WebsocketTrigger {
workspace_id: String,
path: String,
url: String,
script_path: String,
is_flow: bool,
edited_by: String,
email: String,
edited_at: chrono::DateTime<chrono::Utc>,
server_id: Option<String>,
last_server_ping: Option<chrono::DateTime<chrono::Utc>>,
extra_perms: serde_json::Value,
error: Option<String>,
enabled: bool,
filters: Vec<serde_json::Value>,
}
#[derive(Deserialize)]
struct EditWebsocketTrigger {
path: String,
url: String,
script_path: String,
is_flow: bool,
filters: Vec<serde_json::Value>,
}
#[derive(Deserialize)]
pub struct ListWebsocketTriggerQuery {
pub page: Option<usize>,
pub per_page: Option<usize>,
pub path: Option<String>,
pub is_flow: Option<bool>,
pub path_start: Option<String>,
}
async fn list_websocket_triggers(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
Query(lst): Query<ListWebsocketTriggerQuery>,
) -> error::JsonResult<Vec<WebsocketTrigger>> {
let mut tx = user_db.begin(&authed).await?;
let (per_page, offset) = paginate(Pagination { per_page: lst.per_page, page: lst.page });
let mut sqlb = SqlBuilder::select_from("websocket_trigger")
.field("*")
.order_by("edited_at", true)
.and_where("workspace_id = ?".bind(&w_id))
.offset(offset)
.limit(per_page)
.clone();
if let Some(path) = lst.path {
sqlb.and_where_eq("script_path", "?".bind(&path));
}
if let Some(is_flow) = lst.is_flow {
sqlb.and_where_eq("is_flow", "?".bind(&is_flow));
}
if let Some(path_start) = &lst.path_start {
sqlb.and_where_like_left("path", path_start);
}
let sql = sqlb
.sql()
.map_err(|e| error::Error::InternalErr(e.to_string()))?;
let rows = sqlx::query_as::<_, WebsocketTrigger>(&sql)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
Ok(Json(rows))
}
async fn get_websocket_trigger(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> error::JsonResult<WebsocketTrigger> {
let mut tx = user_db.begin(&authed).await?;
let path = path.to_path();
let trigger = sqlx::query_as!(
WebsocketTrigger,
r#"SELECT *
FROM websocket_trigger
WHERE workspace_id = $1 AND path = $2"#,
w_id,
path,
)
.fetch_optional(&mut *tx)
.await?;
tx.commit().await?;
let trigger = not_found_if_none(trigger, "Trigger", path)?;
Ok(Json(trigger))
}
async fn create_websocket_trigger(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
Json(ct): Json<NewWebsocketTrigger>,
) -> error::Result<(StatusCode, String)> {
require_admin(authed.is_admin, &authed.username)?;
let mut tx = user_db.begin(&authed).await?;
sqlx::query_as!(
WebsocketTrigger,
"INSERT INTO websocket_trigger (workspace_id, path, url, script_path, is_flow, enabled, filters, edited_by, email, edited_at) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, now()) RETURNING *",
w_id,
ct.path,
ct.url,
ct.script_path,
ct.is_flow,
ct.enabled.unwrap_or(true),
&ct.filters,
&authed.username,
&authed.email
)
.fetch_one(&mut *tx).await?;
audit_log(
&mut *tx,
&authed,
"websocket_triggers.create",
ActionKind::Create,
&w_id,
Some(ct.path.as_str()),
None,
)
.await?;
tx.commit().await?;
Ok((StatusCode::CREATED, format!("{}", ct.path)))
}
async fn update_websocket_trigger(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
Json(ct): Json<EditWebsocketTrigger>,
) -> error::Result<String> {
let path = path.to_path();
let mut tx = user_db.begin(&authed).await?;
// important to update server_id, last_server_ping and error 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, edited_by = $6, email = $7, edited_at = now(), server_id = NULL, last_server_ping = NULL, error = NULL
WHERE workspace_id = $8 AND path = $9",
ct.url,
ct.script_path,
ct.path,
ct.is_flow,
&ct.filters,
&authed.username,
&authed.email,
w_id,
path,
)
.execute(&mut *tx).await?;
audit_log(
&mut *tx,
&authed,
"websocket_triggers.update",
ActionKind::Create,
&w_id,
Some(path),
None,
)
.await?;
tx.commit().await?;
Ok(path.to_string())
}
#[derive(Deserialize)]
pub struct SetEnabled {
pub enabled: bool,
}
pub async fn set_enabled(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
Json(payload): Json<SetEnabled>,
) -> error::Result<String> {
let mut tx = user_db.begin(&authed).await?;
let path = path.to_path();
// important to set server_id, last_server_ping and error to NULL to stop current websocket listener
let one_o = sqlx::query_scalar!(
"UPDATE websocket_trigger SET enabled = $1, email = $2, edited_by = $3, edited_at = now(), server_id = NULL, last_server_ping = NULL, error = NULL
WHERE path = $4 AND workspace_id = $5 RETURNING 1",
payload.enabled,
&authed.email,
&authed.username,
path,
w_id,
).fetch_optional(&mut *tx).await?;
not_found_if_none(one_o.flatten(), "Websocket trigger", path)?;
audit_log(
&mut *tx,
&authed,
"websocket_triggers.setenabled",
ActionKind::Update,
&w_id,
Some(path),
Some([("enabled", payload.enabled.to_string().as_ref())].into()),
)
.await?;
tx.commit().await?;
Ok(format!(
"succesfully updated websocket trigger at path {} to status {}",
path, payload.enabled
))
}
async fn delete_websocket_trigger(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> error::Result<String> {
require_admin(authed.is_admin, &authed.username)?;
let path = path.to_path();
let mut tx = user_db.begin(&authed).await?;
sqlx::query!(
"DELETE FROM websocket_trigger WHERE workspace_id = $1 AND path = $2",
w_id,
path,
)
.execute(&mut *tx)
.await?;
audit_log(
&mut *tx,
&authed,
"websocket_triggers.delete",
ActionKind::Delete,
&w_id,
Some(path),
None,
)
.await?;
tx.commit().await?;
Ok(format!("Websocket trigger {path} deleted"))
}
async fn exists_websocket_trigger(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<bool> {
let path = path.to_path();
let exists = sqlx::query_scalar!(
"SELECT EXISTS(SELECT 1 FROM websocket_trigger WHERE path = $1 AND workspace_id = $2)",
path,
w_id,
)
.fetch_one(&db)
.await?
.unwrap_or(false);
Ok(Json(exists))
}
pub async fn start_websockets(db: DB, rsmq: Option<rsmq_async::MultiplexedRsmq>) -> () {
tokio::spawn(async move {
loop {
match sqlx::query_as!(
WebsocketTrigger,
r#"SELECT *
FROM websocket_trigger
WHERE enabled IS TRUE AND (server_id IS NULL OR last_server_ping IS NULL OR last_server_ping < now() - interval '15 seconds')"#
)
.fetch_all(&db)
.await
{
Ok(mut triggers) => {
triggers.shuffle(&mut rand::thread_rng());
for trigger in triggers {
maybe_listen_to_websocket(trigger, db.clone(), rsmq.clone()).await;
}
}
Err(err) => {
tracing::error!("Error fetching websocket triggers: {:?}", err);
}
};
tokio::time::sleep(tokio::time::Duration::from_secs(15)).await;
}
});
}
async fn maybe_listen_to_websocket(
ws_trigger: WebsocketTrigger,
db: DB,
rsmq: Option<rsmq_async::MultiplexedRsmq>,
) -> () {
match sqlx::query_scalar!(
"UPDATE websocket_trigger SET server_id = $1, last_server_ping = now() WHERE enabled IS TRUE AND workspace_id = $2 AND path = $3 AND (server_id IS NULL OR last_server_ping IS NULL OR last_server_ping < now() - interval '15 seconds') RETURNING true",
*INSTANCE_NAME,
ws_trigger.workspace_id,
ws_trigger.path,
).fetch_optional(&db).await {
Ok(has_lock) => {
if has_lock.flatten().unwrap_or(false) {
tokio::spawn(listen_to_websocket(ws_trigger, db, rsmq));
} else {
tracing::info!("Websocket {} already being listened to", ws_trigger.url);
}
},
Err(err) => {
tracing::error!("Error acquiring lock for websocket {}: {:?}", ws_trigger.path, err);
}
};
}
struct SupersetVisitor<'a> {
key: &'a str,
value_to_check: &'a Value,
}
impl<'de, 'a> Visitor<'de> for SupersetVisitor<'a> {
type Value = bool;
fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
formatter.write_str("a JSON object with a specific key at the top level")
}
fn visit_map<V>(self, mut map: V) -> Result<Self::Value, V::Error>
where
V: MapAccess<'de>,
{
while let Some(key) = map.next_key::<String>()? {
if key == self.key {
// Deserialize the value for the key and check if it's a superset
let json_value: Value = map.next_value()?;
tracing::info!("json_value: {:?}", json_value);
tracing::info!("value_to_check: {:?}", self.value_to_check);
return Ok(is_superset(&json_value, self.value_to_check));
} else {
// Skip the value if it's not the one we're interested in
let _ = map.next_value::<de::IgnoredAny>()?;
}
}
// If the key was not found, return false
Ok(false)
}
}
// Function to check if json_value is a superset of value_to_check
fn is_superset(json_value: &Value, value_to_check: &Value) -> bool {
match (json_value, value_to_check) {
(Value::Object(json_map), Value::Object(check_map)) => {
// Check that all keys and values in check_map exist and match in json_map
check_map.iter().all(|(k, v)| {
json_map
.get(k)
.map_or(false, |json_val| is_superset(json_val, v))
})
}
(Value::Array(json_array), Value::Array(check_array)) => {
// Check that all elements in check_array exist in json_array
check_array.iter().all(|check_item| {
json_array
.iter()
.any(|json_item| is_superset(json_item, check_item))
})
}
_ => json_value == value_to_check,
}
}
// A function to deserialize and check if the value at the given key is a superset of a passed value
fn is_value_superset<'a, 'de, D>(
deserializer: D,
key: &'a str,
value_to_check: &'a Value,
) -> Result<bool, D::Error>
where
D: Deserializer<'de>,
{
deserializer.deserialize_map(SupersetVisitor { key, value_to_check })
}
async fn listen_to_websocket(
ws_trigger: WebsocketTrigger,
db: DB,
rsmq: Option<rsmq_async::MultiplexedRsmq>,
) -> () {
async fn update_ping(db: DB, ws_trigger: &WebsocketTrigger, error: Option<&str>) -> Option<()> {
match sqlx::query_scalar!(
"UPDATE websocket_trigger SET last_server_ping = now(), error = $1 WHERE workspace_id = $2 AND path = $3 AND server_id = $4 AND enabled IS TRUE RETURNING 1",
error,
ws_trigger.workspace_id,
ws_trigger.path,
*INSTANCE_NAME
).fetch_optional(&db).await {
Ok(updated) => {
if updated.flatten().is_none() {
tracing::info!("Websocket {} changed, disabled, or deleted, stopping...", ws_trigger.url);
return None;
}
},
Err(err) => {
tracing::warn!("Error updating ping of websocket {}: {:?}", ws_trigger.url, err);
}
};
Some(())
}
let url = ws_trigger.url.as_str();
#[derive(Deserialize)]
struct JsonFilter {
key: String,
value: serde_json::Value,
}
#[derive(Deserialize)]
#[serde(untagged)]
enum Filter {
JsonFilter(JsonFilter),
}
let filters: Vec<Filter> = ws_trigger
.filters
.iter()
.filter_map(|m| serde_json::from_value(m.clone()).ok())
.collect_vec();
loop {
match connect_async(url).await {
Ok((ws_stream, _)) => {
tracing::info!("Listening to websocket {}", url);
if let None = update_ping(db.clone(), &ws_trigger, None).await {
return;
}
let (_, mut read) = ws_stream.split();
loop {
tokio::select! {
msg = read.next() => {
if let Some(msg) = msg {
match msg {
Ok(msg) => {
match msg {
tokio_tungstenite::tungstenite::Message::Text(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;
}
}
if should_handle {
let db_ = db.clone();
let rsmq_ = rsmq.clone();
let ws_trigger_ = ws_trigger.clone();
tokio::spawn(async move {
let url = ws_trigger_.url.clone();
if let Err(err) = run_job(db_, rsmq_, ws_trigger_, text).await {
tracing::error!("Error running job on websocket {}: {:?}", url, err);
};
});
}
},
_ => {}
}
},
Err(err) => {
tracing::error!("Error reading from websocket {}: {:?}", url, err);
}
}
} else {
tracing::error!("Websocket {} closed, reconnecting in 5s...", url);
break;
}
},
_ = tokio::time::sleep(tokio::time::Duration::from_secs(5)) => {
if let None = update_ping(db.clone(), &ws_trigger, None).await {
return;
}
},
}
}
}
Err(err) => {
tracing::error!("Error connecting to websocket {}: {:?}", url, err);
if let None =
update_ping(db.clone(), &ws_trigger, Some(err.to_string().as_str())).await
{
return;
}
}
}
tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
}
}
async fn run_job(
db: DB,
rsmq: Option<rsmq_async::MultiplexedRsmq>,
trigger: WebsocketTrigger,
msg: String,
) -> anyhow::Result<()> {
let args = PushArgsOwned {
args: HashMap::from([("msg".to_string(), to_raw_value(&msg))]),
extra: Some(HashMap::from([(
"wm_trigger".to_string(),
to_raw_value(&serde_json::json!({"kind": "websocket"})),
)])),
};
let label_prefix = Some(format!("ws-{}-", trigger.path));
let authed = fetch_api_authed(
trigger.edited_by.clone(),
trigger.email.clone(),
&trigger.workspace_id,
&db,
"anonymous".to_string(),
)
.await?;
let user_db = UserDB::new(db.clone());
let run_query = RunJobQuery::default();
if trigger.is_flow {
run_wait_result_flow_by_path_internal(
db,
run_query,
StripPath(trigger.script_path.to_owned()),
authed,
rsmq,
user_db,
args,
trigger.workspace_id.clone(),
label_prefix,
)
.await?;
} else {
run_wait_result_script_by_path_internal(
db,
run_query,
StripPath(trigger.script_path.to_owned()),
authed,
rsmq,
user_db,
trigger.workspace_id.clone(),
args,
label_prefix,
)
.await?;
}
Ok(())
}
+26 -1
View File
@@ -116,7 +116,8 @@ pub fn workspaced_service() -> Router {
.route("/get_workspace_name", get(get_workspace_name))
.route("/change_workspace_name", post(change_workspace_name))
.route("/change_workspace_id", post(change_workspace_id))
.route("/usage", get(get_usage));
.route("/usage", get(get_usage))
.route("/used_triggers", get(get_used_triggers));
#[cfg(feature = "stripe")]
{
@@ -1488,6 +1489,30 @@ async fn set_encryption_key(
return Ok(());
}
#[derive(Serialize)]
struct UsedTriggers {
pub websocket_used: bool,
pub http_routes_used: bool,
}
async fn get_used_triggers(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
) -> JsonResult<UsedTriggers> {
let mut tx = user_db.begin(&authed).await?;
let websocket_used = sqlx::query_as!(
UsedTriggers,
r#"SELECT EXISTS(SELECT 1 FROM websocket_trigger WHERE workspace_id = $1) as "websocket_used!", EXISTS(SELECT 1 FROM http_trigger WHERE workspace_id = $1) as "http_routes_used!""#,
w_id,
)
.fetch_one(&mut *tx)
.await?;
tx.commit().await?;
Ok(Json(websocket_used))
}
async fn list_workspaces_as_super_admin(
authed: ApiAuthed,
Extension(db): Extension<DB>,
+3 -3
View File
@@ -475,9 +475,9 @@
const testStepStore = writable<Record<string, any>>({})
const selectedIdStore = writable('settings-metadata')
const selectedTriggerStore = writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes'>(
'webhooks'
)
const selectedTriggerStore = writable<
'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets'
>('webhooks')
const primaryScheduleStore = writable<ScheduleTrigger | undefined | false>(undefined)
const triggersCount = writable<TriggersCount | undefined>(undefined)
@@ -461,9 +461,9 @@
}
const selectedIdStore = writable<string>(selectedId ?? 'settings-metadata')
const selectedTriggerStore = writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes'>(
'webhooks'
)
const selectedTriggerStore = writable<
'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets'
>('webhooks')
export function getSelectedId() {
return $selectedIdStore
@@ -483,7 +483,9 @@
selectedIdStore.set(selectedId)
}
function selectTrigger(selectedTrigger: 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes') {
function selectTrigger(
selectedTrigger: 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets'
) {
selectedTriggerStore.set(selectedTrigger)
}
+12 -5
View File
@@ -13,7 +13,8 @@
ScheduleService,
ScriptService,
HttpTriggerService,
VariableService
VariableService,
WebsocketTriggerService
} from '$lib/gen'
import { superadmin, userStore, workspaceStore } from '$lib/stores'
import { createEventDispatcher, getContext } from 'svelte'
@@ -35,6 +36,7 @@
| 'app'
| 'raw_app'
| 'http_trigger'
| 'websocket_trigger'
let meta: Meta | undefined = undefined
export let fullNamePlaceholder: string | undefined = undefined
export let namePlaceholder = ''
@@ -218,6 +220,11 @@
workspace: $workspaceStore!,
path: path
})
} else if (kind == 'websocket_trigger') {
return await WebsocketTriggerService.existsWebsocketTrigger({
workspace: $workspaceStore!,
path: path
})
} else {
return false
}
@@ -231,10 +238,10 @@
error = 'This name is not valid'
return false
} else if (meta.owner == '' && meta.ownerKind == 'folder') {
error = 'Folder need to be chosen'
error = 'Folder needs to be chosen'
return false
} else if (meta.owner == '' && meta.ownerKind == 'group') {
error = 'Group need to be chosen'
error = 'Group needs to be chosen'
return false
} else {
return true
@@ -455,10 +462,10 @@
/>
<!-- <span class="font-mono text-sm break-all">{path}</span> -->
</div>
<div class="text-red-600 dark:text-red-400 text-2xs">{error}</div>
<div class="text-red-600 dark:text-red-400 text-2xs mt-1.5">{error}</div>
</div>
{#if kind != 'app' && kind != 'schedule' && kind != 'http_trigger' && initialPath != '' && initialPath != undefined && initialPath != path}
{#if kind != 'app' && kind != 'schedule' && kind != 'http_trigger' && kind != 'websocket_trigger' && initialPath != '' && initialPath != undefined && initialPath != path}
<Alert type="warning" class="mt-4" title="Moving may break other items relying on it">
You are renaming an item that may be depended upon by other items. This may break apps, flows
or resources. Find if it used elsewhere using the content search. Note that linked variables
@@ -109,9 +109,9 @@
? { schedule_count: 1, primary_schedule: { schedule: savedPrimarySchedule.cron } }
: undefined
)
const selectedTriggerStore = writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes'>(
'webhooks'
)
const selectedTriggerStore = writable<
'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets'
>('webhooks')
export function setPrimarySchedule(schedule: ScheduleTrigger | undefined | false) {
primaryScheduleStore.set(schedule)
@@ -25,6 +25,7 @@
| 'app'
| 'raw_app'
| 'http_trigger'
| 'websocket_trigger'
let kind: Kind
let path: string = ''
@@ -6,7 +6,13 @@
import FlowViewerInner from '../FlowViewerInner.svelte'
import DetailPageTriggerPanel from './DetailPageTriggerPanel.svelte'
export let triggerSelected: 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' = 'webhooks'
export let triggerSelected:
| 'webhooks'
| 'emails'
| 'schedules'
| 'cli'
| 'routes'
| 'websockets' = 'webhooks'
export let flow_json: any | undefined = undefined
export let isOperator: boolean = false
@@ -45,6 +51,7 @@
<DetailPageTriggerPanel bind:triggerSelected>
<slot slot="webhooks" name="webhooks" />
<slot slot="routes" name="routes" />
<slot slot="websockets" name="websockets" />
<slot slot="emails" name="emails" />
<slot slot="schedules" name="schedules" />
<slot slot="cli" name="cli" />
@@ -19,9 +19,9 @@
let clientWidth = window.innerWidth
const primaryScheduleStore = writable<ScheduleTrigger | undefined | false>(undefined)
const selectedTriggerStore = writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes'>(
'webhooks'
)
const selectedTriggerStore = writable<
'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets'
>('webhooks')
setContext<TriggerContext>('TriggerContext', {
selectedTrigger: selectedTriggerStore,
@@ -48,6 +48,7 @@
>
<slot slot="webhooks" name="webhooks" />
<slot slot="routes" name="routes" />
<slot slot="websockets" name="websockets" />
<slot slot="emails" name="emails" />
<slot slot="schedules" name="schedules" />
<slot slot="cli" name="cli" />
@@ -81,6 +82,7 @@
<DetailPageTriggerPanel bind:triggerSelected={$selectedTriggerStore}>
<slot slot="webhooks" name="webhooks" />
<slot slot="routes" name="routes" />
<slot slot="websockets" name="websockets" />
<slot slot="emails" name="emails" />
<slot slot="schedules" name="schedules" />
<slot slot="cli" name="cli" />
@@ -1,10 +1,16 @@
<script lang="ts">
import { Tabs, Tab } from '$lib/components/common'
import { CalendarCheck2, MailIcon, Route, Terminal, Webhook } from 'lucide-svelte'
import { CalendarCheck2, MailIcon, Route, Terminal, Webhook, Unplug } from 'lucide-svelte'
import HighlightTheme from '../HighlightTheme.svelte'
export let triggerSelected: 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' = 'webhooks'
export let triggerSelected:
| 'webhooks'
| 'emails'
| 'schedules'
| 'cli'
| 'routes'
| 'websockets' = 'webhooks'
</script>
<HighlightTheme />
@@ -28,6 +34,12 @@
HTTP
</span>
</Tab>
<Tab value="websockets">
<span class="flex flex-row gap-2 items-center">
<Unplug size={14} />
Websockets
</span>
</Tab>
<Tab value="emails">
<span class="flex flex-row gap-2 items-center">
<MailIcon size={14} />
@@ -52,6 +64,8 @@
<slot name="emails" />
{:else if triggerSelected === 'schedules'}
<slot name="schedules" />
{:else if triggerSelected === 'websockets'}
<slot name="websockets" />
{:else if triggerSelected === 'cli'}
<slot name="cli" />
{/if}
@@ -1,5 +1,5 @@
<script lang="ts">
import { Calendar, Mail, Webhook } from 'lucide-svelte'
import { Calendar, Mail, Webhook, Unplug } from 'lucide-svelte'
import TriggerButton from './TriggerButton.svelte'
import Popover from '$lib/components/Popover.svelte'
@@ -89,6 +89,22 @@
</Popover>
{/if}
{#if !showOnlyWithCount || ($triggersCount?.websocket_count ?? 0) > 0}
<Popover>
<svelte:fragment slot="text">Websockets</svelte:fragment>
<TriggerButton
on:click={() => {
$selectedTrigger = 'websockets'
dispatch('select')
}}
selected={selected && $selectedTrigger === 'websockets'}
>
<TriggerCount count={$triggersCount?.websocket_count} />
<Unplug size={12} />
</TriggerButton>
</Popover>
{/if}
{#if !showOnlyWithCount || ($triggersCount?.email_count ?? 0) > 0}
<Popover>
<svelte:fragment slot="text">Emails</svelte:fragment>
@@ -21,7 +21,8 @@
ServerCog,
Settings,
UserCog,
Plus
Plus,
Unplug
} from 'lucide-svelte'
import Menu from '../common/menu/MenuV2.svelte'
import MenuButton from './MenuButton.svelte'
@@ -77,6 +78,13 @@
icon: Route,
disabled: $userStore?.operator,
kind: 'http'
},
{
label: 'Websockets',
href: '/websocket_triggers',
icon: Unplug,
disabled: $userStore?.operator,
kind: 'ws'
}
]
+1 -1
View File
@@ -10,7 +10,7 @@ export type ScheduleTrigger = {
}
export type TriggerContext = {
selectedTrigger: Writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes'>
selectedTrigger: Writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets'>
primarySchedule: Writable<ScheduleTrigger | undefined | false>
triggersCount: Writable<TriggersCount | undefined>
}
@@ -39,6 +39,7 @@
itemKind = isFlow ? 'flow' : 'script'
edit = true
dirtyPath = false
dirtyRoutePath = false
await loadTrigger()
} catch (err) {
sendUserToast(`Could not load route: ${err}`, true)
@@ -58,6 +59,7 @@
requires_auth = false
initialRoutePath = ''
route_path = ''
dirtyRoutePath = false
http_method = 'post'
initialScriptPath = ''
fixedScriptPath = fixedScriptPath_ ?? ''
@@ -73,6 +75,7 @@
let path: string = ''
let pathError = ''
let routeError = ''
let dirtyRoutePath = false
let is_async = false
let requires_auth = false
let initialRoutePath = ''
@@ -239,6 +242,9 @@
class={routeError === ''
? ''
: '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'}
on:input={() => {
dirtyRoutePath = true
}}
/>
</label>
@@ -272,7 +278,10 @@
}}
/>
</div>
<div class="text-red-600 dark:text-red-400 text-2xs">{routeError}</div>
<div class="text-red-600 dark:text-red-400 text-2xs mt-1.5"
>{dirtyRoutePath ? routeError : ''}</div
>
</div>
</div>
</Section>
@@ -10,6 +10,7 @@
import FlowCard from '../flows/common/FlowCard.svelte'
import { getContext } from 'svelte'
import type { TriggerContext } from '$lib/components/triggers'
import WebsocketTriggersPanel from './WebsocketTriggersPanel.svelte'
export let noEditor: boolean
export let newItem = false
@@ -27,6 +28,7 @@
<Tab value="webhooks" selectedClass="text-primary font-semibold">Webhooks</Tab>
<Tab value="schedules" selectedClass="text-primary text-sm font-semibold">Schedules</Tab>
<Tab value="routes" selectedClass="text-primary text-sm font-semibold">Routes</Tab>
<Tab value="websockets" selectedClass="text-primary text-sm font-semibold">Websockets</Tab>
<Tab value="emails" selectedClass="text-primary text-sm font-semibold">Email</Tab>
<svelte:fragment slot="content">
@@ -60,6 +62,12 @@
</div>
{/if}
{#if $selectedTrigger === 'websockets'}
<div class="p-4">
<WebsocketTriggersPanel {newItem} path={currentPath} {isFlow} />
</div>
{/if}
{#if $selectedTrigger === 'schedules'}
<div class="p-2">
<RunPageSchedules
@@ -0,0 +1,23 @@
<script lang="ts">
import { tick } from 'svelte'
import WebsocketTriggerEditorInner from './WebsocketTriggerEditorInner.svelte'
let open = false
export async function openEdit(ePath: string, isFlow: boolean) {
open = true
await tick()
drawer?.openEdit(ePath, isFlow)
}
export async function openNew(is_flow: boolean, initial_script_path?: string) {
open = true
await tick()
drawer?.openNew(is_flow, initial_script_path)
}
let drawer: WebsocketTriggerEditorInner
</script>
{#if open}
<WebsocketTriggerEditorInner on:update bind:this={drawer} />
{/if}
@@ -0,0 +1,338 @@
<script lang="ts">
import { Alert, Button } from '$lib/components/common'
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'
import Required from '$lib/components/Required.svelte'
import ScriptPicker from '$lib/components/ScriptPicker.svelte'
import { WebsocketTriggerService } from '$lib/gen'
import { usedTriggerKinds, userStore, workspaceStore } from '$lib/stores'
import { canWrite, emptyString, sendUserToast } from '$lib/utils'
import { createEventDispatcher } from 'svelte'
import Section from '$lib/components/Section.svelte'
import { Loader2, Save, X, Plus } from 'lucide-svelte'
import Label from '$lib/components/Label.svelte'
import Toggle from '../Toggle.svelte'
import { fade } from 'svelte/transition'
import JsonEditor from '../apps/editor/settingsPanel/inputEditor/JsonEditor.svelte'
let drawer: Drawer
let is_flow: boolean = false
let initialPath = ''
let edit = true
let itemKind: 'flow' | 'script' = 'script'
let script_path = ''
let initialScriptPath = ''
let fixedScriptPath = ''
let path: string = ''
let pathError = ''
let url = ''
let urlError = ''
let dirtyUrl = false
let enabled = false
let filters: {
key: string
value: any
}[] = []
let dirtyPath = false
let can_write = true
let drawerLoading = true
const dispatch = createEventDispatcher()
$: is_flow = itemKind === 'flow'
export async function openEdit(ePath: string, isFlow: boolean) {
drawerLoading = true
try {
drawer?.openDrawer()
initialPath = ePath
itemKind = isFlow ? 'flow' : 'script'
edit = true
dirtyPath = false
dirtyUrl = false
await loadTrigger()
} catch (err) {
sendUserToast(`Could not load websocket trigger: ${err}`, true)
} finally {
drawerLoading = false
}
}
export async function openNew(nis_flow: boolean, fixedScriptPath_?: string) {
drawerLoading = true
try {
drawer?.openDrawer()
is_flow = nis_flow
edit = false
itemKind = nis_flow ? 'flow' : 'script'
url = ''
dirtyUrl = false
initialScriptPath = ''
fixedScriptPath = fixedScriptPath_ ?? ''
script_path = fixedScriptPath
path = ''
initialPath = ''
filters = []
dirtyPath = false
} finally {
drawerLoading = false
}
}
async function loadTrigger(): Promise<void> {
const s = await WebsocketTriggerService.getWebsocketTrigger({
workspace: $workspaceStore!,
path: initialPath
})
script_path = s.script_path
initialScriptPath = s.script_path
is_flow = s.is_flow
path = s.path
url = s.url
enabled = s.enabled
filters = s.filters
can_write = canWrite(s.path, s.extra_perms, $userStore)
}
async function updateTrigger(): Promise<void> {
if (edit) {
await WebsocketTriggerService.updateWebsocketTrigger({
workspace: $workspaceStore!,
path: initialPath,
requestBody: {
path,
script_path,
is_flow,
url,
filters
}
})
sendUserToast(`Route ${path} updated`)
} else {
await WebsocketTriggerService.createWebsocketTrigger({
workspace: $workspaceStore!,
requestBody: {
path,
script_path,
is_flow,
url,
enabled: true,
filters
}
})
sendUserToast(`Route ${path} created`)
}
if (!$usedTriggerKinds.includes('ws')) {
$usedTriggerKinds = [...$usedTriggerKinds, 'ws']
}
dispatch('update')
drawer.closeDrawer()
}
let validateTimeout: NodeJS.Timeout | undefined = undefined
function validateUrl(url: string) {
urlError = ''
if (validateTimeout) {
clearTimeout(validateTimeout)
}
validateTimeout = setTimeout(() => {
if (/^(ws:|wss:)\/\/[^\s]+$/.test(url) === false) {
urlError = 'Invalid websocket URL'
}
validateTimeout = undefined
}, 500)
}
$: validateUrl(url)
</script>
<Drawer size="700px" bind:this={drawer}>
<DrawerContent
title={edit
? can_write
? `Edit WS trigger ${initialPath}`
: `WS trigger ${initialPath}`
: 'New WS trigger'}
on:close={drawer.closeDrawer}
>
<svelte:fragment slot="actions">
{#if !drawerLoading && can_write}
{#if edit}
<div class="mr-8 center-center -mt-1">
<Toggle
disabled={!can_write}
checked={enabled}
options={{ right: 'enable', left: 'disable' }}
on:change={async (e) => {
await WebsocketTriggerService.setWebsocketTriggerEnabled({
path: initialPath,
workspace: $workspaceStore ?? '',
requestBody: { enabled: e.detail }
})
sendUserToast(
`${e.detail ? 'enabled' : 'disabled'} websocket trigger ${initialPath}`
)
}}
/>
</div>
{/if}
<Button
startIcon={{ icon: Save }}
disabled={pathError != '' || urlError != '' || emptyString(script_path) || !can_write}
on:click={updateTrigger}
>
Save
</Button>
{/if}
</svelte:fragment>
{#if drawerLoading}
<Loader2 class="animate-spin" />
{:else}
<Alert title="Info" type="info">
{#if edit}
Changes can take up to 30 seconds to take effect.
{:else}
New websocket triggers can take up to 30 seconds to start listening.
{/if}
</Alert>
<div class="flex flex-col gap-12 mt-6">
<div class="flex flex-col gap-4">
<Label label="Path">
<Path
bind:dirty={dirtyPath}
bind:error={pathError}
bind:path
{initialPath}
checkInitialPathExistence={!edit}
namePlaceholder="ws_trigger"
kind="websocket_trigger"
disabled={!can_write}
/>
</Label>
</div>
<Section label="Websocket">
<div class="flex flex-col w-full gap-4">
<label class="block grow w-full">
<div class="text-secondary text-sm flex items-center gap-1 w-full justify-between">
<div>
URL
<Required required={true} />
</div>
</div>
<input
type="text"
autocomplete="off"
bind:value={url}
disabled={!can_write}
on:input={() => {
dirtyUrl = true
}}
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'}
/>
<div class="text-red-600 dark:text-red-400 text-2xs mt-1.5">
{dirtyUrl ? urlError : ''}
</div>
</label>
</div>
</Section>
<Section label="Runnable">
<p class="text-xs mb-1 text-tertiary">
Pick a script or flow to be triggered<Required required={true} />
</p>
<div class="flex flex-row mb-2">
<ScriptPicker
disabled={fixedScriptPath != '' || !can_write}
initialPath={fixedScriptPath || initialScriptPath}
kinds={['script']}
allowFlow={true}
bind:itemKind
bind:scriptPath={script_path}
allowRefresh
/>
</div>
</Section>
<Section label="Filters">
<p class="text-xs mb-1 text-tertiary">
Filters will limit the execution of the trigger to only messages that match all
criteria.<br />
The JSON filter checks if the value at the key is equal or a superset of the filter value.
</p>
<div class="flex flex-col gap-4 mt-1">
{#each filters as v, i}
<div class="flex w-full gap-4 items-center">
<div class="w-full flex flex-col gap-2">
<div class="flex flex-row gap-2 w-full">
<label class="flex flex-col w-full">
<div class="text-secondary text-sm">Type</div>
<select
class="w-20"
on:change={(e) => {
if (e.target?.['value']) {
filters[i] = {
key: '',
value: ''
}
}
}}
value={'json'}
>
<option value="json">JSON</option>
</select>
</label>
</div>
<label class="flex flex-col w-full">
<div class="text-secondary text-sm">Key</div>
<input type="text" bind:value={v.key} />
</label>
<!-- svelte-ignore a11y-label-has-associated-control -->
<label class="flex flex-col w-full">
<div class="text-secondary text-sm">Value</div>
<JsonEditor bind:value={v.value} code={JSON.stringify(v.value)} />
</label>
</div>
<button
transition:fade|local={{ duration: 100 }}
class="rounded-full p-1 bg-surface-secondary duration-200 hover:bg-surface-hover"
aria-label="Clear"
on:click={() => {
filters = filters.filter((_, index) => index !== i)
}}
>
<X size={14} />
</button>
</div>
{/each}
<div class="flex items-baseline">
<Button
variant="border"
color="light"
size="xs"
btnClasses="mt-1"
on:click={() => {
if (filters == undefined || !Array.isArray(filters)) {
filters = []
}
filters = filters.concat({
key: '',
value: ''
})
}}
startIcon={{ icon: Plus }}
>
Add item
</Button>
</div>
</div>
</Section>
</div>
{/if}
</DrawerContent>
</Drawer>
@@ -0,0 +1,103 @@
<script lang="ts">
import { Button } from '../common'
import { userStore, workspaceStore } from '$lib/stores'
import { WebsocketTriggerService, type WebsocketTrigger } from '$lib/gen'
import { UnplugIcon } from 'lucide-svelte'
import Skeleton from '../common/skeleton/Skeleton.svelte'
import { canWrite } from '$lib/utils'
import Alert from '../common/alert/Alert.svelte'
import type { TriggerContext } from '../triggers'
import { getContext } from 'svelte'
import WebsocketTriggerEditor from './WebsocketTriggerEditor.svelte'
export let isFlow: boolean
export let path: string
export let newItem: boolean = false
let wsTriggerEditor: WebsocketTriggerEditor
$: path && loadTriggers()
const { triggersCount } = getContext<TriggerContext>('TriggerContext')
let wsTriggers: (WebsocketTrigger & { canWrite: boolean })[] | undefined = undefined
export async function loadTriggers() {
try {
wsTriggers = (
await WebsocketTriggerService.listWebsocketTriggers({
workspace: $workspaceStore ?? '',
path,
isFlow
})
).map((x) => {
return { canWrite: canWrite(x.path, x.extra_perms!, $userStore), ...x }
})
$triggersCount = { ...($triggersCount ?? {}), websocket_count: wsTriggers?.length }
} catch (e) {
console.error('impossible to load WS triggers', e)
}
}
</script>
<WebsocketTriggerEditor
on:update={() => {
loadTriggers()
}}
bind:this={wsTriggerEditor}
/>
<div class="flex flex-col gap-4">
{#if !newItem}
{#if $userStore?.is_admin || $userStore?.is_super_admin}
<Button
on:click={() => wsTriggerEditor?.openNew(isFlow, path)}
variant="border"
color="light"
size="xs"
startIcon={{ icon: UnplugIcon }}
>
New WS Trigger
</Button>
{:else}
<Alert title="Only workspace admins can create routes" type="warning" size="xs" />
{/if}
{/if}
{#if wsTriggers}
{#if wsTriggers.length == 0}
<div class="text-xs text-secondary"> No WS triggers </div>
{:else}
<div class="flex flex-col divide-y pt-2">
{#each wsTriggers as wsTriggers (wsTriggers.path)}
<div class="grid grid-cols-5 text-2xs items-center py-2">
<div class="col-span-2 truncate">{wsTriggers.path}</div>
<div class="col-span-2 truncate">
{wsTriggers.url}
</div>
<div class="flex justify-end">
<button
on:click={() => wsTriggerEditor?.openEdit(wsTriggers.path, isFlow)}
class="px-2"
>
{#if wsTriggers.canWrite}
Edit
{:else}
View
{/if}
</button>
</div>
</div>
{/each}
</div>
{/if}
{:else}
<Skeleton layout={[[8]]} />
{/if}
{#if newItem}
<Alert title="Triggers disabled" type="warning" size="xs">
Deploy the {isFlow ? 'flow' : 'script'} to add WS triggers.
</Alert>
{/if}
</div>
+3 -3
View File
@@ -514,7 +514,7 @@ export async function main(approver?: string) {
const BUN_PREPROCESSOR_MODULE_CODE = `
export async function preprocessor(
wm_trigger: {
kind: 'http' | 'email' | 'webhook',
kind: 'http' | 'email' | 'webhook' | 'websocket',
http?: {
route: string // The route path, e.g. "/users/:id"
path: string // The actual path called, e.g. "/users/123"
@@ -535,7 +535,7 @@ export async function preprocessor(
const DENO_PREPROCESSOR_MODULE_CODE = `
export async function preprocessor(
wm_trigger: {
kind: 'http' | 'email' | 'wehbook',
kind: 'http' | 'email' | 'wehbook' | 'websocket',
http?: {
route: string // The route path, e.g. "/users/:id"
path: string // The actual path called, e.g. "/users/123"
@@ -591,7 +591,7 @@ class Http(TypedDict):
headers: dict[str, str]
class WmTrigger(TypedDict):
kind: Literal["http", "email", "webhook"]
kind: Literal["http", "email", "webhook", "websocket"]
http: Http | None
def preprocessor(
@@ -8,7 +8,6 @@
RawAppService,
ScriptService,
SettingService,
HttpTriggerService,
UserService,
WorkspaceService
} from '$lib/gen'
@@ -187,8 +186,17 @@
}
async function loadUsedTriggerKinds() {
const httpUsed = await HttpTriggerService.used({ workspace: $workspaceStore ?? '' })
$usedTriggerKinds = httpUsed ? ['http'] : []
let usedKinds: string[] = []
const { http_routes_used, websocket_used } = await WorkspaceService.getUsedTriggers({
workspace: $workspaceStore ?? ''
})
if (http_routes_used) {
usedKinds.push('http')
}
if (websocket_used) {
usedKinds.push('ws')
}
$usedTriggerKinds = usedKinds
}
function pathInAppMode(pathname: string | undefined): boolean {
@@ -60,6 +60,7 @@
import json from 'svelte-highlight/languages/json'
import { writable } from 'svelte/store'
import TriggersBadge from '$lib/components/graph/renderers/triggers/TriggersBadge.svelte'
import WebsocketTriggersPanel from '$lib/components/triggers/WebsocketTriggersPanel.svelte'
let flow: Flow | undefined
let can_write = false
@@ -533,6 +534,11 @@
<RoutesPanel path={flow.path ?? ''} isFlow />
</div>
</svelte:fragment>
<svelte:fragment slot="websockets">
<div class="p-2">
<WebsocketTriggersPanel path={flow.path ?? ''} isFlow />
</div>
</svelte:fragment>
<svelte:fragment slot="emails">
<div class="p-2">
<EmailTriggerPanel
@@ -54,7 +54,7 @@
const TRIGGER_PATH_KIND_FILTER_SETTING = 'filter_path_of'
const FILTER_USER_FOLDER_SETTING_NAME = 'user_and_folders_only'
let selectedFilterKind =
(getLocalSetting(TRIGGER_PATH_KIND_FILTER_SETTING) as 'route' | 'script_flow') ?? 'route'
(getLocalSetting(TRIGGER_PATH_KIND_FILTER_SETTING) as 'trigger' | 'script_flow') ?? 'trigger'
let filterUserFolders = getLocalSetting(FILTER_USER_FOLDER_SETTING_NAME) == 'true'
$: storeLocalSetting(TRIGGER_PATH_KIND_FILTER_SETTING, selectedFilterKind)
@@ -62,12 +62,12 @@
function filterItemsPathsBaseOnUserFilters(
item: TriggerW,
selectedFilterKind: 'route' | 'script_flow',
selectedFilterKind: 'trigger' | 'script_flow',
filterUserFolders: boolean
) {
if ($workspaceStore == 'admins') return true
if (filterUserFolders) {
if (selectedFilterKind === 'route') {
if (selectedFilterKind === 'trigger') {
return (
!item.path.startsWith('u/') || item.path.startsWith('u/' + $userStore?.username + '/')
)
@@ -84,15 +84,15 @@
$: preFilteredItems =
ownerFilter != undefined
? selectedFilterKind === 'route'
? selectedFilterKind === 'trigger'
? triggers?.filter(
(x) =>
x.path.startsWith(ownerFilter + '/' ?? '') &&
x.path.startsWith(ownerFilter + '/') &&
filterItemsPathsBaseOnUserFilters(x, selectedFilterKind, filterUserFolders)
)
: triggers?.filter(
(x) =>
x.script_path.startsWith(ownerFilter + '/' ?? '') &&
x.script_path.startsWith(ownerFilter + '/') &&
filterItemsPathsBaseOnUserFilters(x, selectedFilterKind, filterUserFolders)
)
: triggers?.filter((x) =>
@@ -104,7 +104,7 @@
}
$: owners =
selectedFilterKind === 'route'
selectedFilterKind === 'trigger'
? Array.from(
new Set(filteredItems?.map((x) => x.path.split('/').slice(0, 2).join('/')) ?? [])
).sort()
@@ -133,7 +133,7 @@
let queryFilterKind = url.searchParams.get(TRIGGER_PATH_KIND_FILTER_SETTING)
let queryFilterUserFolders = url.searchParams.get(FILTER_USER_FOLDER_SETTING_NAME)
if (queryFilterKind) {
selectedFilterKind = queryFilterKind as 'route' | 'script_flow'
selectedFilterKind = queryFilterKind as 'trigger' | 'script_flow'
}
if (queryFilterUserFolders) {
filterUserFolders = queryFilterUserFolders == 'true'
@@ -174,7 +174,7 @@
<div class="flex flex-row items-center gap-2 mt-6">
<div class="text-sm shrink-0"> Filter by path of </div>
<ToggleButtonGroup bind:selected={selectedFilterKind}>
<ToggleButton small value="route" label="Route" icon={Route} />
<ToggleButton small value="trigger" label="Route" icon={Route} />
<ToggleButton small value="script_flow" label="Script/Flow" icon={Code} />
</ToggleButtonGroup>
</div>
@@ -172,13 +172,13 @@
? selectedFilterKind === 'schedule'
? schedules?.filter(
(x) =>
x.path.startsWith(ownerFilter + '/' ?? '') &&
x.path.startsWith(ownerFilter + '/') &&
filterItemsPathsBaseOnUserFilters(x, selectedFilterKind, filterUserFolders) &&
filterItemsBasedOnEnabledDisabled(x, filterEnabledDisabled)
)
: schedules?.filter(
(x) =>
x.script_path.startsWith(ownerFilter + '/' ?? '') &&
x.script_path.startsWith(ownerFilter + '/') &&
filterItemsPathsBaseOnUserFilters(x, selectedFilterKind, filterUserFolders) &&
filterItemsBasedOnEnabledDisabled(x, filterEnabledDisabled)
)
@@ -84,6 +84,7 @@
import json from 'svelte-highlight/languages/json'
import { writable } from 'svelte/store'
import TriggersBadge from '$lib/components/graph/renderers/triggers/TriggersBadge.svelte'
import WebsocketTriggersPanel from '$lib/components/triggers/WebsocketTriggersPanel.svelte'
let script: Script | undefined
let topHash: string | undefined
@@ -727,6 +728,11 @@
<RoutesPanel path={script.path ?? ''} isFlow={false} />
</div>
</svelte:fragment>
<svelte:fragment slot="websockets">
<div class="p-2">
<WebsocketTriggersPanel path={script.path ?? ''} isFlow={false} />
</div>
</svelte:fragment>
<svelte:fragment slot="emails">
<div class="p-2">
<EmailTriggerPanel
@@ -0,0 +1,5 @@
export function load() {
return {
stuff: { title: 'WS triggers' }
}
}
@@ -0,0 +1,414 @@
<script lang="ts">
import { WebsocketTriggerService, type WebsocketTrigger } from '$lib/gen'
import {
canWrite,
displayDate,
getLocalSetting,
sendUserToast,
storeLocalSetting
} from '$lib/utils'
import { base } from '$app/paths'
import CenteredPage from '$lib/components/CenteredPage.svelte'
import { Button, Skeleton } from '$lib/components/common'
import Dropdown from '$lib/components/DropdownV2.svelte'
import PageHeader from '$lib/components/PageHeader.svelte'
import SharedBadge from '$lib/components/SharedBadge.svelte'
import ShareModal from '$lib/components/ShareModal.svelte'
import Toggle from '$lib/components/Toggle.svelte'
import { userStore, workspaceStore } from '$lib/stores'
import { Unplug, Code, Eye, Pen, Plus, Share, Trash, Circle } from 'lucide-svelte'
import { goto } from '$lib/navigation'
import SearchItems from '$lib/components/SearchItems.svelte'
import NoItemFound from '$lib/components/home/NoItemFound.svelte'
import RowIcon from '$lib/components/common/table/RowIcon.svelte'
import ListFilters from '$lib/components/home/ListFilters.svelte'
import ToggleButtonGroup from '$lib/components/common/toggleButton-v2/ToggleButtonGroup.svelte'
import ToggleButton from '$lib/components/common/toggleButton-v2/ToggleButton.svelte'
import { setQuery } from '$lib/navigation'
import { onDestroy, onMount } from 'svelte'
import WebsocketTriggerEditor from '$lib/components/triggers/WebsocketTriggerEditor.svelte'
import Popover from '$lib/components/Popover.svelte'
type TriggerW = WebsocketTrigger & { canWrite: boolean }
let triggers: TriggerW[] = []
let shareModal: ShareModal
let loading = true
async function loadTriggers(): Promise<void> {
triggers = (
await WebsocketTriggerService.listWebsocketTriggers({ workspace: $workspaceStore! })
).map((x) => {
return { canWrite: canWrite(x.path, x.extra_perms!, $userStore), ...x }
})
loading = false
}
let status: {
[path: string]: { error: string | undefined; last_server_ping: string | undefined }
} = {}
let interval = setInterval(async () => {
try {
status = (
await WebsocketTriggerService.listWebsocketTriggers({
workspace: $workspaceStore!
})
).reduce((acc, x) => {
acc[x.path] = x
return acc
}, {})
} catch (err) {
console.error(err)
}
}, 5000)
onDestroy(() => {
clearInterval(interval)
})
async function setTriggerEnabled(path: string, enabled: boolean): Promise<void> {
try {
await WebsocketTriggerService.setWebsocketTriggerEnabled({
path,
workspace: $workspaceStore!,
requestBody: { enabled }
})
} catch (err) {
sendUserToast(
`Cannot ` + (enabled ? 'enable' : 'disable') + ` websocket trigger: ${err.body}`,
true
)
} finally {
loadTriggers()
}
}
$: {
if ($workspaceStore && $userStore) {
loadTriggers()
}
}
let websocketTriggerEditor: WebsocketTriggerEditor
let filteredItems: (TriggerW & { marked?: any })[] | undefined = []
let items: typeof filteredItems | undefined = []
let preFilteredItems: typeof filteredItems | undefined = []
let filter = ''
let ownerFilter: string | undefined = undefined
let nbDisplayed = 15
const TRIGGER_PATH_KIND_FILTER_SETTING = 'filter_path_of'
const FILTER_USER_FOLDER_SETTING_NAME = 'user_and_folders_only'
let selectedFilterKind =
(getLocalSetting(TRIGGER_PATH_KIND_FILTER_SETTING) as 'trigger' | 'script_flow') ?? 'trigger'
let filterUserFolders = getLocalSetting(FILTER_USER_FOLDER_SETTING_NAME) == 'true'
$: storeLocalSetting(TRIGGER_PATH_KIND_FILTER_SETTING, selectedFilterKind)
$: storeLocalSetting(FILTER_USER_FOLDER_SETTING_NAME, filterUserFolders ? 'true' : undefined)
function filterItemsPathsBaseOnUserFilters(
item: TriggerW,
selectedFilterKind: 'trigger' | 'script_flow',
filterUserFolders: boolean
) {
if ($workspaceStore == 'admins') return true
if (filterUserFolders) {
if (selectedFilterKind === 'trigger') {
return (
!item.path.startsWith('u/') || item.path.startsWith('u/' + $userStore?.username + '/')
)
} else {
return (
!item.script_path.startsWith('u/') ||
item.script_path.startsWith('u/' + $userStore?.username + '/')
)
}
} else {
return true
}
}
$: preFilteredItems =
ownerFilter != undefined
? selectedFilterKind === 'trigger'
? triggers?.filter(
(x) =>
x.path.startsWith(ownerFilter + '/') &&
filterItemsPathsBaseOnUserFilters(x, selectedFilterKind, filterUserFolders)
)
: triggers?.filter(
(x) =>
x.script_path.startsWith(ownerFilter + '/') &&
filterItemsPathsBaseOnUserFilters(x, selectedFilterKind, filterUserFolders)
)
: triggers?.filter((x) =>
filterItemsPathsBaseOnUserFilters(x, selectedFilterKind, filterUserFolders)
)
$: if ($workspaceStore) {
ownerFilter = undefined
}
$: owners =
selectedFilterKind === 'trigger'
? Array.from(
new Set(filteredItems?.map((x) => x.path.split('/').slice(0, 2).join('/')) ?? [])
).sort()
: Array.from(
new Set(filteredItems?.map((x) => x.script_path.split('/').slice(0, 2).join('/')) ?? [])
).sort()
$: items = filter !== '' ? filteredItems : preFilteredItems
function updateQueryFilters(selectedFilterKind, filterUserFolders) {
setQuery(
new URL(window.location.href),
TRIGGER_PATH_KIND_FILTER_SETTING,
selectedFilterKind
).then(() => {
setQuery(
new URL(window.location.href),
FILTER_USER_FOLDER_SETTING_NAME,
String(filterUserFolders)
)
})
}
function loadQueryFilters() {
let url = new URL(window.location.href)
let queryFilterKind = url.searchParams.get(TRIGGER_PATH_KIND_FILTER_SETTING)
let queryFilterUserFolders = url.searchParams.get(FILTER_USER_FOLDER_SETTING_NAME)
if (queryFilterKind) {
selectedFilterKind = queryFilterKind as 'trigger' | 'script_flow'
}
if (queryFilterUserFolders) {
filterUserFolders = queryFilterUserFolders == 'true'
}
}
onMount(() => {
loadQueryFilters()
})
$: updateQueryFilters(selectedFilterKind, filterUserFolders)
</script>
<WebsocketTriggerEditor on:update={loadTriggers} bind:this={websocketTriggerEditor} />
<SearchItems
{filter}
items={preFilteredItems}
bind:filteredItems
f={(x) => (x.summary ?? '') + ' ' + x.path + ' (' + x.script_path + ')'}
/>
<CenteredPage>
<PageHeader
title="Websocket triggers"
tooltip="Windmill can listen to websocket events and trigger scripts or flows based on them."
>
<Button
size="md"
startIcon={{ icon: Plus }}
on:click={() => websocketTriggerEditor.openNew(false)}
>
New&nbsp;WS trigger
</Button>
</PageHeader>
<div class="w-full h-full flex flex-col">
<div class="w-full pb-4 pt-6">
<input type="text" placeholder="Search WS triggers" bind:value={filter} class="search-item" />
<div class="flex flex-row items-center gap-2 mt-6">
<div class="text-sm shrink-0"> Filter by path of </div>
<ToggleButtonGroup bind:selected={selectedFilterKind}>
<ToggleButton small value="trigger" label="WS Trigger" icon={Unplug} />
<ToggleButton small value="script_flow" label="Script/Flow" icon={Code} />
</ToggleButtonGroup>
</div>
<ListFilters syncQuery bind:selectedFilter={ownerFilter} filters={owners} />
<div class="flex flex-row items-center justify-end gap-4">
{#if $userStore?.is_super_admin && $userStore.username.includes('@')}
<Toggle size="xs" bind:checked={filterUserFolders} options={{ right: 'Only f/*' }} />
{:else if $userStore?.is_admin || $userStore?.is_super_admin}
<Toggle
size="xs"
bind:checked={filterUserFolders}
options={{ right: `Only u/${$userStore.username} and f/*` }}
/>
{/if}
</div>
</div>
{#if loading}
{#each new Array(6) as _}
<Skeleton layout={[[6], 0.4]} />
{/each}
{:else if !triggers?.length}
<div class="text-center text-sm text-tertiary mt-2"> No websocket 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, url, is_flow, extra_perms, canWrite, marked, error, last_server_ping, enabled } (path)}
{@const href = `${is_flow ? '/flows/get' : '/scripts/get'}/${script_path}`}
{@const wsStatus = status[path] ?? { error, last_server_ping }}
<div
class="hover:bg-surface-hover w-full items-center px-4 py-2 gap-4 first-of-type:!border-t-0
first-of-type:rounded-t-md last-of-type:rounded-b-md flex flex-col"
>
<div class="w-full flex gap-5 items-center">
<RowIcon kind={is_flow ? 'flow' : 'script'} />
<a
href="#{path}"
on:click={() => websocketTriggerEditor?.openEdit(path, is_flow)}
class="min-w-0 grow hover:underline decoration-gray-400"
>
<div class="text-primary flex-wrap text-left text-md font-semibold mb-1 truncate">
{#if marked}
<span class="text-xs">
{@html marked}
</span>
{:else}
{url}
{/if}
</div>
<div class="text-secondary text-xs truncate text-left font-light">
{path}
</div>
<div class="text-secondary text-xs truncate text-left font-light">
runnable: {script_path}
</div>
</a>
<div class="hidden lg:flex flex-row gap-1 items-center">
<SharedBadge {canWrite} extraPerms={extra_perms} />
</div>
<div class="w-10">
{#if enabled}
{@const ping = wsStatus.last_server_ping
? new Date(wsStatus.last_server_ping)
: undefined}
{#if !ping || ping.getTime() < new Date().getTime() - 15 * 1000 || wsStatus.error}
<Popover notClickable>
<span class="flex h-4 w-4">
<Circle
class="text-red-600 animate-ping absolute inline-flex fill-current"
size={12}
/>
<Circle class="text-red-600 relative inline-flex fill-current" size={12} />
</span>
<div slot="text">
Websocket is not connected{wsStatus.error ? ': ' + wsStatus.error : ''}
</div>
</Popover>
{:else}
<Popover notClickable>
<span class="flex h-4 w-4">
<Circle
class="text-green-600 relative inline-flex fill-current"
size={12}
/>
</span>
<div slot="text"> Websocket is connected </div>
</Popover>
{/if}
{/if}
</div>
<Toggle
checked={enabled}
disabled={!canWrite}
on:change={(e) => {
setTriggerEnabled(path, e.detail)
}}
/>
<div class="flex gap-2 items-center justify-end">
<Button
on:click={() => websocketTriggerEditor?.openEdit(path, is_flow)}
size="xs"
startIcon={canWrite
? { icon: Pen }
: {
icon: Eye
}}
color="dark"
>
{canWrite ? 'Edit' : 'View'}
</Button>
<Dropdown
items={[
{
displayName: `View ${is_flow ? 'Flow' : 'Script'}`,
icon: Eye,
action: () => {
goto(href)
}
},
{
displayName: 'Delete',
type: 'delete',
icon: Trash,
disabled: !canWrite,
action: async () => {
await WebsocketTriggerService.deleteWebsocketTrigger({
workspace: $workspaceStore ?? '',
path
})
loadTriggers()
}
},
{
displayName: canWrite ? 'Edit' : 'View',
icon: canWrite ? Pen : Eye,
action: () => {
websocketTriggerEditor?.openEdit(path, is_flow)
}
},
{
displayName: 'Audit logs',
icon: Eye,
href: `${base}/audit_logs?resource=${path}`
},
{
displayName: canWrite ? 'Share' : 'See Permissions',
icon: Share,
action: () => {
shareModal.openDrawer(path, 'websocket_trigger')
}
}
]}
/>
</div>
</div>
<div class="w-full flex justify-between items-baseline">
<div
class="flex flex-wrap text-[0.7em] text-tertiary gap-1 items-center justify-end truncate pr-2"
><div class="truncate">edited by {edited_by}</div><div class="truncate"
>the {displayDate(edited_at)}</div
></div
></div
>
</div>
{/each}
</div>
{:else}
<NoItemFound />
{/if}
</div>
{#if items && items?.length > 15 && nbDisplayed < items.length}
<span class="text-xs"
>{nbDisplayed} items out of {items.length}
<button class="ml-4" on:click={() => (nbDisplayed += 30)}>load 30 more</button></span
>
{/if}
</CenteredPage>
<ShareModal
bind:this={shareModal}
on:change={() => {
loadTriggers()
}}
/>
+3 -3
View File
@@ -73,9 +73,9 @@
const selectedIdStore = writable('settings-metadata')
const primaryScheduleStore = writable<ScheduleTrigger | undefined | false>(undefined)
const triggersCount = writable<TriggersCount | undefined>(undefined)
const selectedTriggerStore = writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes'>(
'webhooks'
)
const selectedTriggerStore = writable<
'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets'
>('webhooks')
setContext<TriggerContext>('TriggerContext', {
primarySchedule: primaryScheduleStore,
selectedTrigger: selectedTriggerStore,