diff --git a/backend/.sqlx/query-02424907504848e983bfa89eec343061932dc5b4b17cf13d5cf8d833aedbe6d5.json b/backend/.sqlx/query-02424907504848e983bfa89eec343061932dc5b4b17cf13d5cf8d833aedbe6d5.json new file mode 100644 index 0000000000..c8b5e3086f --- /dev/null +++ b/backend/.sqlx/query-02424907504848e983bfa89eec343061932dc5b4b17cf13d5cf8d833aedbe6d5.json @@ -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" +} diff --git a/backend/.sqlx/query-12a86755706ce030a0a9142da78d035d0c7d361240b60005cc864e4345eb0bc7.json b/backend/.sqlx/query-12a86755706ce030a0a9142da78d035d0c7d361240b60005cc864e4345eb0bc7.json new file mode 100644 index 0000000000..dc357c2169 --- /dev/null +++ b/backend/.sqlx/query-12a86755706ce030a0a9142da78d035d0c7d361240b60005cc864e4345eb0bc7.json @@ -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" +} diff --git a/backend/.sqlx/query-9ebb9c16948a695a053068d9a1df0152691271ac33d8c82db47be11746ccbcef.json b/backend/.sqlx/query-2238aaed46031f71b17c1cabcef935a0d3b2ca6d057f446ed9b02d4386e2ddd9.json similarity index 58% rename from backend/.sqlx/query-9ebb9c16948a695a053068d9a1df0152691271ac33d8c82db47be11746ccbcef.json rename to backend/.sqlx/query-2238aaed46031f71b17c1cabcef935a0d3b2ca6d057f446ed9b02d4386e2ddd9.json index ae30d3bd98..aa32886ba8 100644 --- a/backend/.sqlx/query-9ebb9c16948a695a053068d9a1df0152691271ac33d8c82db47be11746ccbcef.json +++ b/backend/.sqlx/query-2238aaed46031f71b17c1cabcef935a0d3b2ca6d057f446ed9b02d4386e2ddd9.json @@ -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" } diff --git a/backend/.sqlx/query-2e6165543e34216dfaedf6e10729f733cff812e62dea1c1bce8cefd3a3979b14.json b/backend/.sqlx/query-2e6165543e34216dfaedf6e10729f733cff812e62dea1c1bce8cefd3a3979b14.json new file mode 100644 index 0000000000..4a15e2d4fa --- /dev/null +++ b/backend/.sqlx/query-2e6165543e34216dfaedf6e10729f733cff812e62dea1c1bce8cefd3a3979b14.json @@ -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" +} diff --git a/backend/.sqlx/query-4e9668a46bad9e82baa51422946d373b18b6577198df7545c94bd19be3446775.json b/backend/.sqlx/query-4e9668a46bad9e82baa51422946d373b18b6577198df7545c94bd19be3446775.json new file mode 100644 index 0000000000..2f96d31ecb --- /dev/null +++ b/backend/.sqlx/query-4e9668a46bad9e82baa51422946d373b18b6577198df7545c94bd19be3446775.json @@ -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" +} diff --git a/backend/.sqlx/query-acbf74cf3302bfcf7615285070d3f8958932bb8a2dda715f1b9152ab44442780.json b/backend/.sqlx/query-acbf74cf3302bfcf7615285070d3f8958932bb8a2dda715f1b9152ab44442780.json new file mode 100644 index 0000000000..eef46ce41d --- /dev/null +++ b/backend/.sqlx/query-acbf74cf3302bfcf7615285070d3f8958932bb8a2dda715f1b9152ab44442780.json @@ -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" +} diff --git a/backend/.sqlx/query-d490fef418e8567fa40aad60e5d46f233ce478086430d8070fe7b2b8a4f9580e.json b/backend/.sqlx/query-d490fef418e8567fa40aad60e5d46f233ce478086430d8070fe7b2b8a4f9580e.json new file mode 100644 index 0000000000..07d30c5881 --- /dev/null +++ b/backend/.sqlx/query-d490fef418e8567fa40aad60e5d46f233ce478086430d8070fe7b2b8a4f9580e.json @@ -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" +} diff --git a/backend/.sqlx/query-d54840373df5da9662ed11ed0a605cddac517bdc06dca0a8be071330338e948a.json b/backend/.sqlx/query-d54840373df5da9662ed11ed0a605cddac517bdc06dca0a8be071330338e948a.json new file mode 100644 index 0000000000..5b298f8db3 --- /dev/null +++ b/backend/.sqlx/query-d54840373df5da9662ed11ed0a605cddac517bdc06dca0a8be071330338e948a.json @@ -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" +} diff --git a/backend/.sqlx/query-db24f1f6ef1e26b7de7926bd3f32d1fa673c4970c25ff9ad98f0fdf199801b53.json b/backend/.sqlx/query-db24f1f6ef1e26b7de7926bd3f32d1fa673c4970c25ff9ad98f0fdf199801b53.json new file mode 100644 index 0000000000..70f4b91b71 --- /dev/null +++ b/backend/.sqlx/query-db24f1f6ef1e26b7de7926bd3f32d1fa673c4970c25ff9ad98f0fdf199801b53.json @@ -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" +} diff --git a/backend/.sqlx/query-f2baee15e6d1fecd6d2d7b39fda1b50a15ec8bd349f1081af54da5dc2f5e3021.json b/backend/.sqlx/query-f2baee15e6d1fecd6d2d7b39fda1b50a15ec8bd349f1081af54da5dc2f5e3021.json new file mode 100644 index 0000000000..2e75b4e74b --- /dev/null +++ b/backend/.sqlx/query-f2baee15e6d1fecd6d2d7b39fda1b50a15ec8bd349f1081af54da5dc2f5e3021.json @@ -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" +} diff --git a/backend/.sqlx/query-febd70e6d5304ad912c232d3d8fc3f8ce0133d8a1a9e5ca1124eff9a988dc2d7.json b/backend/.sqlx/query-febd70e6d5304ad912c232d3d8fc3f8ce0133d8a1a9e5ca1124eff9a988dc2d7.json new file mode 100644 index 0000000000..7958f9d2f3 --- /dev/null +++ b/backend/.sqlx/query-febd70e6d5304ad912c232d3d8fc3f8ce0133d8a1a9e5ca1124eff9a988dc2d7.json @@ -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" +} diff --git a/backend/Cargo.lock b/backend/Cargo.lock index c02d6788c4..8e6649aa16 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -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", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index ac9026769d..9596afef69 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -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"] } \ No newline at end of file diff --git a/backend/migrations/20241002163207_add_websocket_triggers.down.sql b/backend/migrations/20241002163207_add_websocket_triggers.down.sql new file mode 100644 index 0000000000..bb579574a8 --- /dev/null +++ b/backend/migrations/20241002163207_add_websocket_triggers.down.sql @@ -0,0 +1,2 @@ +-- Add down migration script here +DROP TABLE websocket_trigger; \ No newline at end of file diff --git a/backend/migrations/20241002163207_add_websocket_triggers.up.sql b/backend/migrations/20241002163207_add_websocket_triggers.up.sql new file mode 100644 index 0000000000..4d6b44cb8f --- /dev/null +++ b/backend/migrations/20241002163207_add_websocket_triggers.up.sql @@ -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)); \ No newline at end of file diff --git a/backend/windmill-api/Cargo.toml b/backend/windmill-api/Cargo.toml index e4b6b5d71b..b7c4a41aed 100644 --- a/backend/windmill-api/Cargo.toml +++ b/backend/windmill-api/Cargo.toml @@ -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 diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index b4611bdd3d..bc19dda924 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -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 diff --git a/backend/windmill-api/src/http_triggers.rs b/backend/windmill-api/src/http_triggers.rs index 9e30743fbc..61757cb164 100644 --- a/backend/windmill-api/src/http_triggers.rs +++ b/backend/windmill-api/src/http_triggers.rs @@ -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, Path(w_id): Path) -> JsonResult { - 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, 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 { - 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, diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 817a6f0d02..53cd958c6f 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -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) diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index 78108a61a4..342ecd70a1 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -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( diff --git a/backend/windmill-api/src/triggers.rs b/backend/windmill-api/src/triggers.rs index dbc5306027..337a20364b 100644 --- a/backend/windmill-api/src/triggers.rs +++ b/backend/windmill-api/src/triggers.rs @@ -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, })) } diff --git a/backend/windmill-api/src/websocket_triggers.rs b/backend/windmill-api/src/websocket_triggers.rs new file mode 100644 index 0000000000..3258f87261 --- /dev/null +++ b/backend/windmill-api/src/websocket_triggers.rs @@ -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, + filters: Vec, +} + +#[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, + server_id: Option, + last_server_ping: Option>, + extra_perms: serde_json::Value, + error: Option, + enabled: bool, + filters: Vec, +} + +#[derive(Deserialize)] +struct EditWebsocketTrigger { + path: String, + url: String, + script_path: String, + is_flow: bool, + filters: Vec, +} + +#[derive(Deserialize)] +pub struct ListWebsocketTriggerQuery { + pub page: Option, + pub per_page: Option, + pub path: Option, + pub is_flow: Option, + pub path_start: Option, +} + +async fn list_websocket_triggers( + authed: ApiAuthed, + Extension(user_db): Extension, + Path(w_id): Path, + Query(lst): Query, +) -> error::JsonResult> { + 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, + Path((w_id, path)): Path<(String, StripPath)>, +) -> error::JsonResult { + 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, + Path(w_id): Path, + Json(ct): Json, +) -> 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, + Path((w_id, path)): Path<(String, StripPath)>, + Json(ct): Json, +) -> error::Result { + 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, + Path((w_id, path)): Path<(String, StripPath)>, + Json(payload): Json, +) -> error::Result { + 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, + Path((w_id, path)): Path<(String, StripPath)>, +) -> error::Result { + 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, + Path((w_id, path)): Path<(String, StripPath)>, +) -> JsonResult { + 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) -> () { + 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, +) -> () { + 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(self, mut map: V) -> Result + where + V: MapAccess<'de>, + { + while let Some(key) = map.next_key::()? { + 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::()?; + } + } + // 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 +where + D: Deserializer<'de>, +{ + deserializer.deserialize_map(SupersetVisitor { key, value_to_check }) +} + +async fn listen_to_websocket( + ws_trigger: WebsocketTrigger, + db: DB, + rsmq: Option, +) -> () { + 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 = 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, + 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(()) +} diff --git a/backend/windmill-api/src/workspaces.rs b/backend/windmill-api/src/workspaces.rs index 95cd97fad0..cf48781331 100644 --- a/backend/windmill-api/src/workspaces.rs +++ b/backend/windmill-api/src/workspaces.rs @@ -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, + Path(w_id): Path, +) -> JsonResult { + 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, diff --git a/frontend/src/lib/components/Dev.svelte b/frontend/src/lib/components/Dev.svelte index 158f7bb762..5e3b52d867 100644 --- a/frontend/src/lib/components/Dev.svelte +++ b/frontend/src/lib/components/Dev.svelte @@ -475,9 +475,9 @@ const testStepStore = writable>({}) 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(undefined) const triggersCount = writable(undefined) diff --git a/frontend/src/lib/components/FlowBuilder.svelte b/frontend/src/lib/components/FlowBuilder.svelte index 49d7055219..6da81b66d3 100644 --- a/frontend/src/lib/components/FlowBuilder.svelte +++ b/frontend/src/lib/components/FlowBuilder.svelte @@ -461,9 +461,9 @@ } const selectedIdStore = writable(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) } diff --git a/frontend/src/lib/components/Path.svelte b/frontend/src/lib/components/Path.svelte index c6d4e916a2..2d112387af 100644 --- a/frontend/src/lib/components/Path.svelte +++ b/frontend/src/lib/components/Path.svelte @@ -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 @@ /> -
{error}
+
{error}
- {#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} 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 diff --git a/frontend/src/lib/components/ScriptBuilder.svelte b/frontend/src/lib/components/ScriptBuilder.svelte index 8be6e5104f..45a48e66d8 100644 --- a/frontend/src/lib/components/ScriptBuilder.svelte +++ b/frontend/src/lib/components/ScriptBuilder.svelte @@ -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) diff --git a/frontend/src/lib/components/ShareModal.svelte b/frontend/src/lib/components/ShareModal.svelte index 23605ce132..8f0192d5e4 100644 --- a/frontend/src/lib/components/ShareModal.svelte +++ b/frontend/src/lib/components/ShareModal.svelte @@ -25,6 +25,7 @@ | 'app' | 'raw_app' | 'http_trigger' + | 'websocket_trigger' let kind: Kind let path: string = '' diff --git a/frontend/src/lib/components/details/DetailPageDetailPanel.svelte b/frontend/src/lib/components/details/DetailPageDetailPanel.svelte index 4ce13fca2f..a5927c2938 100644 --- a/frontend/src/lib/components/details/DetailPageDetailPanel.svelte +++ b/frontend/src/lib/components/details/DetailPageDetailPanel.svelte @@ -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 @@ + diff --git a/frontend/src/lib/components/details/DetailPageLayout.svelte b/frontend/src/lib/components/details/DetailPageLayout.svelte index 60022e4e62..8e9dc119a0 100644 --- a/frontend/src/lib/components/details/DetailPageLayout.svelte +++ b/frontend/src/lib/components/details/DetailPageLayout.svelte @@ -19,9 +19,9 @@ let clientWidth = window.innerWidth const primaryScheduleStore = writable(undefined) - const selectedTriggerStore = writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes'>( - 'webhooks' - ) + const selectedTriggerStore = writable< + 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets' + >('webhooks') setContext('TriggerContext', { selectedTrigger: selectedTriggerStore, @@ -48,6 +48,7 @@ > + @@ -81,6 +82,7 @@ + diff --git a/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte b/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte index cbf1e704b0..71824b14d0 100644 --- a/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte +++ b/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte @@ -1,10 +1,16 @@ @@ -28,6 +34,12 @@ HTTP + + + + Websockets + + @@ -52,6 +64,8 @@ {:else if triggerSelected === 'schedules'} + {:else if triggerSelected === 'websockets'} + {:else if triggerSelected === 'cli'} {/if} diff --git a/frontend/src/lib/components/graph/renderers/triggers/TriggersBadge.svelte b/frontend/src/lib/components/graph/renderers/triggers/TriggersBadge.svelte index 471efae4c9..ae4be3a994 100644 --- a/frontend/src/lib/components/graph/renderers/triggers/TriggersBadge.svelte +++ b/frontend/src/lib/components/graph/renderers/triggers/TriggersBadge.svelte @@ -1,5 +1,5 @@ + +{#if open} + +{/if} diff --git a/frontend/src/lib/components/triggers/WebsocketTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/WebsocketTriggerEditorInner.svelte new file mode 100644 index 0000000000..39f9ce3cae --- /dev/null +++ b/frontend/src/lib/components/triggers/WebsocketTriggerEditorInner.svelte @@ -0,0 +1,338 @@ + + + + + + {#if !drawerLoading && can_write} + {#if edit} +
+ { + await WebsocketTriggerService.setWebsocketTriggerEnabled({ + path: initialPath, + workspace: $workspaceStore ?? '', + requestBody: { enabled: e.detail } + }) + sendUserToast( + `${e.detail ? 'enabled' : 'disabled'} websocket trigger ${initialPath}` + ) + }} + /> +
+ {/if} + + {/if} +
+ {#if drawerLoading} + + {:else} + + {#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} + +
+
+ +
+ +
+
+ +
+
+ +
+

+ Pick a script or flow to be triggered +

+
+ +
+
+ +
+

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

+
+ {#each filters as v, i} +
+
+
+ +
+ + + +
+ +
+ {/each} + +
+ +
+
+
+
+ {/if} +
+
diff --git a/frontend/src/lib/components/triggers/WebsocketTriggersPanel.svelte b/frontend/src/lib/components/triggers/WebsocketTriggersPanel.svelte new file mode 100644 index 0000000000..3369e87899 --- /dev/null +++ b/frontend/src/lib/components/triggers/WebsocketTriggersPanel.svelte @@ -0,0 +1,103 @@ + + + { + loadTriggers() + }} + bind:this={wsTriggerEditor} +/> + +
+ {#if !newItem} + {#if $userStore?.is_admin || $userStore?.is_super_admin} + + {:else} + + {/if} + {/if} + + {#if wsTriggers} + {#if wsTriggers.length == 0} +
No WS triggers
+ {:else} +
+ {#each wsTriggers as wsTriggers (wsTriggers.path)} +
+
{wsTriggers.path}
+
+ {wsTriggers.url} +
+
+ +
+
+ {/each} +
+ {/if} + {:else} + + {/if} + + {#if newItem} + + Deploy the {isFlow ? 'flow' : 'script'} to add WS triggers. + + {/if} +
diff --git a/frontend/src/lib/script_helpers.ts b/frontend/src/lib/script_helpers.ts index b940aba39e..f584ee9f19 100644 --- a/frontend/src/lib/script_helpers.ts +++ b/frontend/src/lib/script_helpers.ts @@ -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( diff --git a/frontend/src/routes/(root)/(logged)/+layout.svelte b/frontend/src/routes/(root)/(logged)/+layout.svelte index 3eae42916f..ca7a98cfaf 100644 --- a/frontend/src/routes/(root)/(logged)/+layout.svelte +++ b/frontend/src/routes/(root)/(logged)/+layout.svelte @@ -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 { diff --git a/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte b/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte index dfb55daff5..64b003a8a8 100644 --- a/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte @@ -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 @@ + +
+ +
+
- 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 @@
Filter by path of
- +
diff --git a/frontend/src/routes/(root)/(logged)/schedules/+page.svelte b/frontend/src/routes/(root)/(logged)/schedules/+page.svelte index e9225d405d..1f8c9e84a5 100644 --- a/frontend/src/routes/(root)/(logged)/schedules/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/schedules/+page.svelte @@ -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) ) diff --git a/frontend/src/routes/(root)/(logged)/scripts/get/[...hash]/+page.svelte b/frontend/src/routes/(root)/(logged)/scripts/get/[...hash]/+page.svelte index 152f2dde3a..de4a75fc4e 100644 --- a/frontend/src/routes/(root)/(logged)/scripts/get/[...hash]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/scripts/get/[...hash]/+page.svelte @@ -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 @@
+ +
+ +
+
+ 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 { + 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 { + 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) + + + + + (x.summary ?? '') + ' ' + x.path + ' (' + x.script_path + ')'} +/> + + + + + +
+
+ +
+
Filter by path of
+ + + + +
+ + +
+ {#if $userStore?.is_super_admin && $userStore.username.includes('@')} + + {:else if $userStore?.is_admin || $userStore?.is_super_admin} + + {/if} +
+
+ {#if loading} + {#each new Array(6) as _} + + {/each} + {:else if !triggers?.length} +
No websocket triggers
+ {:else if items?.length} +
+ {#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 }} + +
+
+ + + websocketTriggerEditor?.openEdit(path, is_flow)} + class="min-w-0 grow hover:underline decoration-gray-400" + > +
+ {#if marked} + + {@html marked} + + {:else} + {url} + {/if} +
+
+ {path} +
+
+ runnable: {script_path} +
+
+ + + +
+ {#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} + + + + + +
+ Websocket is not connected{wsStatus.error ? ': ' + wsStatus.error : ''} +
+
+ {:else} + + + + +
Websocket is connected
+
+ {/if} + {/if} +
+ + { + setTriggerEnabled(path, e.detail) + }} + /> + +
+ + { + 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') + } + } + ]} + /> +
+
+
+
edited by {edited_by}
the {displayDate(edited_at)}
+
+ {/each} +
+ {:else} + + {/if} +
+ {#if items && items?.length > 15 && nbDisplayed < items.length} + {nbDisplayed} items out of {items.length} + + {/if} +
+ + { + loadTriggers() + }} +/> diff --git a/frontend/src/routes/flows/dev/+page.svelte b/frontend/src/routes/flows/dev/+page.svelte index c2d2359816..2db4421e8a 100644 --- a/frontend/src/routes/flows/dev/+page.svelte +++ b/frontend/src/routes/flows/dev/+page.svelte @@ -73,9 +73,9 @@ const selectedIdStore = writable('settings-metadata') const primaryScheduleStore = writable(undefined) const triggersCount = writable(undefined) - const selectedTriggerStore = writable<'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes'>( - 'webhooks' - ) + const selectedTriggerStore = writable< + 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets' + >('webhooks') setContext('TriggerContext', { primarySchedule: primaryScheduleStore, selectedTrigger: selectedTriggerStore,