From 2a3c0d55fd5654ea8ea2cdf67645f90ecdbfc916 Mon Sep 17 00:00:00 2001 From: HugoCasa Date: Mon, 18 Nov 2024 17:03:50 +0100 Subject: [PATCH] feat: kafka triggers (#4713) * feat: kafka triggers * sqlx * fix build * improve error messages on windows * nit * fix build * missing action * nit * maybe fix ssl on windows * nit * update ee ref --------- Co-authored-by: Ruben Fiszel Co-authored-by: Ruben Fiszel --- .github/workflows/build-publish-rh-image.yml | 11 +- .github/workflows/build-staging-image.yml | 2 +- .github/workflows/build_windows_worker_.yml | 2 +- .github/workflows/docker-image.yml | 10 +- .github/workflows/publish_windows_worker.yml | 2 +- ...c5c32f582a3911efcf88e437aa90d7d5a49b5.json | 24 + ...81d870b847a8ba07e6197212fdf85ec901b09.json | 15 + ...4e82c1b7060ba0ca288dc6a7c7afcd76212e8.json | 16 + ...f5664130ae7500279055c08db42d67b2822f6.json | 23 + ...957c87f2307a218115ca749ea0e71eca2002e.json | 23 + ...973e8fc108d9f32be8a8501391052d76e191e.json | 25 ++ ...64814f95dbe5376f0539195d24a8ca96ca3e2.json | 26 ++ ...1fe73aba65cfdc8329c02b4d257de4d2d168a.json | 34 ++ ...160ad1962668fc3305d8e80ae91ef73614a80.json | 24 + ...2d1fa673c4970c25ff9ad98f0fdf199801b53.json | 28 -- backend/Cargo.lock | 78 +++- backend/Cargo.toml | 2 + backend/ee-repo-ref.txt | 2 +- .../20241112104946_kafka_triggers.down.sql | 2 + .../20241112104946_kafka_triggers.up.sql | 69 +++ backend/windmill-api/Cargo.toml | 2 + backend/windmill-api/openapi.yaml | 279 ++++++++++++ backend/windmill-api/src/apps.rs | 25 +- backend/windmill-api/src/granular_acls.rs | 4 +- backend/windmill-api/src/http_triggers.rs | 2 +- backend/windmill-api/src/kafka_triggers_ee.rs | 14 + backend/windmill-api/src/lib.rs | 33 +- backend/windmill-api/src/triggers.rs | 12 + backend/windmill-api/src/users.rs | 6 +- .../windmill-api/src/websocket_triggers.rs | 8 +- backend/windmill-api/src/workspaces.rs | 6 +- frontend/src/lib/components/Path.svelte | 9 +- frontend/src/lib/components/ShareModal.svelte | 1 + .../details/DetailPageDetailPanel.svelte | 1 + .../details/DetailPageLayout.svelte | 2 + .../details/DetailPageTriggerPanel.svelte | 10 + .../renderers/triggers/TriggersBadge.svelte | 12 +- .../src/lib/components/icons/KafkaIcon.svelte | 20 + .../components/sidebar/SidebarContent.svelte | 16 +- frontend/src/lib/components/triggers.ts | 9 +- .../triggers/KafkaTriggerEditor.svelte | 23 + .../triggers/KafkaTriggerEditorInner.svelte | 310 +++++++++++++ .../triggers/KafkaTriggersPanel.svelte | 108 +++++ .../components/triggers/TriggersEditor.svelte | 8 + .../WebsocketTriggerEditorInner.svelte | 4 +- frontend/src/lib/script_helpers.ts | 34 +- .../src/routes/(root)/(logged)/+layout.svelte | 11 +- .../(logged)/flows/get/[...path]/+page.svelte | 8 + .../(root)/(logged)/kafka_triggers/+page.js | 5 + .../(logged)/kafka_triggers/+page.svelte | 422 ++++++++++++++++++ .../scripts/get/[...hash]/+page.svelte | 6 + 51 files changed, 1736 insertions(+), 92 deletions(-) create mode 100644 backend/.sqlx/query-2139f1fb1877294bbf55d786000c5c32f582a3911efcf88e437aa90d7d5a49b5.json create mode 100644 backend/.sqlx/query-3997dcf2c11817e59bf50fd896381d870b847a8ba07e6197212fdf85ec901b09.json create mode 100644 backend/.sqlx/query-3e64c894c89ef82c4527180b82c4e82c1b7060ba0ca288dc6a7c7afcd76212e8.json create mode 100644 backend/.sqlx/query-5962733746c81480abe9ab3a6ccf5664130ae7500279055c08db42d67b2822f6.json create mode 100644 backend/.sqlx/query-6d418b5cd7a4df54cfe3ec06e2a957c87f2307a218115ca749ea0e71eca2002e.json create mode 100644 backend/.sqlx/query-6e70ebf078ac04a2933d2f83791973e8fc108d9f32be8a8501391052d76e191e.json create mode 100644 backend/.sqlx/query-6eb076afcf11dba0dfe35ee108c64814f95dbe5376f0539195d24a8ca96ca3e2.json create mode 100644 backend/.sqlx/query-cc8a71abdeccf0787695af32aa21fe73aba65cfdc8329c02b4d257de4d2d168a.json create mode 100644 backend/.sqlx/query-d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80.json delete mode 100644 backend/.sqlx/query-db24f1f6ef1e26b7de7926bd3f32d1fa673c4970c25ff9ad98f0fdf199801b53.json create mode 100644 backend/migrations/20241112104946_kafka_triggers.down.sql create mode 100644 backend/migrations/20241112104946_kafka_triggers.up.sql create mode 100644 backend/windmill-api/src/kafka_triggers_ee.rs create mode 100644 frontend/src/lib/components/icons/KafkaIcon.svelte create mode 100644 frontend/src/lib/components/triggers/KafkaTriggerEditor.svelte create mode 100644 frontend/src/lib/components/triggers/KafkaTriggerEditorInner.svelte create mode 100644 frontend/src/lib/components/triggers/KafkaTriggersPanel.svelte create mode 100644 frontend/src/routes/(root)/(logged)/kafka_triggers/+page.js create mode 100644 frontend/src/routes/(root)/(logged)/kafka_triggers/+page.svelte diff --git a/.github/workflows/build-publish-rh-image.yml b/.github/workflows/build-publish-rh-image.yml index b957ac5148..8d4ab5343f 100644 --- a/.github/workflows/build-publish-rh-image.yml +++ b/.github/workflows/build-publish-rh-image.yml @@ -3,8 +3,7 @@ env: IMAGE_NAME: ${{ github.repository }} name: Build and publish windmill for RHEL9 -on: - workflow_dispatch +on: workflow_dispatch permissions: write-all @@ -65,7 +64,7 @@ jobs: platforms: linux/amd64 push: true build-args: | - features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core + features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core,kafka secrets: | rh_username=${{ secrets.RH_USERNAME }} rh_password=${{ secrets.RH_PASSWORD }} @@ -74,7 +73,7 @@ jobs: labels: | ${{ steps.meta-ee-public.outputs.labels }}-amd64 org.opencontainers.image.licenses=Windmill-Enterprise-License - + - name: Build and push publicly ee arm64 uses: depot/build-push-action@v1 with: @@ -82,7 +81,7 @@ jobs: platforms: linux/arm64 push: true build-args: | - features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core + features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core,kafka secrets: | rh_username=${{ secrets.RH_USERNAME }} rh_password=${{ secrets.RH_PASSWORD }} @@ -108,7 +107,7 @@ jobs: run: | mv "${{ steps.extract-ee-amd64.outputs.destination }}/windmill" "${{ steps.extract-ee-amd64.outputs.destination }}/windmill-ee-amd64-rhel9" mv "${{ steps.extract-ee-arm64.outputs.destination }}/windmill" "${{ steps.extract-ee-arm64.outputs.destination }}/windmill-ee-arm64-rhel9" - + - uses: actions/upload-artifact@v4 with: name: RHEL9-amd64 build diff --git a/.github/workflows/build-staging-image.yml b/.github/workflows/build-staging-image.yml index 724fde4e59..197c219b09 100644 --- a/.github/workflows/build-staging-image.yml +++ b/.github/workflows/build-staging-image.yml @@ -62,7 +62,7 @@ jobs: platforms: linux/amd64,linux/arm64 push: true build-args: | - features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core + features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,deno_core,kafka tags: | ${{ steps.meta-ee-public.outputs.tags }} labels: | diff --git a/.github/workflows/build_windows_worker_.yml b/.github/workflows/build_windows_worker_.yml index 256ce0de1a..674a97e85f 100644 --- a/.github/workflows/build_windows_worker_.yml +++ b/.github/workflows/build_windows_worker_.yml @@ -45,7 +45,7 @@ jobs: $env:OPENSSL_DIR="${Env:VCPKG_INSTALLATION_ROOT}\installed\x64-windows-static" mkdir frontend/build && cd backend New-Item -Path . -Name "windmill-api/openapi-deref.yaml" -ItemType "File" -Force - cargo build --release --features=enterprise,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core + cargo build --release --features=enterprise,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core,kafka - name: Rename binary with corresponding architecture run: | diff --git a/.github/workflows/docker-image.yml b/.github/workflows/docker-image.yml index f1f3cae9f8..42cc21264b 100644 --- a/.github/workflows/docker-image.yml +++ b/.github/workflows/docker-image.yml @@ -1,10 +1,8 @@ env: REGISTRY: ghcr.io - IMAGE_NAME: - ${{ github.event_name != 'pull_request' && github.repository || + IMAGE_NAME: ${{ github.event_name != 'pull_request' && github.repository || 'windmill-labs/windmill-test' }} - DEV_SHA: - ${{ github.event_name != 'pull_request' && 'dev' || format('pr-{0}', + DEV_SHA: ${{ github.event_name != 'pull_request' && 'dev' || format('pr-{0}', github.event.number) }} name: Build windmill:main @@ -140,7 +138,7 @@ jobs: platforms: linux/amd64,linux/arm64 push: true build-args: | - features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core + features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core,kafka tags: | ${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-ee:${{ env.DEV_SHA }} ${{ steps.meta-ee-public.outputs.tags }} @@ -202,7 +200,7 @@ jobs: platforms: linux/amd64 push: true build-args: | - features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core + features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core,kafka PYTHON_IMAGE=python:3.12.2-slim-bookworm tags: | ${{ steps.meta-ee-public-py312.outputs.tags }} diff --git a/.github/workflows/publish_windows_worker.yml b/.github/workflows/publish_windows_worker.yml index 01c38dd2b1..f6d73a7621 100644 --- a/.github/workflows/publish_windows_worker.yml +++ b/.github/workflows/publish_windows_worker.yml @@ -47,7 +47,7 @@ jobs: $env:OPENSSL_DIR="${Env:VCPKG_INSTALLATION_ROOT}\installed\x64-windows-static" mkdir frontend/build && cd backend New-Item -Path . -Name "windmill-api/openapi-deref.yaml" -ItemType "File" -Force - cargo build --release --features=enterprise,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core + cargo build --release --features=enterprise,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core,kafka - name: Rename binary with corresponding architecture run: | diff --git a/backend/.sqlx/query-2139f1fb1877294bbf55d786000c5c32f582a3911efcf88e437aa90d7d5a49b5.json b/backend/.sqlx/query-2139f1fb1877294bbf55d786000c5c32f582a3911efcf88e437aa90d7d5a49b5.json new file mode 100644 index 0000000000..fef0a60ca7 --- /dev/null +++ b/backend/.sqlx/query-2139f1fb1877294bbf55d786000c5c32f582a3911efcf88e437aa90d7d5a49b5.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT COUNT(*) FROM kafka_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": "2139f1fb1877294bbf55d786000c5c32f582a3911efcf88e437aa90d7d5a49b5" +} diff --git a/backend/.sqlx/query-3997dcf2c11817e59bf50fd896381d870b847a8ba07e6197212fdf85ec901b09.json b/backend/.sqlx/query-3997dcf2c11817e59bf50fd896381d870b847a8ba07e6197212fdf85ec901b09.json new file mode 100644 index 0000000000..4c90c24dd6 --- /dev/null +++ b/backend/.sqlx/query-3997dcf2c11817e59bf50fd896381d870b847a8ba07e6197212fdf85ec901b09.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM kafka_trigger WHERE workspace_id = $1 AND path = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "3997dcf2c11817e59bf50fd896381d870b847a8ba07e6197212fdf85ec901b09" +} diff --git a/backend/.sqlx/query-3e64c894c89ef82c4527180b82c4e82c1b7060ba0ca288dc6a7c7afcd76212e8.json b/backend/.sqlx/query-3e64c894c89ef82c4527180b82c4e82c1b7060ba0ca288dc6a7c7afcd76212e8.json new file mode 100644 index 0000000000..8831155670 --- /dev/null +++ b/backend/.sqlx/query-3e64c894c89ef82c4527180b82c4e82c1b7060ba0ca288dc6a7c7afcd76212e8.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE kafka_trigger SET enabled = FALSE, error = $1, server_id = NULL, last_server_ping = NULL WHERE workspace_id = $2 AND path = $3", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "3e64c894c89ef82c4527180b82c4e82c1b7060ba0ca288dc6a7c7afcd76212e8" +} diff --git a/backend/.sqlx/query-5962733746c81480abe9ab3a6ccf5664130ae7500279055c08db42d67b2822f6.json b/backend/.sqlx/query-5962733746c81480abe9ab3a6ccf5664130ae7500279055c08db42d67b2822f6.json new file mode 100644 index 0000000000..0d64b110f8 --- /dev/null +++ b/backend/.sqlx/query-5962733746c81480abe9ab3a6ccf5664130ae7500279055c08db42d67b2822f6.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE kafka_trigger SET kafka_resource_path = $1, group_id = $2, topics = $3, script_path = $4, path = $5, is_flow = $6, edited_by = $7, email = $8, edited_at = now(), server_id = NULL, last_server_ping = NULL, error = NULL\n WHERE workspace_id = $9 AND path = $10", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "VarcharArray", + "Varchar", + "Varchar", + "Bool", + "Varchar", + "Varchar", + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "5962733746c81480abe9ab3a6ccf5664130ae7500279055c08db42d67b2822f6" +} diff --git a/backend/.sqlx/query-6d418b5cd7a4df54cfe3ec06e2a957c87f2307a218115ca749ea0e71eca2002e.json b/backend/.sqlx/query-6d418b5cd7a4df54cfe3ec06e2a957c87f2307a218115ca749ea0e71eca2002e.json new file mode 100644 index 0000000000..c20295fe97 --- /dev/null +++ b/backend/.sqlx/query-6d418b5cd7a4df54cfe3ec06e2a957c87f2307a218115ca749ea0e71eca2002e.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT EXISTS(SELECT 1 FROM kafka_trigger WHERE path = $1 AND workspace_id = $2)", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "exists", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "6d418b5cd7a4df54cfe3ec06e2a957c87f2307a218115ca749ea0e71eca2002e" +} diff --git a/backend/.sqlx/query-6e70ebf078ac04a2933d2f83791973e8fc108d9f32be8a8501391052d76e191e.json b/backend/.sqlx/query-6e70ebf078ac04a2933d2f83791973e8fc108d9f32be8a8501391052d76e191e.json new file mode 100644 index 0000000000..d95c8c9b5c --- /dev/null +++ b/backend/.sqlx/query-6e70ebf078ac04a2933d2f83791973e8fc108d9f32be8a8501391052d76e191e.json @@ -0,0 +1,25 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE kafka_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": "6e70ebf078ac04a2933d2f83791973e8fc108d9f32be8a8501391052d76e191e" +} diff --git a/backend/.sqlx/query-6eb076afcf11dba0dfe35ee108c64814f95dbe5376f0539195d24a8ca96ca3e2.json b/backend/.sqlx/query-6eb076afcf11dba0dfe35ee108c64814f95dbe5376f0539195d24a8ca96ca3e2.json new file mode 100644 index 0000000000..96825b9905 --- /dev/null +++ b/backend/.sqlx/query-6eb076afcf11dba0dfe35ee108c64814f95dbe5376f0539195d24a8ca96ca3e2.json @@ -0,0 +1,26 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE kafka_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": "6eb076afcf11dba0dfe35ee108c64814f95dbe5376f0539195d24a8ca96ca3e2" +} diff --git a/backend/.sqlx/query-cc8a71abdeccf0787695af32aa21fe73aba65cfdc8329c02b4d257de4d2d168a.json b/backend/.sqlx/query-cc8a71abdeccf0787695af32aa21fe73aba65cfdc8329c02b4d257de4d2d168a.json new file mode 100644 index 0000000000..c8f2b72b83 --- /dev/null +++ b/backend/.sqlx/query-cc8a71abdeccf0787695af32aa21fe73aba65cfdc8329c02b4d257de4d2d168a.json @@ -0,0 +1,34 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT \n EXISTS(SELECT 1 FROM websocket_trigger WHERE workspace_id = $1) as \"websocket_used!\", \n EXISTS(SELECT 1 FROM http_trigger WHERE workspace_id = $1) as \"http_routes_used!\",\n EXISTS(SELECT 1 FROM kafka_trigger WHERE workspace_id = $1) as \"kafka_used!\"", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "websocket_used!", + "type_info": "Bool" + }, + { + "ordinal": 1, + "name": "http_routes_used!", + "type_info": "Bool" + }, + { + "ordinal": 2, + "name": "kafka_used!", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Text" + ] + }, + "nullable": [ + null, + null, + null + ] + }, + "hash": "cc8a71abdeccf0787695af32aa21fe73aba65cfdc8329c02b4d257de4d2d168a" +} diff --git a/backend/.sqlx/query-d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80.json b/backend/.sqlx/query-d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80.json new file mode 100644 index 0000000000..271c395f9a --- /dev/null +++ b/backend/.sqlx/query-d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE kafka_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": "d6615719bf8db4b333ed55c9a3c160ad1962668fc3305d8e80ae91ef73614a80" +} diff --git a/backend/.sqlx/query-db24f1f6ef1e26b7de7926bd3f32d1fa673c4970c25ff9ad98f0fdf199801b53.json b/backend/.sqlx/query-db24f1f6ef1e26b7de7926bd3f32d1fa673c4970c25ff9ad98f0fdf199801b53.json deleted file mode 100644 index 70f4b91b71..0000000000 --- a/backend/.sqlx/query-db24f1f6ef1e26b7de7926bd3f32d1fa673c4970c25ff9ad98f0fdf199801b53.json +++ /dev/null @@ -1,28 +0,0 @@ -{ - "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/Cargo.lock b/backend/Cargo.lock index 6e7053f97d..71a05984ab 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -1378,7 +1378,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2593a3b8b938bd68373196c9832f516be11fa487ef4ae745eb282e6a56a7244" dependencies = [ "once_cell", - "proc-macro-crate", + "proc-macro-crate 3.2.0", "proc-macro2", "quote", "syn 2.0.87", @@ -5313,7 +5313,7 @@ dependencies = [ "darling 0.20.10", "heck 0.5.0", "num-bigint", - "proc-macro-crate", + "proc-macro-crate 3.2.0", "proc-macro-error2", "proc-macro2", "quote", @@ -5556,6 +5556,27 @@ dependencies = [ "libc", ] +[[package]] +name = "num_enum" +version = "0.5.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1f646caf906c20226733ed5b1374287eb97e3c2a5c227ce668c1f2ce20ae57c9" +dependencies = [ + "num_enum_derive", +] + +[[package]] +name = "num_enum_derive" +version = "0.5.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dcbff9bc912032c62bf65ef1d5aea88983b420f4f839db1e9b0c281a25c9c799" +dependencies = [ + "proc-macro-crate 1.3.1", + "proc-macro2", + "quote", + "syn 1.0.109", +] + [[package]] name = "number_prefix" version = "0.4.0" @@ -5736,6 +5757,15 @@ version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff011a302c396a5197692431fc1948019154afc178baf7d8e37367442a4601cf" +[[package]] +name = "openssl-src" +version = "300.4.0+3.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a709e02f2b4aca747929cca5ed248880847c650233cf8b8cdc48f40aaf4898a6" +dependencies = [ + "cc", +] + [[package]] name = "openssl-sys" version = "0.9.104" @@ -5744,6 +5774,7 @@ checksum = "45abf306cbf99debc8195b66b7346498d7b10c210de50418b5ccd7ceba08c741" dependencies = [ "cc", "libc", + "openssl-src", "pkg-config", "vcpkg", ] @@ -6239,6 +6270,16 @@ dependencies = [ "elliptic-curve", ] +[[package]] +name = "proc-macro-crate" +version = "1.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7f4c021e1093a56626774e81216a4ce732a735e5bad4868a03f3ed65ca0c3919" +dependencies = [ + "once_cell", + "toml_edit 0.19.15", +] + [[package]] name = "proc-macro-crate" version = "3.2.0" @@ -6711,6 +6752,38 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "rdkafka" +version = "0.36.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1beea247b9a7600a81d4cc33f659ce1a77e1988323d7d2809c7ed1c21f4c316d" +dependencies = [ + "futures-channel", + "futures-util", + "libc", + "log", + "rdkafka-sys", + "serde", + "serde_derive", + "serde_json", + "slab", + "tokio", +] + +[[package]] +name = "rdkafka-sys" +version = "4.7.0+2.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "55e0d2f9ba6253f6ec72385e453294f8618e9e15c2c6aba2a5c01ccf9622d615" +dependencies = [ + "cmake", + "libc", + "libz-sys", + "num_enum", + "openssl-sys", + "pkg-config", +] + [[package]] name = "reborrow" version = "0.5.5" @@ -10695,6 +10768,7 @@ dependencies = [ "prometheus", "quick_cache", "rand 0.8.5", + "rdkafka", "regex", "reqwest 0.12.9", "rsa 0.7.2", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index ff795b3954..a14e5be0a1 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -62,6 +62,7 @@ jemalloc = ["windmill-common/jemalloc", "dep:tikv-jemallocator", "dep:tikv-jemal tantivy = ["dep:windmill-indexer", "windmill-api/tantivy"] sqlx = ["windmill-worker/sqlx"] deno_core = ["windmill-worker/deno_core", "dep:deno_core"] +kafka = ["windmill-api/kafka"] [dependencies] anyhow.workspace = true @@ -266,6 +267,7 @@ tokio-native-tls = "^0" openssl = "=0.10" mail-parser = "^0" matchit = "=0.7.3" +rdkafka = { version = "0.36.2", features = ["cmake-build", "ssl-vendored"] } datafusion = "39.0.0" object_store = { version = "0.10.0", features = ["aws", "azure"] } diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index caed423d9a..92fb7bac27 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -51dcbf93b0d127af9f33fa346cc63fcd2475d4fa +7cdaac656fde3f3e267dfb9e94dda7f3843aa44e \ No newline at end of file diff --git a/backend/migrations/20241112104946_kafka_triggers.down.sql b/backend/migrations/20241112104946_kafka_triggers.down.sql new file mode 100644 index 0000000000..27cbeb9c50 --- /dev/null +++ b/backend/migrations/20241112104946_kafka_triggers.down.sql @@ -0,0 +1,2 @@ +-- Add down migration script here +DROP TABLE kafka_trigger; \ No newline at end of file diff --git a/backend/migrations/20241112104946_kafka_triggers.up.sql b/backend/migrations/20241112104946_kafka_triggers.up.sql new file mode 100644 index 0000000000..9a4a985e42 --- /dev/null +++ b/backend/migrations/20241112104946_kafka_triggers.up.sql @@ -0,0 +1,69 @@ +-- Add up migration script here +-- Add up migration script here + +CREATE TABLE kafka_trigger ( + path VARCHAR(255) NOT NULL, + kafka_resource_path VARCHAR(255) NOT NULL, + topics VARCHAR(255)[] NOT NULL, + group_id 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, + PRIMARY KEY (path, workspace_id) +); + +GRANT ALL ON kafka_trigger TO windmill_user; +GRANT ALL ON kafka_trigger TO windmill_admin; + +ALTER TABLE kafka_trigger ENABLE ROW LEVEL SECURITY; + +CREATE POLICY admin_policy ON kafka_trigger FOR ALL TO windmill_admin USING (true); + +CREATE POLICY see_folder_extra_perms_user_select ON kafka_trigger FOR SELECT TO windmill_user +USING (SPLIT_PART(kafka_trigger.path, '/', 1) = 'f' AND SPLIT_PART(kafka_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_read'), ',')::text[])); +CREATE POLICY see_folder_extra_perms_user_insert ON kafka_trigger FOR INSERT TO windmill_user +WITH CHECK (SPLIT_PART(kafka_trigger.path, '/', 1) = 'f' AND SPLIT_PART(kafka_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[])); +CREATE POLICY see_folder_extra_perms_user_update ON kafka_trigger FOR UPDATE TO windmill_user +USING (SPLIT_PART(kafka_trigger.path, '/', 1) = 'f' AND SPLIT_PART(kafka_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[])); +CREATE POLICY see_folder_extra_perms_user_delete ON kafka_trigger FOR DELETE TO windmill_user +USING (SPLIT_PART(kafka_trigger.path, '/', 1) = 'f' AND SPLIT_PART(kafka_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.folders_write'), ',')::text[])); + +CREATE POLICY see_own ON kafka_trigger FOR ALL TO windmill_user +USING (SPLIT_PART(kafka_trigger.path, '/', 1) = 'u' AND SPLIT_PART(kafka_trigger.path, '/', 2) = current_setting('session.user')); +CREATE POLICY see_member ON kafka_trigger FOR ALL TO windmill_user +USING (SPLIT_PART(kafka_trigger.path, '/', 1) = 'g' AND SPLIT_PART(kafka_trigger.path, '/', 2) = any(regexp_split_to_array(current_setting('session.groups'), ',')::text[])); + +CREATE POLICY see_extra_perms_user_select ON kafka_trigger FOR SELECT TO windmill_user +USING (extra_perms ? CONCAT('u/', current_setting('session.user'))); +CREATE POLICY see_extra_perms_user_insert ON kafka_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 kafka_trigger FOR UPDATE TO windmill_user +USING ((extra_perms ->> CONCAT('u/', current_setting('session.user')))::boolean); +CREATE POLICY see_extra_perms_user_delete ON kafka_trigger FOR DELETE TO windmill_user +USING ((extra_perms ->> CONCAT('u/', current_setting('session.user')))::boolean); + +CREATE POLICY see_extra_perms_groups_select ON kafka_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 kafka_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 kafka_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 kafka_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 b7c4a41aed..02a2de9265 100644 --- a/backend/windmill-api/Cargo.toml +++ b/backend/windmill-api/Cargo.toml @@ -19,6 +19,7 @@ parquet = ["dep:datafusion", "dep:object_store", "dep:url", "windmill-common/par prometheus = ["windmill-common/prometheus", "windmill-queue/prometheus", "dep:prometheus"] openidconnect = ["dep:openidconnect"] tantivy = ["dep:windmill-indexer"] +kafka = ["dep:rdkafka"] [dependencies] windmill-queue.workspace = true @@ -96,6 +97,7 @@ url = { workspace = true, optional = true} jsonwebtoken = { workspace = true } matchit.workspace = true tokio-tungstenite.workspace = true +rdkafka = { workspace = true, optional = true } pin-project.workspace = true http.workspace = true diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 1685505324..5bcf481f7a 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -2332,9 +2332,12 @@ paths: type: boolean websocket_used: type: boolean + kafka_used: + type: boolean required: - http_routes_used - websocket_used + - kafka_used /w/{workspace}/users/list: get: @@ -7685,6 +7688,169 @@ paths: schema: type: string + /w/{workspace}/kafka_triggers/create: + post: + summary: create kafka trigger + operationId: createKafkaTrigger + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + requestBody: + description: new kafka trigger + required: true + content: + application/json: + schema: + $ref: "#/components/schemas/NewKafkaTrigger" + responses: + "201": + description: kafka trigger created + content: + text/plain: + schema: + type: string + + /w/{workspace}/kafka_triggers/update/{path}: + post: + summary: update kafka trigger + operationId: updateKafkaTrigger + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/Path" + requestBody: + description: updated trigger + required: true + content: + application/json: + schema: + $ref: "#/components/schemas/EditKafkaTrigger" + responses: + "200": + description: kafka trigger updated + content: + text/plain: + schema: + type: string + + /w/{workspace}/kafka_triggers/delete/{path}: + delete: + summary: delete kafka trigger + operationId: deleteKafkaTrigger + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/Path" + responses: + "200": + description: kafka trigger deleted + content: + text/plain: + schema: + type: string + + /w/{workspace}/kafka_triggers/get/{path}: + get: + summary: get kafka trigger + operationId: getKafkaTrigger + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/Path" + responses: + "200": + description: kafka trigger deleted + content: + application/json: + schema: + $ref: "#/components/schemas/KafkaTrigger" + + + /w/{workspace}/kafka_triggers/list: + get: + summary: list kafka triggers + operationId: listKafkaTriggers + tags: + - kafka_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: kafka trigger list + content: + application/json: + schema: + type: array + items: + $ref: "#/components/schemas/KafkaTrigger" + + + /w/{workspace}/kafka_triggers/exists/{path}: + get: + summary: does kafka trigger exists + operationId: existsKafkaTrigger + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/Path" + responses: + "200": + description: kafka trigger exists + content: + application/json: + schema: + type: boolean + + /w/{workspace}/kafka_triggers/setenabled/{path}: + post: + summary: set enabled kafka trigger + operationId: setKafkaTriggerEnabled + tags: + - kafka_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/Path" + requestBody: + description: updated kafka trigger enable + required: true + content: + application/json: + schema: + type: object + properties: + enabled: + type: boolean + required: + - enabled + responses: + "200": + description: kafka trigger enabled set + content: + text/plain: + schema: + type: string + /groups/list: get: @@ -8554,6 +8720,7 @@ paths: raw_app, http_trigger, websocket_trigger, + kafka_trigger, ] responses: "200": @@ -8592,6 +8759,7 @@ paths: raw_app, http_trigger, websocket_trigger, + kafka_trigger, ] requestBody: description: acl to add @@ -8641,6 +8809,7 @@ paths: raw_app, http_trigger, websocket_trigger, + kafka_trigger, ] requestBody: description: acl to add @@ -11932,6 +12101,8 @@ components: type: number websocket_count: type: number + kafka_count: + type: number WebsocketTrigger: type: object @@ -12103,6 +12274,114 @@ components: required: - runnable_result + KafkaTrigger: + type: object + properties: + path: + type: string + edited_by: + type: string + edited_at: + type: string + format: date-time + script_path: + type: string + kafka_resource_path: + type: string + group_id: + type: string + topics: + type: array + items: + 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 + + required: + - path + - edited_by + - edited_at + - script_path + - kafka_resource_path + - group_id + - topics + - extra_perms + - is_flow + - email + - workspace_id + - enabled + + NewKafkaTrigger: + type: object + properties: + path: + type: string + script_path: + type: string + is_flow: + type: boolean + kafka_resource_path: + type: string + group_id: + type: string + topics: + type: array + items: + type: string + enabled: + type: boolean + + required: + - path + - script_path + - is_flow + - kafka_resource_path + - group_id + - topics + + EditKafkaTrigger: + type: object + properties: + kafka_resource_path: + type: string + group_id: + type: string + topics: + type: array + items: + type: string + path: + type: string + script_path: + type: string + is_flow: + type: boolean + + required: + - path + - script_path + - kafka_resource_path + - group_id + - topics + - is_flow + Group: type: object properties: diff --git a/backend/windmill-api/src/apps.rs b/backend/windmill-api/src/apps.rs index 4ec07393dc..1b6ce7fb76 100644 --- a/backend/windmill-api/src/apps.rs +++ b/backend/windmill-api/src/apps.rs @@ -8,11 +8,6 @@ use std::collections::HashMap; * LICENSE-AGPL for a copy of the license. */ -#[cfg(feature = "parquet")] -use crate::{job_helpers_ee::{ - get_random_file_name, get_s3_resource, get_workspace_s3_resource, upload_file_internal, - UploadFileResponse, -}, users::fetch_api_authed_from_permissioned_as}; use crate::{ db::{ApiAuthed, DB}, resources::get_resource_value_interpolated_internal, @@ -22,6 +17,14 @@ use crate::{ webhook_util::{WebhookMessage, WebhookShared}, HTTP_CLIENT, }; +#[cfg(feature = "parquet")] +use crate::{ + job_helpers_ee::{ + get_random_file_name, get_s3_resource, get_workspace_s3_resource, upload_file_internal, + UploadFileResponse, + }, + users::fetch_api_authed_from_permissioned_as, +}; use axum::{ extract::{Extension, Json, Path, Query}, response::IntoResponse, @@ -1288,7 +1291,6 @@ async fn upload_s3_file_from_app() -> Result<()> { )); } - #[cfg(feature = "parquet")] #[derive(Debug, Deserialize, Clone)] struct UploadFileToS3Query { @@ -1353,9 +1355,14 @@ async fn upload_s3_file_from_app( let (username, permissioned_as, email) = get_on_behalf_details_from_policy_and_authed(&policy, &opt_authed).await?; - let on_behalf_authed = - fetch_api_authed_from_permissioned_as(permissioned_as, email, &w_id, &db, username) - .await?; + let on_behalf_authed = fetch_api_authed_from_permissioned_as( + permissioned_as, + email, + &w_id, + &db, + Some(username), + ) + .await?; if let Some(file_key) = query.file_key { // file key is provided => requires workspace, user or list policy and must match the regex diff --git a/backend/windmill-api/src/granular_acls.rs b/backend/windmill-api/src/granular_acls.rs index ad5793816b..3ae01cc293 100644 --- a/backend/windmill-api/src/granular_acls.rs +++ b/backend/windmill-api/src/granular_acls.rs @@ -23,7 +23,7 @@ use windmill_common::{ utils::{not_found_if_none, StripPath}, }; -const KINDS: [&str; 10] = [ +const KINDS: [&str; 12] = [ "script", "group_", "resource", @@ -34,6 +34,8 @@ const KINDS: [&str; 10] = [ "app", "raw_app", "http_trigger", + "websocket_trigger", + "kafka_trigger", ]; pub fn workspaced_service() -> Router { diff --git a/backend/windmill-api/src/http_triggers.rs b/backend/windmill-api/src/http_triggers.rs index 2c323c0e8c..feeb9977a8 100644 --- a/backend/windmill-api/src/http_triggers.rs +++ b/backend/windmill-api/src/http_triggers.rs @@ -507,7 +507,7 @@ async fn get_http_route_trigger( trigger.email.clone(), &trigger.workspace_id, &db, - username_override.unwrap_or("anonymous".to_string()), + Some(username_override.unwrap_or("anonymous".to_string())), ) .await?; diff --git a/backend/windmill-api/src/kafka_triggers_ee.rs b/backend/windmill-api/src/kafka_triggers_ee.rs new file mode 100644 index 0000000000..0b3db94c4a --- /dev/null +++ b/backend/windmill-api/src/kafka_triggers_ee.rs @@ -0,0 +1,14 @@ +use crate::db::DB; +use axum::Router; + +pub fn workspaced_service() -> Router { + Router::new() +} + +pub async fn start_kafka_consumers( + _db: DB, + _rsmq: Option, + mut _killpill_rx: tokio::sync::broadcast::Receiver<()>, +) -> () { + // implementation is not open source +} diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index aa6e6955e6..ba1ab8df59 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -68,6 +68,8 @@ mod ai; mod job_helpers_ee; pub mod job_metrics; pub mod jobs; +#[cfg(all(feature = "enterprise", feature = "kafka"))] +mod kafka_triggers_ee; pub mod oauth2_ee; mod oidc_ee; mod raw_apps; @@ -235,6 +237,9 @@ pub async fn run_server( } } + // #[cfg(feature = "kafka")] + // start_listening().await; + let job_helpers_service = { #[cfg(feature = "parquet")] { @@ -247,9 +252,27 @@ pub async fn run_server( } }; + let kafka_triggers_service = { + #[cfg(all(feature = "enterprise", feature = "kafka"))] + { + kafka_triggers_ee::workspaced_service() + } + + #[cfg(not(all(feature = "enterprise", feature = "kafka")))] + { + Router::new() + } + }; + if !*CLOUD_HOSTED { let ws_killpill_rx = rx.resubscribe(); - websocket_triggers::start_websockets(db.clone(), rsmq, ws_killpill_rx).await; + websocket_triggers::start_websockets(db.clone(), rsmq.clone(), ws_killpill_rx).await; + + #[cfg(all(feature = "enterprise", feature = "kafka"))] + { + let kafka_killpill_rx = rx.resubscribe(); + kafka_triggers_ee::start_kafka_consumers(db.clone(), rsmq, kafka_killpill_rx).await; + } } // build our application with a route @@ -296,7 +319,8 @@ pub async fn run_server( .nest( "/websocket_triggers", websocket_triggers::workspaced_service(), - ), + ) + .nest("/kafka_triggers", kafka_triggers_service), ) .nest("/workspaces", workspaces::global_service()) .nest( @@ -321,10 +345,7 @@ pub async fn run_server( "/srch/w/:workspace_id/index", indexer_ee::workspaced_service(), ) - .nest( - "/srch/index", - indexer_ee::global_service(), - ) + .nest("/srch/index", indexer_ee::global_service()) .nest("/oidc", oidc_ee::global_service()) .nest( "/saml", diff --git a/backend/windmill-api/src/triggers.rs b/backend/windmill-api/src/triggers.rs index 337a20364b..bc734fc6ea 100644 --- a/backend/windmill-api/src/triggers.rs +++ b/backend/windmill-api/src/triggers.rs @@ -18,6 +18,7 @@ pub struct TriggersCount { webhook_count: i64, email_count: i64, websocket_count: i64, + kafka_count: i64, } pub(crate) async fn get_triggers_count_internal( db: &DB, @@ -64,6 +65,16 @@ pub(crate) async fn get_triggers_count_internal( .await? .unwrap_or(0); + let kafka_count = sqlx::query_scalar!( + "SELECT COUNT(*) FROM kafka_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[]", @@ -105,6 +116,7 @@ pub(crate) async fn get_triggers_count_internal( webhook_count, email_count, websocket_count, + kafka_count, })) } diff --git a/backend/windmill-api/src/users.rs b/backend/windmill-api/src/users.rs index bf95a632b1..6d0f0b4ad1 100644 --- a/backend/windmill-api/src/users.rs +++ b/backend/windmill-api/src/users.rs @@ -746,7 +746,7 @@ pub async fn fetch_api_authed( email: String, w_id: &str, db: &DB, - username_override: String, + username_override: Option, ) -> error::Result { let permissioned_as = username_to_permissioned_as(username.as_str()); fetch_api_authed_from_permissioned_as(permissioned_as, email, w_id, db, username_override).await @@ -757,7 +757,7 @@ pub async fn fetch_api_authed_from_permissioned_as( email: String, w_id: &str, db: &DB, - username_override: String, + username_override: Option, ) -> error::Result { let authed = fetch_authed_from_permissioned_as(permissioned_as, email.clone(), w_id, db).await?; @@ -769,7 +769,7 @@ pub async fn fetch_api_authed_from_permissioned_as( groups: authed.groups, folders: authed.folders, scopes: authed.scopes, - username_override: Some(username_override), + username_override: username_override, }) } diff --git a/backend/windmill-api/src/websocket_triggers.rs b/backend/windmill-api/src/websocket_triggers.rs index ecdc990061..c9b16b4958 100644 --- a/backend/windmill-api/src/websocket_triggers.rs +++ b/backend/windmill-api/src/websocket_triggers.rs @@ -527,7 +527,7 @@ async fn wait_runnable_result( ws_trigger.email.clone(), &ws_trigger.workspace_id, &db, - username_override, + Some(username_override), ) .await?; @@ -773,7 +773,9 @@ async fn listen_to_websocket( rsmq: Option, mut killpill_rx: tokio::sync::broadcast::Receiver<()>, ) -> () { - update_ping(&db, &ws_trigger, Some("Connecting...")).await; + if let None = update_ping(&db, &ws_trigger, Some("Connecting...")).await { + return; + } let url = ws_trigger.url.as_str(); @@ -970,7 +972,7 @@ async fn run_job( trigger.email.clone(), &trigger.workspace_id, db, - "anonymous".to_string(), + Some("anonymous".to_string()), ) .await?; diff --git a/backend/windmill-api/src/workspaces.rs b/backend/windmill-api/src/workspaces.rs index 36bef54655..0fc98863d7 100644 --- a/backend/windmill-api/src/workspaces.rs +++ b/backend/windmill-api/src/workspaces.rs @@ -1296,6 +1296,7 @@ async fn set_encryption_key( struct UsedTriggers { pub websocket_used: bool, pub http_routes_used: bool, + pub kafka_used: bool, } async fn get_used_triggers( @@ -1306,7 +1307,10 @@ async fn get_used_triggers( 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!""#, + 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!", + EXISTS(SELECT 1 FROM kafka_trigger WHERE workspace_id = $1) as "kafka_used!""#, w_id, ) .fetch_one(&mut *tx) diff --git a/frontend/src/lib/components/Path.svelte b/frontend/src/lib/components/Path.svelte index 2d112387af..0ff3ff56bc 100644 --- a/frontend/src/lib/components/Path.svelte +++ b/frontend/src/lib/components/Path.svelte @@ -14,7 +14,8 @@ ScriptService, HttpTriggerService, VariableService, - WebsocketTriggerService + WebsocketTriggerService, + KafkaTriggerService } from '$lib/gen' import { superadmin, userStore, workspaceStore } from '$lib/stores' import { createEventDispatcher, getContext } from 'svelte' @@ -37,6 +38,7 @@ | 'raw_app' | 'http_trigger' | 'websocket_trigger' + | 'kafka_trigger' let meta: Meta | undefined = undefined export let fullNamePlaceholder: string | undefined = undefined export let namePlaceholder = '' @@ -225,6 +227,11 @@ workspace: $workspaceStore!, path: path }) + } else if (kind == 'kafka_trigger') { + return await KafkaTriggerService.existsKafkaTrigger({ + workspace: $workspaceStore!, + path: path + }) } else { return false } diff --git a/frontend/src/lib/components/ShareModal.svelte b/frontend/src/lib/components/ShareModal.svelte index 8f0192d5e4..64712812f3 100644 --- a/frontend/src/lib/components/ShareModal.svelte +++ b/frontend/src/lib/components/ShareModal.svelte @@ -26,6 +26,7 @@ | 'raw_app' | 'http_trigger' | 'websocket_trigger' + | 'kafka_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 7d998aa21a..f9e4f927b9 100644 --- a/frontend/src/lib/components/details/DetailPageDetailPanel.svelte +++ b/frontend/src/lib/components/details/DetailPageDetailPanel.svelte @@ -51,6 +51,7 @@ + diff --git a/frontend/src/lib/components/details/DetailPageLayout.svelte b/frontend/src/lib/components/details/DetailPageLayout.svelte index 2fb070cbce..0889963609 100644 --- a/frontend/src/lib/components/details/DetailPageLayout.svelte +++ b/frontend/src/lib/components/details/DetailPageLayout.svelte @@ -52,6 +52,7 @@ + @@ -94,6 +95,7 @@ + diff --git a/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte b/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte index 144cd1d865..e0a3d0b7f0 100644 --- a/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte +++ b/frontend/src/lib/components/details/DetailPageTriggerPanel.svelte @@ -3,6 +3,7 @@ import { CalendarCheck2, MailIcon, Route, Terminal, Webhook, Unplug } from 'lucide-svelte' import HighlightTheme from '../HighlightTheme.svelte' + import KafkaIcon from '../icons/KafkaIcon.svelte' export let triggerSelected: | 'webhooks' @@ -11,6 +12,7 @@ | 'cli' | 'routes' | 'websockets' + | 'kafka' | 'scheduledPoll' = 'webhooks' export let simplfiedPoll: boolean = false @@ -43,6 +45,12 @@ Websockets + + + + Kafka + + @@ -69,6 +77,8 @@ {:else if triggerSelected === 'websockets'} + {:else if triggerSelected === 'kafka'} + {: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 ace322ad65..fb3254d80a 100644 --- a/frontend/src/lib/components/graph/renderers/triggers/TriggersBadge.svelte +++ b/frontend/src/lib/components/graph/renderers/triggers/TriggersBadge.svelte @@ -10,6 +10,7 @@ import { type TriggerContext } from '$lib/components/triggers' import { FlowService, ScriptService } from '$lib/gen' import { workspaceStore } from '$lib/stores' + import KafkaIcon from '$lib/components/icons/KafkaIcon.svelte' const { selectedTrigger, triggersCount } = getContext('TriggerContext') @@ -18,8 +19,14 @@ export let isFlow: boolean export let selected: boolean export let showOnlyWithCount: boolean - export let triggersToDisplay: ('webhooks' | 'schedules' | 'routes' | 'websockets' | 'emails')[] = - ['webhooks', 'schedules', 'routes', 'websockets', 'emails'] + export let triggersToDisplay: ( + | 'webhooks' + | 'schedules' + | 'routes' + | 'websockets' + | 'kafka' + | 'emails' + )[] = ['webhooks', 'schedules', 'routes', 'websockets', 'kafka', 'emails'] const dispatch = createEventDispatcher() onMount(() => { @@ -47,6 +54,7 @@ schedules: { icon: Calendar, countKey: 'schedule_count' }, routes: { icon: Route, countKey: 'http_routes_count' }, websockets: { icon: Unplug, countKey: 'websocket_count' }, + kafka: { icon: KafkaIcon, countKey: 'kafka_count' }, emails: { icon: Mail, countKey: 'email_count' } } diff --git a/frontend/src/lib/components/icons/KafkaIcon.svelte b/frontend/src/lib/components/icons/KafkaIcon.svelte new file mode 100644 index 0000000000..9b4eda8ba5 --- /dev/null +++ b/frontend/src/lib/components/icons/KafkaIcon.svelte @@ -0,0 +1,20 @@ + + + + + diff --git a/frontend/src/lib/components/sidebar/SidebarContent.svelte b/frontend/src/lib/components/sidebar/SidebarContent.svelte index 4e34071256..e5ab2635ea 100644 --- a/frontend/src/lib/components/sidebar/SidebarContent.svelte +++ b/frontend/src/lib/components/sidebar/SidebarContent.svelte @@ -49,6 +49,7 @@ import { type Changelog, changelogs } from './changelogs' import { page } from '$app/stores' import SideBarNotification from './SideBarNotification.svelte' + import KafkaIcon from '../icons/KafkaIcon.svelte' export let numUnacknowledgedCriticalAlerts = 0 @@ -97,6 +98,13 @@ icon: Unplug, disabled: $userStore?.operator, kind: 'ws' + }, + { + label: 'Kafka', + href: '/kafka_triggers', + icon: KafkaIcon, + disabled: $userStore?.operator, + kind: 'kafka' } ] @@ -205,10 +213,10 @@ action: () => { isCriticalAlertsUIOpen.set(true) }, - icon: AlertCircle, - notificationCount: numUnacknowledgedCriticalAlerts - } - ] + icon: AlertCircle, + notificationCount: numUnacknowledgedCriticalAlerts + } + ] : []) ] } diff --git a/frontend/src/lib/components/triggers.ts b/frontend/src/lib/components/triggers.ts index 4c06603f20..efc8c8a901 100644 --- a/frontend/src/lib/components/triggers.ts +++ b/frontend/src/lib/components/triggers.ts @@ -11,7 +11,14 @@ export type ScheduleTrigger = { export type TriggerContext = { selectedTrigger: Writable< - 'webhooks' | 'emails' | 'schedules' | 'cli' | 'routes' | 'websockets' | 'scheduledPoll' + | 'webhooks' + | 'emails' + | 'schedules' + | 'cli' + | 'routes' + | 'websockets' + | 'scheduledPoll' + | 'kafka' > primarySchedule: Writable triggersCount: Writable diff --git a/frontend/src/lib/components/triggers/KafkaTriggerEditor.svelte b/frontend/src/lib/components/triggers/KafkaTriggerEditor.svelte new file mode 100644 index 0000000000..33083d7ec6 --- /dev/null +++ b/frontend/src/lib/components/triggers/KafkaTriggerEditor.svelte @@ -0,0 +1,23 @@ + + +{#if open} + +{/if} diff --git a/frontend/src/lib/components/triggers/KafkaTriggerEditorInner.svelte b/frontend/src/lib/components/triggers/KafkaTriggerEditorInner.svelte new file mode 100644 index 0000000000..c65bdbc32c --- /dev/null +++ b/frontend/src/lib/components/triggers/KafkaTriggerEditorInner.svelte @@ -0,0 +1,310 @@ + + + + + + {#if !drawerLoading && can_write} + {#if edit} +
+ { + await KafkaTriggerService.setKafkaTriggerEnabled({ + path: initialPath, + workspace: $workspaceStore ?? '', + requestBody: { enabled: e.detail } + }) + sendUserToast(`${e.detail ? 'enabled' : 'disabled'} kafka trigger ${initialPath}`) + }} + /> +
+ {/if} + + {/if} +
+ {#if drawerLoading} + + {:else} + + {#if edit} + Changes can take up to 30 seconds to take effect. + {:else} + Kafka consumers can take up to 30 seconds to start. + {/if} + +
+
+ +
+ +
+
+
+
+ Resource + +
+ +
+ + + +
+
+ +
+

+ Pick a script or flow to be triggered +

+
+ +
+
+
+ {/if} +
+
diff --git a/frontend/src/lib/components/triggers/KafkaTriggersPanel.svelte b/frontend/src/lib/components/triggers/KafkaTriggersPanel.svelte new file mode 100644 index 0000000000..6b38beab73 --- /dev/null +++ b/frontend/src/lib/components/triggers/KafkaTriggersPanel.svelte @@ -0,0 +1,108 @@ + + + { + loadTriggers() + }} + bind:this={kafkaTriggerEditor} +/> + +
+ {#if !newItem} + {#if isCloudHosted()} + + Kafka triggers are disabled in the multi-tenant cloud. + + {:else if $userStore?.is_admin || $userStore?.is_super_admin} + + {:else} + + {/if} + {/if} + + {#if kafkaTriggers} + {#if kafkaTriggers.length == 0} +
No kafka triggers
+ {:else} +
+ {#each kafkaTriggers as kafkaTrigger (kafkaTrigger.path)} +
+
{kafkaTrigger.path}
+
+ {kafkaTrigger.kafka_resource_path} +
+
+ +
+
+ {/each} +
+ {/if} + {:else} + + {/if} + + {#if newItem} + + Deploy the {isFlow ? 'flow' : 'script'} to add kafka triggers. + + {/if} +
diff --git a/frontend/src/lib/components/triggers/TriggersEditor.svelte b/frontend/src/lib/components/triggers/TriggersEditor.svelte index 5e0cebef7a..787378657b 100644 --- a/frontend/src/lib/components/triggers/TriggersEditor.svelte +++ b/frontend/src/lib/components/triggers/TriggersEditor.svelte @@ -12,6 +12,7 @@ import type { TriggerContext } from '$lib/components/triggers' import WebsocketTriggersPanel from './WebsocketTriggersPanel.svelte' import ScheduledPollPanel from './ScheduledPollPanel.svelte' + import KafkaTriggersPanel from './KafkaTriggersPanel.svelte' export let noEditor: boolean export let newItem = false @@ -31,6 +32,7 @@ Schedules Routes Websockets + Kafka Email {#if isFlow} {/if} + {#if $selectedTrigger === 'kafka'} +
+ +
+ {/if} + {#if $selectedTrigger === 'schedules'}
- query: Record + params: Record // path parameters + query: Record // query parameters headers: Record }, websocket?: { url: string // The websocket url + }, + kafka?: { + brokers: string[] + topic: string + group_id: string } }, /* your other args */ @@ -538,17 +543,22 @@ export async function preprocessor( const DENO_PREPROCESSOR_MODULE_CODE = ` export async function preprocessor( wm_trigger: { - kind: 'http' | 'email' | 'wehbook' | 'websocket', + kind: 'http' | 'email' | 'webhook' | 'websocket' | 'kafka', http?: { route: string // The route path, e.g. "/users/:id" path: string // The actual path called, e.g. "/users/123" method: string - params: Record - query: Record + params: Record // path parameters + query: Record // query parameters headers: Record }, websocket?: { url: string // The websocket url + }, + kafka?: { + brokers: string[] + topic: string + group_id: string } }, /* your other args */ @@ -599,10 +609,16 @@ class Http(TypedDict): class Websocket(TypedDict): url: str # The websocket url +class Kafka(TypedDict): + topic: str + brokers: list[str] + group_id: str + class WmTrigger(TypedDict): - kind: Literal["http", "email", "webhook", "websocket"] - http: Http | None - websocket: Websocket | None + kind: Literal["http", "email", "webhook", "websocket", "kafka"] + http: Http | None + websocket: Websocket | None + kafka: Kafka | None def preprocessor( wm_trigger: WmTrigger, diff --git a/frontend/src/routes/(root)/(logged)/+layout.svelte b/frontend/src/routes/(root)/(logged)/+layout.svelte index 66b8e32fd8..fa60e2065d 100644 --- a/frontend/src/routes/(root)/(logged)/+layout.svelte +++ b/frontend/src/routes/(root)/(logged)/+layout.svelte @@ -190,15 +190,20 @@ async function loadUsedTriggerKinds() { let usedKinds: string[] = [] - const { http_routes_used, websocket_used } = await WorkspaceService.getUsedTriggers({ - workspace: $workspaceStore ?? '' - }) + const { http_routes_used, websocket_used, kafka_used } = await WorkspaceService.getUsedTriggers( + { + workspace: $workspaceStore ?? '' + } + ) if (http_routes_used) { usedKinds.push('http') } if (websocket_used) { usedKinds.push('ws') } + if (kafka_used) { + usedKinds.push('kafka') + } $usedTriggerKinds = usedKinds } 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 f91f68b7be..b69de584a7 100644 --- a/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/flows/get/[...path]/+page.svelte @@ -61,6 +61,7 @@ import { writable } from 'svelte/store' import TriggersBadge from '$lib/components/graph/renderers/triggers/TriggersBadge.svelte' import WebsocketTriggersPanel from '$lib/components/triggers/WebsocketTriggersPanel.svelte' + import KafkaTriggersPanel from '$lib/components/triggers/KafkaTriggersPanel.svelte' let flow: Flow | undefined let can_write = false @@ -539,6 +540,13 @@
+ + +
+ +
+
+
+ import { KafkaTriggerService, type KafkaTrigger } 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 { Alert, 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 { 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 KafkaTriggerEditor from '$lib/components/triggers/KafkaTriggerEditor.svelte' + import Popover from '$lib/components/Popover.svelte' + import { isCloudHosted } from '$lib/cloud' + import KafkaIcon from '$lib/components/icons/KafkaIcon.svelte' + + type TriggerW = KafkaTrigger & { canWrite: boolean } + + let triggers: TriggerW[] = [] + let shareModal: ShareModal + let loading = true + + async function loadTriggers(): Promise { + triggers = (await KafkaTriggerService.listKafkaTriggers({ workspace: $workspaceStore! })).map( + (x) => { + return { canWrite: canWrite(x.path, x.extra_perms!, $userStore), ...x } + } + ) + loading = false + } + + let interval = setInterval(async () => { + try { + const newTriggers = await KafkaTriggerService.listKafkaTriggers({ + workspace: $workspaceStore! + }) + for (let i = 0; i < triggers.length; i++) { + const newTrigger = newTriggers.find((x) => x.path === triggers[i].path) + if (newTrigger) { + triggers[i] = { + ...triggers[i], + error: newTrigger.error, + last_server_ping: newTrigger.last_server_ping, + enabled: newTrigger.enabled + } + } + } + } catch (err) { + console.error(err) + } + }, 5000) + + onDestroy(() => { + clearInterval(interval) + }) + + async function setTriggerEnabled(path: string, enabled: boolean): Promise { + try { + await KafkaTriggerService.setKafkaTriggerEnabled({ + path, + workspace: $workspaceStore!, + requestBody: { enabled } + }) + } catch (err) { + sendUserToast( + `Cannot ` + (enabled ? 'enable' : 'disable') + ` kafka trigger: ${err.body}`, + true + ) + } finally { + loadTriggers() + } + } + + $: { + if ($workspaceStore && $userStore) { + loadTriggers() + } + } + let kafkaTriggerEditor: KafkaTriggerEditor + + 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 + ')'} +/> + + + + + + + {#if isCloudHosted()} + + Kafka triggers are disabled in the multi-tenant cloud. + +
+ {/if} +
+
+ +
+
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 kafka triggers
+ {:else if items?.length} +
+ {#each items.slice(0, nbDisplayed) as { path, edited_by, edited_at, script_path, is_flow, kafka_resource_path, topics, extra_perms, canWrite, marked, error, last_server_ping, enabled } (path)} + {@const href = `${is_flow ? '/flows/get' : '/scripts/get'}/${script_path}`} + {@const ping = last_server_ping ? new Date(last_server_ping) : undefined} + +
+
+ + + kafkaTriggerEditor?.openEdit(path, is_flow)} + class="min-w-0 grow hover:underline decoration-gray-400" + > +
+ {#if marked} + + {@html marked} + + {:else} + {kafka_resource_path} - {topics.join(', ')} + {/if} +
+
+ {path} +
+
+ runnable: {script_path} +
+
+ + + +
+ {#if (enabled && (!ping || ping.getTime() < new Date().getTime() - 15 * 1000 || error)) || (!enabled && error)} + + + + + +
+ {#if enabled} + Consumer is not connected{error ? ': ' + error : ''} + {:else} + Consumer was disabled because of an error: {error} + {/if} +
+
+ {:else if enabled} + + + + +
Consumer is connected
+
+ {/if} +
+ + { + setTriggerEnabled(path, e.detail) + }} + /> + +
+ + { + goto(href) + } + }, + { + displayName: 'Delete', + type: 'delete', + icon: Trash, + disabled: !canWrite, + action: async () => { + await KafkaTriggerService.deleteKafkaTrigger({ + workspace: $workspaceStore ?? '', + path + }) + loadTriggers() + } + }, + { + displayName: canWrite ? 'Edit' : 'View', + icon: canWrite ? Pen : Eye, + action: () => { + kafkaTriggerEditor?.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, 'kafka_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/(root)/(logged)/scripts/get/[...hash]/+page.svelte b/frontend/src/routes/(root)/(logged)/scripts/get/[...hash]/+page.svelte index b90ed9d2b3..7ac04486ed 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 { writable } from 'svelte/store' import TriggersBadge from '$lib/components/graph/renderers/triggers/TriggersBadge.svelte' import WebsocketTriggersPanel from '$lib/components/triggers/WebsocketTriggersPanel.svelte' + import KafkaTriggersPanel from '$lib/components/triggers/KafkaTriggersPanel.svelte' let script: Script | undefined let topHash: string | undefined @@ -721,6 +722,11 @@
+ +
+ +
+