mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-18 00:02:27 +00:00
Compare commits
26
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8d3b18f523 | ||
|
|
8e19f9652d | ||
|
|
45980f0220 | ||
|
|
8363ff1eee | ||
|
|
cf596f370a | ||
|
|
b12304d834 | ||
|
|
b5f8d931bd | ||
|
|
f927c673b5 | ||
|
|
4f653ca957 | ||
|
|
08031640a0 | ||
|
|
244c09fcd7 | ||
|
|
1456f73c16 | ||
|
|
4343b73485 | ||
|
|
b4b3edb64d | ||
|
|
4631e58c66 | ||
|
|
eb827ce5e8 | ||
|
|
dd421845ba | ||
|
|
b1d6ac91bd | ||
|
|
894d8a94f8 | ||
|
|
7d37a83d4f | ||
|
|
750926d8da | ||
|
|
0b0696459e | ||
|
|
a3fc27b232 | ||
|
|
ff70a4e9d1 | ||
|
|
76377a00a6 | ||
|
|
dc5e69481d |
@@ -39,7 +39,7 @@ jobs:
|
||||
- name: install xmlsec1 and gssapi
|
||||
run: |
|
||||
sudo apt-get update
|
||||
sudo apt-get install -y libxml2-dev libxmlsec1-dev libkrb5-dev libsasl2-dev mold clang
|
||||
sudo apt-get install -y libxml2-dev libxmlsec1-dev libkrb5-dev libsasl2-dev libcurl4-openssl-dev mold clang
|
||||
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
@@ -109,7 +109,7 @@ jobs:
|
||||
- name: install xmlsec1 and gssapi
|
||||
run: |
|
||||
sudo apt-get update
|
||||
sudo apt-get install -y libxml2-dev libxmlsec1-dev libkrb5-dev libsasl2-dev mold clang
|
||||
sudo apt-get install -y libxml2-dev libxmlsec1-dev libkrb5-dev libsasl2-dev libcurl4-openssl-dev mold clang
|
||||
|
||||
- name: Substitute EE code (EE logic is behind feature flag)
|
||||
run: |
|
||||
|
||||
@@ -72,7 +72,7 @@ jobs:
|
||||
bundler-cache: false
|
||||
- name: Install PowerShell, mold and clang
|
||||
run: |
|
||||
sudo apt-get update && sudo apt-get install -y powershell mold clang
|
||||
sudo apt-get update && sudo apt-get install -y powershell mold clang libcurl4-openssl-dev
|
||||
working-directory: /
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
@@ -222,4 +222,4 @@ jobs:
|
||||
run: |
|
||||
deno --version && bun -v && node --version && go version && python3 --version && php --version && ruby --version && pwsh --version && dotnet --version
|
||||
cd windmill-duckdb-ffi-internal && ./build_dev.sh && cd ..
|
||||
DENO_PATH=$(which deno) BUN_PATH=$(which bun) NODE_BIN_PATH=$(which node) GO_PATH=$(which go) UV_PATH=$(which uv) PHP_PATH=$(which php) COMPOSER_PATH=$(which composer) RUBY_PATH=$(which ruby) RUBY_BUNDLE_PATH=$(which bundle) RUBY_GEM_PATH=$(which gem) POWERSHELL_PATH=$(which pwsh) DOTNET_PATH=$(which dotnet) cargo test --features enterprise,deno_core,duckdb,license,python,rust,scoped_cache,parquet,private,private_registry_test,csharp,php,ruby,mysql,quickjs,mcp --all -- --nocapture --test-threads=16
|
||||
DENO_PATH=$(which deno) BUN_PATH=$(which bun) NODE_BIN_PATH=$(which node) GO_PATH=$(which go) UV_PATH=$(which uv) PHP_PATH=$(which php) COMPOSER_PATH=$(which composer) RUBY_PATH=$(which ruby) RUBY_BUNDLE_PATH=$(which bundle) RUBY_GEM_PATH=$(which gem) POWERSHELL_PATH=$(which pwsh) DOTNET_PATH=$(which dotnet) cargo test --features enterprise,deno_core,duckdb,license,python,rust,scoped_cache,parquet,private,private_registry_test,csharp,php,ruby,mysql,quickjs,mcp --all -- --nocapture --test-threads=10
|
||||
|
||||
@@ -11,7 +11,7 @@ env:
|
||||
|
||||
jobs:
|
||||
cargo_build_windows:
|
||||
runs-on: windows-latest
|
||||
runs-on: blacksmith-16vcpu-windows-2025
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
|
||||
@@ -30,33 +30,40 @@ jobs:
|
||||
token: ${{ secrets.WINDMILL_EE_PRIVATE_ACCESS }}
|
||||
fetch-depth: 0
|
||||
|
||||
- name: Setup Rust
|
||||
uses: actions-rs/toolchain@v1
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
cache-workspaces: backend
|
||||
toolchain: 1.90.0
|
||||
override: true
|
||||
|
||||
- name: Substitute EE code
|
||||
shell: bash
|
||||
run: |
|
||||
./backend/substitute_ee_code.sh --copy --dir ./windmill-ee-private
|
||||
|
||||
- name: Cargo check (fail fast on warnings)
|
||||
timeout-minutes: 60
|
||||
env:
|
||||
RUSTFLAGS: "-D warnings"
|
||||
run: |
|
||||
mkdir frontend/build && cd backend
|
||||
New-Item -Path . -Name "windmill-api/openapi-deref.yaml" -ItemType "File" -Force
|
||||
cargo check --features=ee_windows
|
||||
|
||||
- name: Cargo build dynamic libraries windows
|
||||
timeout-minutes: 90
|
||||
timeout-minutes: 180
|
||||
run: |
|
||||
cd backend/windmill-duckdb-ffi-internal
|
||||
cargo build --release -p windmill_duckdb_ffi_internal
|
||||
|
||||
- name: Cargo build binary windows
|
||||
timeout-minutes: 90
|
||||
timeout-minutes: 180
|
||||
run: |
|
||||
vcpkg.exe install openssl-windows:x64-windows
|
||||
vcpkg.exe install openssl:x64-windows-static
|
||||
vcpkg.exe integrate install
|
||||
$env:VCPKGRS_DYNAMIC=1
|
||||
$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
|
||||
cd backend
|
||||
cargo build --release --features=ee_windows
|
||||
- name: Rename binary with corresponding architecture
|
||||
run: |
|
||||
|
||||
@@ -78,7 +78,7 @@ jobs:
|
||||
- name: install xmlsec1 and gssapi
|
||||
run: |
|
||||
sudo apt-get update
|
||||
sudo apt-get install -y libxml2-dev libxmlsec1-dev libkrb5-dev libsasl2-dev
|
||||
sudo apt-get install -y libxml2-dev libxmlsec1-dev libkrb5-dev libsasl2-dev libcurl4-openssl-dev
|
||||
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
|
||||
@@ -13,7 +13,7 @@ env:
|
||||
|
||||
jobs:
|
||||
cargo_build_windows:
|
||||
runs-on: windows-latest
|
||||
runs-on: blacksmith-16vcpu-windows-2025
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
|
||||
@@ -32,11 +32,10 @@ jobs:
|
||||
token: ${{ secrets.WINDMILL_EE_PRIVATE_ACCESS }}
|
||||
fetch-depth: 0
|
||||
|
||||
- name: Setup Rust
|
||||
uses: actions-rs/toolchain@v1
|
||||
- uses: actions-rust-lang/setup-rust-toolchain@v1
|
||||
with:
|
||||
cache-workspaces: backend
|
||||
toolchain: 1.90.0
|
||||
override: true
|
||||
|
||||
- name: Substitute EE code
|
||||
shell: bash
|
||||
@@ -44,13 +43,13 @@ jobs:
|
||||
./backend/substitute_ee_code.sh --copy --dir ./windmill-ee-private
|
||||
|
||||
- name: Cargo build dynamic libraries windows
|
||||
timeout-minutes: 90
|
||||
timeout-minutes: 180
|
||||
run: |
|
||||
cd backend/windmill-duckdb-ffi-internal
|
||||
cargo build --release -p windmill_duckdb_ffi_internal
|
||||
|
||||
- name: Cargo build windows
|
||||
timeout-minutes: 90
|
||||
timeout-minutes: 180
|
||||
run: |
|
||||
vcpkg.exe install openssl-windows:x64-windows
|
||||
vcpkg.exe install openssl:x64-windows-static
|
||||
|
||||
@@ -1,5 +1,41 @@
|
||||
# Changelog
|
||||
|
||||
## [1.629.1](https://github.com/windmill-labs/windmill/compare/v1.629.0...v1.629.1) (2026-02-10)
|
||||
|
||||
|
||||
### Bug Fixes
|
||||
|
||||
* remove unecessary drop index on labeled_jobs_on_jobs ([0803164](https://github.com/windmill-labs/windmill/commit/08031640a02ebd5971793942e8534d69f4f71d28))
|
||||
|
||||
## [1.629.0](https://github.com/windmill-labs/windmill/compare/v1.628.3...v1.629.0) (2026-02-09)
|
||||
|
||||
|
||||
### Features
|
||||
|
||||
* customer portal extra workspace stats ([#7841](https://github.com/windmill-labs/windmill/issues/7841)) ([153dd32](https://github.com/windmill-labs/windmill/commit/153dd32187a3e32e3f26ab88b62195a0f9a359b9))
|
||||
|
||||
|
||||
### Bug Fixes
|
||||
|
||||
* adapt mysql and ruby test assertions ([477832d](https://github.com/windmill-labs/windmill/commit/477832dbeeafb88fd16c174d5d1df8ed042e6f31))
|
||||
* add missing :name param to groups/is_owner route ([fa53a87](https://github.com/windmill-labs/windmill/commit/fa53a87107158c9e8a7e5b522242b0076a78ce46))
|
||||
* **backend:** prevent sqs hanging ([#7857](https://github.com/windmill-labs/windmill/issues/7857)) ([a3fc27b](https://github.com/windmill-labs/windmill/commit/a3fc27b23224aef2949c19f7c123d28e6cfaf968))
|
||||
* box push() future to prevent stack overflow in nested async chains ([67c8aef](https://github.com/windmill-labs/windmill/commit/67c8aef9faea5afe8ce330b477f16b4aed1779a5))
|
||||
* **frontend:** reorganize workspace settings ([#7788](https://github.com/windmill-labs/windmill/issues/7788)) ([dd42184](https://github.com/windmill-labs/windmill/commit/dd421845ba148bba70bcbafbc6a39f3012eb037c))
|
||||
* improve scheduling reliability in extreme pool contention conditions ([#7825](https://github.com/windmill-labs/windmill/issues/7825)) ([bbb397b](https://github.com/windmill-labs/windmill/commit/bbb397b6ad954052f0bd33cc4ff8897eed66e4db))
|
||||
* improve tracing behavior with NO_PROXY ([4cce13f](https://github.com/windmill-labs/windmill/commit/4cce13f5228a05da1bbce43bed7e856ce0bcf979))
|
||||
* incorrect raw app public workspaceStore derived ([edb0d4a](https://github.com/windmill-labs/windmill/commit/edb0d4a05da567b3b0be5d94c9b2856d68ecb0ff))
|
||||
* increase test thread stack size to 8MB in CI ([5548098](https://github.com/windmill-labs/windmill/commit/5548098e083af76a0b7d6f645a5458592d9c8ddc))
|
||||
* install mold+clang in Docker for cargo linker config ([99bc383](https://github.com/windmill-labs/windmill/commit/99bc383f9e94a415ff1dcef1c45ccc4c8dab1a9e))
|
||||
* make V8 runtime init idempotent and auto-initialize before isolate creation ([aa9f3da](https://github.com/windmill-labs/windmill/commit/aa9f3da429da92a059aaabb28481d33b8dacd37b))
|
||||
* parse Python datetime.datetime and datetime.date type annotations ([#7856](https://github.com/windmill-labs/windmill/issues/7856)) ([ff70a4e](https://github.com/windmill-labs/windmill/commit/ff70a4e9d105cac58c0fb0aba8fbec9875533aa4))
|
||||
* prevent V8 SIGSEGV by serializing isolate creation and fixing use-after-free ([05106d7](https://github.com/windmill-labs/windmill/commit/05106d7deeda92b7ae0e1708554f6dcb088c4a08))
|
||||
* reduce DB pool contention by eliminating dual-connection patterns ([#7861](https://github.com/windmill-labs/windmill/issues/7861)) ([4343b73](https://github.com/windmill-labs/windmill/commit/4343b73485843c3b482c21e60052f171ada9b843))
|
||||
* remove mold linker config that breaks Docker builds ([fea0954](https://github.com/windmill-labs/windmill/commit/fea0954f20f9f7c5a43b25b23df530faeac94999))
|
||||
* restart after empty branchone + improve UI ([#7838](https://github.com/windmill-labs/windmill/issues/7838)) ([b1d6ac9](https://github.com/windmill-labs/windmill/commit/b1d6ac91bd3af073feac0b31d97f7b4414d8786e))
|
||||
* use unprotected V8 platform to prevent SIGSEGV on x86_64 Linux ([90d0103](https://github.com/windmill-labs/windmill/commit/90d010347c65086b17f9802dd9a7d2da90dc68eb))
|
||||
* wmill workspace list to list local profiles ([#7843](https://github.com/windmill-labs/windmill/issues/7843)) ([f924a82](https://github.com/windmill-labs/windmill/commit/f924a8268461c49a0fec26e3216ec9546601b8de))
|
||||
|
||||
## [1.628.3](https://github.com/windmill-labs/windmill/compare/v1.628.2...v1.628.3) (2026-02-06)
|
||||
|
||||
|
||||
|
||||
+1
-1
@@ -97,7 +97,7 @@ ARG features=""
|
||||
|
||||
COPY --from=planner /windmill/recipe.json recipe.json
|
||||
|
||||
RUN apt-get update && apt-get install -y libxml2-dev=2.9.* libxmlsec1-dev=1.2.* libkrb5-dev libsasl2-dev clang=1:14.0-55.* libclang-dev=1:14.0-55.* cmake=3.25.* && \
|
||||
RUN apt-get update && apt-get install -y libxml2-dev=2.9.* libxmlsec1-dev=1.2.* libkrb5-dev libsasl2-dev libcurl4-openssl-dev clang=1:14.0-55.* libclang-dev=1:14.0-55.* cmake=3.25.* && \
|
||||
apt-get clean && \
|
||||
rm -rf /var/lib/apt/lists/*
|
||||
|
||||
|
||||
+21
@@ -0,0 +1,21 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n INSERT INTO email_trigger (\n path, local_part, workspaced_local_part, script_path,\n is_flow, workspace_id, edited_by, email\n )\n VALUES ($1, $2, $3, $4, $5, $6, $7, $8)\n ",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Bool",
|
||||
"Varchar",
|
||||
"Bool",
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Varchar"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "1074c6c98e6a0c83ac04172a39abea21c793f58947051d39931d4da0868a1d77"
|
||||
}
|
||||
-15
@@ -1,15 +0,0 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "WITH inserted_concurrency_counter AS (\n INSERT INTO concurrency_counter (concurrency_id, job_uuids) \n VALUES ($1, '{}'::jsonb)\n ON CONFLICT DO NOTHING\n )\n INSERT INTO concurrency_key(key, job_id) VALUES ($1, $2)",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Uuid"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "1bceaf6e9f25745b7f70128054ca81d68f3d56d4782e99e05b4f1cb362683514"
|
||||
}
|
||||
+17
@@ -0,0 +1,17 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE email_trigger SET script_path = $1, local_part = $2 WHERE workspace_id = $3 AND path = $4",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "2031c5138a785367e5127180ccb6734efa41ee6cb3b4819c1c517798b2b23e4a"
|
||||
}
|
||||
+87
@@ -0,0 +1,87 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n SELECT id, runnable_path, trigger_kind AS \"trigger_kind: String\",\n args AS \"args: sqlx::types::Json<serde_json::Value>\"\n FROM v2_job\n WHERE runnable_path = $1\n AND trigger_kind = $2::job_trigger_kind\n ORDER BY created_at DESC\n LIMIT 1\n ",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "id",
|
||||
"type_info": "Uuid"
|
||||
},
|
||||
{
|
||||
"ordinal": 1,
|
||||
"name": "runnable_path",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 2,
|
||||
"name": "trigger_kind: String",
|
||||
"type_info": {
|
||||
"Custom": {
|
||||
"name": "job_trigger_kind",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"webhook",
|
||||
"http",
|
||||
"websocket",
|
||||
"kafka",
|
||||
"email",
|
||||
"nats",
|
||||
"schedule",
|
||||
"app",
|
||||
"ui",
|
||||
"postgres",
|
||||
"sqs",
|
||||
"gcp",
|
||||
"mqtt",
|
||||
"nextcloud",
|
||||
"google"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"ordinal": 3,
|
||||
"name": "args: sqlx::types::Json<serde_json::Value>",
|
||||
"type_info": "Jsonb"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
{
|
||||
"Custom": {
|
||||
"name": "job_trigger_kind",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"webhook",
|
||||
"http",
|
||||
"websocket",
|
||||
"kafka",
|
||||
"email",
|
||||
"nats",
|
||||
"schedule",
|
||||
"app",
|
||||
"ui",
|
||||
"postgres",
|
||||
"sqs",
|
||||
"gcp",
|
||||
"mqtt",
|
||||
"nextcloud",
|
||||
"google"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false,
|
||||
true,
|
||||
true,
|
||||
true
|
||||
]
|
||||
},
|
||||
"hash": "212553c83e4dcdc6d045eb2fe2dadbb2860ce52d37a56b2861de1215260ecff8"
|
||||
}
|
||||
+23
@@ -0,0 +1,23 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "SELECT value FROM variable WHERE workspace_id = $1 AND path = $2 AND is_secret = true",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "value",
|
||||
"type_info": "Varchar"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "302bc55d0c227c5b12458ccde6569c4b531ff494d1ede0a655872ae05215f8a6"
|
||||
}
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE variable SET value = $1 WHERE workspace_id = $2 AND path = $3 AND is_secret = true",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "33fff66efe810ad7e92b36ca9b287938437182d7817387707e800519d9f5bafc"
|
||||
}
|
||||
+2
-2
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM raw_app WHERE path = $1 AND workspace_id = $2",
|
||||
"query": "DELETE FROM mqtt_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
@@ -11,5 +11,5 @@
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "b2b2b7251be2b80207f47b10eeff78d61c84161caaf16b8fd25d82b97aac1186"
|
||||
"hash": "3481e65196e500ad914f10b87884c19d8a3636ae955788c3366365194fe2ef57"
|
||||
}
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM email_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "40971d637c5b4d2af8e67872722880058cc067fca5e807ab3e1ed17d180cb7f7"
|
||||
}
|
||||
+15
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "WITH inserted_concurrency_counter AS (\n INSERT INTO concurrency_counter (concurrency_id, job_uuids)\n VALUES ($1, '{}'::jsonb)\n ON CONFLICT DO NOTHING\n )\n INSERT INTO concurrency_key(key, job_id) VALUES ($1, $2)",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Uuid"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "5a9cf9cc229f7b7ddbee2485cda4e9f7a91240a00f10d24d3f28d5b722f68768"
|
||||
}
|
||||
+39
@@ -0,0 +1,39 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n SELECT delivery_type AS \"delivery_type: String\",\n delivery_config\n FROM gcp_trigger\n WHERE workspace_id = $1 AND path = $2\n ",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "delivery_type: String",
|
||||
"type_info": {
|
||||
"Custom": {
|
||||
"name": "delivery_mode",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"push",
|
||||
"pull"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"ordinal": 1,
|
||||
"name": "delivery_config",
|
||||
"type_info": "Jsonb"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false,
|
||||
true
|
||||
]
|
||||
},
|
||||
"hash": "60a666a7cffce2d7631682095ccfd31e99997e88374496d4800eaa7b5152a464"
|
||||
}
|
||||
+29
@@ -0,0 +1,29 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "SELECT script_path, local_part FROM email_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "script_path",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 1,
|
||||
"name": "local_part",
|
||||
"type_info": "Varchar"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false,
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "8002f7f1521846056278f18e8b003fe48a4d3d827ac8b212e0f3d3bc2ed42f37"
|
||||
}
|
||||
+4
-4
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "SELECT EXISTS(SELECT 1 FROM raw_app WHERE path = $1 AND workspace_id = $2)",
|
||||
"query": "SELECT COUNT(*) FROM email_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "exists",
|
||||
"type_info": "Bool"
|
||||
"name": "count",
|
||||
"type_info": "Int8"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
@@ -19,5 +19,5 @@
|
||||
null
|
||||
]
|
||||
},
|
||||
"hash": "c4b3ab3d87f158b20f0557f78c6307c50620425aeb43636e19931c142b5ce0d7"
|
||||
"hash": "8364505b9a4d0d6f7406e4a5425360682ac2295bbf6746c7641141ed7bf47bf6"
|
||||
}
|
||||
-17
@@ -1,17 +0,0 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "INSERT INTO raw_app\n (workspace_id, path, summary, extra_perms, data)\n VALUES ($1, $2, $3, '{}', $4)",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "84bbd11f7bb0c65dcfb2e12038b8a9cb51d303480ac3c5f8a5e833a53c40ed9b"
|
||||
}
|
||||
+62
@@ -0,0 +1,62 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n SELECT mqtt_resource_path, client_version AS \"client_version: String\",\n script_path, mode AS \"mode: String\"\n FROM mqtt_trigger\n WHERE workspace_id = $1 AND path = $2\n ",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "mqtt_resource_path",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 1,
|
||||
"name": "client_version: String",
|
||||
"type_info": {
|
||||
"Custom": {
|
||||
"name": "mqtt_client_version",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"v3",
|
||||
"v5"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"ordinal": 2,
|
||||
"name": "script_path",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 3,
|
||||
"name": "mode: String",
|
||||
"type_info": {
|
||||
"Custom": {
|
||||
"name": "trigger_mode",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"enabled",
|
||||
"disabled",
|
||||
"suspended"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "899a162192d93c7fb7d47d38f4fb0fadcafc2bb0c61efa8ee496af7252d18270"
|
||||
}
|
||||
+22
@@ -0,0 +1,22 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n INSERT INTO kafka_trigger (\n path, kafka_resource_path, topics, group_id,\n script_path, is_flow, workspace_id, edited_by, email\n )\n VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)\n ",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"VarcharArray",
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Bool",
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Varchar"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "8c3cc09c1bbb6209467c75723dd02e97dddb99789422012e85e86c3151a9f2e9"
|
||||
}
|
||||
+23
@@ -0,0 +1,23 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "SELECT script_path FROM mqtt_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "script_path",
|
||||
"type_info": "Varchar"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "a086a5fa28dcff80911c7e5de73f204fc42625e58b967b0376ff7b382e38fa11"
|
||||
}
|
||||
+20
@@ -0,0 +1,20 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM job_result_stream_v2 WHERE job_id NOT IN (SELECT id FROM v2_job_queue) RETURNING job_id",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "job_id",
|
||||
"type_info": "Uuid"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": []
|
||||
},
|
||||
"nullable": [
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "a3e75f0309be42aca0fd74834f34b3f18dbb388bd8b9bc88b99aebedae9c3fec"
|
||||
}
|
||||
+84
@@ -0,0 +1,84 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n SELECT gcp_resource_path, topic_id, subscription_id,\n delivery_type AS \"delivery_type: String\",\n subscription_mode AS \"subscription_mode: String\",\n mode AS \"mode: String\"\n FROM gcp_trigger\n WHERE workspace_id = $1 AND path = $2\n ",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "gcp_resource_path",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 1,
|
||||
"name": "topic_id",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 2,
|
||||
"name": "subscription_id",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 3,
|
||||
"name": "delivery_type: String",
|
||||
"type_info": {
|
||||
"Custom": {
|
||||
"name": "delivery_mode",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"push",
|
||||
"pull"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"ordinal": 4,
|
||||
"name": "subscription_mode: String",
|
||||
"type_info": {
|
||||
"Custom": {
|
||||
"name": "gcp_subscription_mode",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"create_update",
|
||||
"existing"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"ordinal": 5,
|
||||
"name": "mode: String",
|
||||
"type_info": {
|
||||
"Custom": {
|
||||
"name": "trigger_mode",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"enabled",
|
||||
"disabled",
|
||||
"suspended"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "a80a17cda8ad722a97f16526d8576b76c7377066a06dba98a2921a68928aa644"
|
||||
}
|
||||
-14
@@ -1,14 +0,0 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM job_result_stream_v2 WHERE job_id = $1",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Uuid"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "bb46e5dcf5490ef3511faa131ad5693dedf34366e51044ddf30695995d194090"
|
||||
}
|
||||
+20
@@ -0,0 +1,20 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "DELETE FROM job_perms\nWHERE job_id NOT IN (SELECT id FROM v2_job_queue)\nRETURNING job_id",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "job_id",
|
||||
"type_info": "Uuid"
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": []
|
||||
},
|
||||
"nullable": [
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "c825fa5c6e287068aeaad994c0b42b8ad59b9129f032c6b918c27426ab304f2b"
|
||||
}
|
||||
+52
@@ -0,0 +1,52 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n SELECT local_part, workspaced_local_part, script_path,\n mode AS \"mode: String\"\n FROM email_trigger\n WHERE workspace_id = $1 AND path = $2\n ",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"ordinal": 0,
|
||||
"name": "local_part",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 1,
|
||||
"name": "workspaced_local_part",
|
||||
"type_info": "Bool"
|
||||
},
|
||||
{
|
||||
"ordinal": 2,
|
||||
"name": "script_path",
|
||||
"type_info": "Varchar"
|
||||
},
|
||||
{
|
||||
"ordinal": 3,
|
||||
"name": "mode: String",
|
||||
"type_info": {
|
||||
"Custom": {
|
||||
"name": "trigger_mode",
|
||||
"kind": {
|
||||
"Enum": [
|
||||
"enabled",
|
||||
"disabled",
|
||||
"suspended"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": [
|
||||
false,
|
||||
false,
|
||||
false,
|
||||
false
|
||||
]
|
||||
},
|
||||
"hash": "cabef30aa1e308bc656660c4ac0d3e77b020de05465f81a4fb4b0c07bd0d7439"
|
||||
}
|
||||
+2
-2
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "SELECT COUNT(*) as count FROM variable WHERE is_secret = true AND value != 'CLEARED'",
|
||||
"query": "SELECT COUNT(*) FROM variable WHERE is_secret = true",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
@@ -16,5 +16,5 @@
|
||||
null
|
||||
]
|
||||
},
|
||||
"hash": "052d42b46d5faba6b41f1fdcbf6a012db51b9e5a255ec0da9a8a0999d668d336"
|
||||
"hash": "d7d11b2c288fa0282539edde0c5450e4b5f8355821617b6c1ce3db64f6ceaf69"
|
||||
}
|
||||
+21
@@ -0,0 +1,21 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n INSERT INTO websocket_trigger (\n path, url, script_path, is_flow, workspace_id,\n edited_by, email, initial_messages\n )\n VALUES ($1, $2, $3, $4, $5, $6, $7, $8)\n ",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Bool",
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"Varchar",
|
||||
"JsonbArray"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "e485c82978d10379c6d1b7cd850f3bf764e6c5bf775f3d710ba36a79f5e424eb"
|
||||
}
|
||||
+2
-2
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "VACUUM v2_job, v2_job_completed, job_result_stream_v2, job_stats, job_logs, concurrency_key, log_file, metrics",
|
||||
"query": "VACUUM v2_job, v2_job_completed, job_result_stream_v2, job_stats, job_logs, job_perms, concurrency_key, log_file, metrics",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
@@ -8,5 +8,5 @@
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "f8ac5379ecfbff7b8ae75c821680737b249a64c8d9e8f7dbcc46fce98e874571"
|
||||
"hash": "f0070b36f7c4fc84dc9c23bb6c73d8ba80993a28b2c2e5df70968acf6d7cebe4"
|
||||
}
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE mqtt_trigger SET script_path = $1 WHERE workspace_id = $2 AND path = $3",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Varchar",
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "f0eb9d49021668673880aa20227c24a130cada07fba3f43e32d904fde97106d3"
|
||||
}
|
||||
+3
-2
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "SELECT COUNT(*) FROM raw_app WHERE workspace_id = $1",
|
||||
"query": "SELECT COUNT(*) FROM mqtt_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
@@ -11,6 +11,7 @@
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text"
|
||||
]
|
||||
},
|
||||
@@ -18,5 +19,5 @@
|
||||
null
|
||||
]
|
||||
},
|
||||
"hash": "3b5295a7c4b99aefa52c9a8ae1e0dd12bf4a0be1bf755caf7a1fa863e7950562"
|
||||
"hash": "f149d740301a65a62664fd407abe89c972e50837b6c54ac7a3ccd5f60ca05d3c"
|
||||
}
|
||||
Generated
+239
-78
@@ -1029,15 +1029,16 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "aws-sdk-sso"
|
||||
version = "1.77.0"
|
||||
version = "1.93.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "18f2f37fea82468fe3f5a059542c05392ef680c4f7f00e0db02df8b6e5c7d0c6"
|
||||
checksum = "9dcb38bb33fc0a11f1ffc3e3e85669e0a11a37690b86f77e75306d8f369146a0"
|
||||
dependencies = [
|
||||
"aws-credential-types",
|
||||
"aws-runtime",
|
||||
"aws-smithy-async",
|
||||
"aws-smithy-http 0.62.6",
|
||||
"aws-smithy-json 0.61.9",
|
||||
"aws-smithy-http 0.63.3",
|
||||
"aws-smithy-json 0.62.3",
|
||||
"aws-smithy-observability",
|
||||
"aws-smithy-runtime",
|
||||
"aws-smithy-runtime-api",
|
||||
"aws-smithy-types",
|
||||
@@ -1045,21 +1046,23 @@ dependencies = [
|
||||
"bytes",
|
||||
"fastrand",
|
||||
"http 0.2.12",
|
||||
"http 1.4.0",
|
||||
"regex-lite",
|
||||
"tracing",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "aws-sdk-ssooidc"
|
||||
version = "1.78.0"
|
||||
version = "1.95.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ecb4f6eada20e0193450cd48b12ed05e1e66baac86f39160191651b932f2b7d9"
|
||||
checksum = "2ada8ffbea7bd1be1f53df1dadb0f8fdb04badb13185b3321b929d1ee3caad09"
|
||||
dependencies = [
|
||||
"aws-credential-types",
|
||||
"aws-runtime",
|
||||
"aws-smithy-async",
|
||||
"aws-smithy-http 0.62.6",
|
||||
"aws-smithy-json 0.61.9",
|
||||
"aws-smithy-http 0.63.3",
|
||||
"aws-smithy-json 0.62.3",
|
||||
"aws-smithy-observability",
|
||||
"aws-smithy-runtime",
|
||||
"aws-smithy-runtime-api",
|
||||
"aws-smithy-types",
|
||||
@@ -1067,6 +1070,7 @@ dependencies = [
|
||||
"bytes",
|
||||
"fastrand",
|
||||
"http 0.2.12",
|
||||
"http 1.4.0",
|
||||
"regex-lite",
|
||||
"tracing",
|
||||
]
|
||||
@@ -6174,6 +6178,19 @@ dependencies = [
|
||||
"wasm-bindgen",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "getrandom"
|
||||
version = "0.4.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "139ef39800118c7683f2fd3c98c1b23c09ae076556b435f8e9064ae108aaeeec"
|
||||
dependencies = [
|
||||
"cfg-if",
|
||||
"libc",
|
||||
"r-efi",
|
||||
"wasip2",
|
||||
"wasip3",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "ghash"
|
||||
version = "0.5.1"
|
||||
@@ -7214,6 +7231,12 @@ dependencies = [
|
||||
"zerovec",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "id-arena"
|
||||
version = "2.3.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3d3067d79b975e8844ca9eb072e16b31c3c1c36928edf9c6789548c524d0d954"
|
||||
|
||||
[[package]]
|
||||
name = "ident_case"
|
||||
version = "1.0.1"
|
||||
@@ -7871,6 +7894,12 @@ version = "1.3.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "830d08ce1d1d941e6b30645f1a0eb5643013d835ce3779a5fc208261dbe10f55"
|
||||
|
||||
[[package]]
|
||||
name = "leb128fmt"
|
||||
version = "0.1.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "09edd9e8b54e49e587e4f6295a7d29c3ea94d469cb40ab8ca70b288248a81db2"
|
||||
|
||||
[[package]]
|
||||
name = "levenshtein_automata"
|
||||
version = "0.2.1"
|
||||
@@ -7936,9 +7965,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "libc"
|
||||
version = "0.2.180"
|
||||
version = "0.2.181"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "bcc35a38544a891a5f7c865aca548a982ccb3b8650a5b06d0fd33a10283c56fc"
|
||||
checksum = "459427e2af2b9c839b132acb702a1c654d95e10f8c326bfc2ad11310e458b1c5"
|
||||
|
||||
[[package]]
|
||||
name = "libffi"
|
||||
@@ -8413,9 +8442,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "memchr"
|
||||
version = "2.7.6"
|
||||
version = "2.8.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "f52b00d39961fc5b2736ea853c9cc86238e165017a493d1d5c8eac6bdc4cc273"
|
||||
checksum = "f8ca58f447f06ed17d5fc4043ce1b10dd205e060fb3ce5b979b8ed8e59ff3f79"
|
||||
|
||||
[[package]]
|
||||
name = "memmap2"
|
||||
@@ -10505,9 +10534,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "psm"
|
||||
version = "0.1.29"
|
||||
version = "0.1.30"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "1fa96cb91275ed31d6da3e983447320c4eb219ac180fa1679a0889ff32861e2d"
|
||||
checksum = "3852766467df634d74f0b2d7819bf8dc483a0eb2e3b0f50f756f9cfe8b0d18d8"
|
||||
dependencies = [
|
||||
"ar_archive_writer",
|
||||
"cc",
|
||||
@@ -10910,9 +10939,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rdkafka-sys"
|
||||
version = "4.9.0+2.10.0"
|
||||
version = "4.10.0+2.12.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "5230dca48bc354d718269f3e4353280e188b610f7af7e2fcf54b7a79d5802872"
|
||||
checksum = "e234cf318915c1059d4921ef7f75616b5219b10b46e9f3a511a15eb4b56a3f77"
|
||||
dependencies = [
|
||||
"cmake",
|
||||
"libc",
|
||||
@@ -11846,9 +11875,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "ryu"
|
||||
version = "1.0.22"
|
||||
version = "1.0.23"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "a50f4cf475b65d88e057964e0e9bb1f0aa9bbb2036dc65c64596b42932536984"
|
||||
checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f"
|
||||
|
||||
[[package]]
|
||||
name = "ryu-js"
|
||||
@@ -13006,9 +13035,9 @@ checksum = "6ce2be8dc25455e1f91df71bfa12ad37d7af1092ae736f3a6cd0e37bc7810596"
|
||||
|
||||
[[package]]
|
||||
name = "stacker"
|
||||
version = "0.1.22"
|
||||
version = "0.1.23"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "e1f8b29fb42aafcea4edeeb6b2f2d7ecd0d969c48b4cf0d2e64aafc471dd6e59"
|
||||
checksum = "08d74a23609d509411d10e2176dc2a4346e3b4aea2e7b1869f19fdedbc71c013"
|
||||
dependencies = [
|
||||
"cc",
|
||||
"cfg-if",
|
||||
@@ -13809,12 +13838,12 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tempfile"
|
||||
version = "3.24.0"
|
||||
version = "3.25.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "655da9c7eb6305c55742045d5a8d2037996d61d8de95806335c7c86ce0f82e9c"
|
||||
checksum = "0136791f7c95b1f6dd99f9cc786b91bb81c3800b639b3478e561ddb7be95e5f1"
|
||||
dependencies = [
|
||||
"fastrand",
|
||||
"getrandom 0.3.4",
|
||||
"getrandom 0.4.1",
|
||||
"once_cell",
|
||||
"rustix 1.1.3",
|
||||
"windows-sys 0.61.2",
|
||||
@@ -14958,9 +14987,9 @@ checksum = "81b79ad29b5e19de4260020f8919b443b2ef0277d242ce532ec7b7a2cc8b6007"
|
||||
|
||||
[[package]]
|
||||
name = "unicode-ident"
|
||||
version = "1.0.22"
|
||||
version = "1.0.23"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9312f7c4f6ff9069b165498234ce8be658059c6728633667c526e27dc2cf1df5"
|
||||
checksum = "537dd038a89878be9b64dd4bd1b260315c1bb94f4d784956b81e27a088d9a09e"
|
||||
|
||||
[[package]]
|
||||
name = "unicode-linebreak"
|
||||
@@ -15288,6 +15317,15 @@ dependencies = [
|
||||
"wit-bindgen",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wasip3"
|
||||
version = "0.4.0+wasi-0.3.0-rc-2026-01-06"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "5428f8bf88ea5ddc08faddef2ac4a67e390b88186c703ce6dbd955e1c145aca5"
|
||||
dependencies = [
|
||||
"wit-bindgen",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wasite"
|
||||
version = "0.1.0"
|
||||
@@ -15390,6 +15428,28 @@ dependencies = [
|
||||
"syn 2.0.114",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wasm-encoder"
|
||||
version = "0.244.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "990065f2fe63003fe337b932cfb5e3b80e0b4d0f5ff650e6985b1048f62c8319"
|
||||
dependencies = [
|
||||
"leb128fmt",
|
||||
"wasmparser",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wasm-metadata"
|
||||
version = "0.244.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "bb0e353e6a2fbdc176932bbaab493762eb1255a7900fe0fea1a2f96c296cc909"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"indexmap 2.11.1",
|
||||
"wasm-encoder",
|
||||
"wasmparser",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wasm-streams"
|
||||
version = "0.4.2"
|
||||
@@ -15413,6 +15473,18 @@ dependencies = [
|
||||
"thiserror 2.0.18",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wasmparser"
|
||||
version = "0.244.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "47b807c72e1bac69382b3a6fb3dbe8ea4c0ed87ff5629b8685ae6b9a611028fe"
|
||||
dependencies = [
|
||||
"bitflags 2.9.4",
|
||||
"hashbrown 0.15.5",
|
||||
"indexmap 2.11.1",
|
||||
"semver 1.0.27",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wasmtimer"
|
||||
version = "0.4.3"
|
||||
@@ -15653,9 +15725,13 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-nats",
|
||||
"aws-config",
|
||||
"aws-credential-types",
|
||||
"aws-sdk-sqs",
|
||||
"axum 0.7.9",
|
||||
"base64 0.22.1",
|
||||
"chrono",
|
||||
@@ -15669,7 +15745,9 @@ dependencies = [
|
||||
"once_cell",
|
||||
"prometheus",
|
||||
"rand 0.9.0",
|
||||
"rdkafka",
|
||||
"reqwest 0.13.1",
|
||||
"rumqttc",
|
||||
"rustls 0.23.35",
|
||||
"serde",
|
||||
"serde_derive",
|
||||
@@ -15705,7 +15783,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"argon2",
|
||||
@@ -15829,7 +15907,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-auth"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum 0.7.9",
|
||||
@@ -15855,7 +15933,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-client"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"reqwest 0.12.28",
|
||||
"serde",
|
||||
@@ -15865,7 +15943,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-embeddings"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum 0.7.9",
|
||||
@@ -15888,7 +15966,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-groups"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"axum 0.7.9",
|
||||
"chrono",
|
||||
@@ -15908,7 +15986,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-jobs"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum 0.7.9",
|
||||
@@ -15933,7 +16011,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-scripts"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"axum 0.7.9",
|
||||
"chrono",
|
||||
@@ -15962,7 +16040,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-sse"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"lazy_static",
|
||||
"serde",
|
||||
@@ -15974,7 +16052,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-users"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"argon2",
|
||||
"axum 0.7.9",
|
||||
@@ -15997,7 +16075,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-api-workspaces"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"axum 0.7.9",
|
||||
"chrono",
|
||||
@@ -16025,7 +16103,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-audit"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"lazy_static",
|
||||
@@ -16039,7 +16117,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-autoscaling"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum 0.7.9",
|
||||
@@ -16058,7 +16136,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-common"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"aes-gcm",
|
||||
"anyhow",
|
||||
@@ -16157,7 +16235,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-dep-map"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"itertools 0.14.0",
|
||||
@@ -16176,7 +16254,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-git-sync"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"regex",
|
||||
"serde",
|
||||
@@ -16191,7 +16269,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-indexer"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"astral-tokio-tar",
|
||||
@@ -16215,7 +16293,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-jseval"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"futures",
|
||||
@@ -16232,7 +16310,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-macros"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"itertools 0.14.0",
|
||||
"lazy_static",
|
||||
@@ -16248,7 +16326,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-mcp"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16269,7 +16347,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-native-triggers"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16300,7 +16378,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-oauth"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-oauth2",
|
||||
@@ -16324,7 +16402,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"convert_case 0.6.0",
|
||||
"serde",
|
||||
@@ -16333,7 +16411,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-bash"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"lazy_static",
|
||||
@@ -16345,7 +16423,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-csharp"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"serde_json",
|
||||
@@ -16357,7 +16435,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-go"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"gosyn",
|
||||
@@ -16369,7 +16447,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-graphql"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"lazy_static",
|
||||
@@ -16381,7 +16459,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-java"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"serde_json",
|
||||
@@ -16393,7 +16471,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-nu"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"nu-parser",
|
||||
@@ -16404,7 +16482,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-php"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"itertools 0.14.0",
|
||||
@@ -16415,7 +16493,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-py"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"itertools 0.14.0",
|
||||
@@ -16428,7 +16506,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-py-imports"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-recursion",
|
||||
@@ -16452,7 +16530,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-ruby"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"lazy_static",
|
||||
@@ -16466,7 +16544,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-rust"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"convert_case 0.6.0",
|
||||
@@ -16483,7 +16561,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-sql"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"lazy_static",
|
||||
@@ -16497,7 +16575,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-ts"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"lazy_static",
|
||||
@@ -16516,7 +16594,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-parser-yaml"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"serde",
|
||||
@@ -16527,7 +16605,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-queue"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-recursion",
|
||||
@@ -16564,7 +16642,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-runtime-nativets"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"const_format",
|
||||
@@ -16602,7 +16680,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-sql-datatype-parser-wasm"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"wasm-bindgen",
|
||||
"wasm-bindgen-test",
|
||||
@@ -16612,7 +16690,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-store"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-recursion",
|
||||
@@ -16641,7 +16719,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16674,7 +16752,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-email"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16694,7 +16772,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-gcp"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16728,7 +16806,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-http"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16762,7 +16840,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-kafka"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16785,7 +16863,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-mqtt"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16809,7 +16887,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-nats"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-nats",
|
||||
@@ -16833,7 +16911,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-postgres"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16868,7 +16946,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-sqs"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16876,6 +16954,7 @@ dependencies = [
|
||||
"aws-credential-types",
|
||||
"aws-sdk-sqs",
|
||||
"aws-sdk-sts",
|
||||
"aws-smithy-types",
|
||||
"axum 0.7.9",
|
||||
"backon",
|
||||
"chrono",
|
||||
@@ -16895,7 +16974,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-trigger-websocket"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
@@ -16918,7 +16997,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-types"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
@@ -16935,7 +17014,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "windmill-worker"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-once-cell",
|
||||
@@ -17649,6 +17728,88 @@ name = "wit-bindgen"
|
||||
version = "0.51.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "d7249219f66ced02969388cf2bb044a09756a083d0fab1e566056b04d9fbcaa5"
|
||||
dependencies = [
|
||||
"wit-bindgen-rust-macro",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wit-bindgen-core"
|
||||
version = "0.51.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ea61de684c3ea68cb082b7a88508a8b27fcc8b797d738bfc99a82facf1d752dc"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"heck 0.5.0",
|
||||
"wit-parser",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wit-bindgen-rust"
|
||||
version = "0.51.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b7c566e0f4b284dd6561c786d9cb0142da491f46a9fbed79ea69cdad5db17f21"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"heck 0.5.0",
|
||||
"indexmap 2.11.1",
|
||||
"prettyplease",
|
||||
"syn 2.0.114",
|
||||
"wasm-metadata",
|
||||
"wit-bindgen-core",
|
||||
"wit-component",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wit-bindgen-rust-macro"
|
||||
version = "0.51.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "0c0f9bfd77e6a48eccf51359e3ae77140a7f50b1e2ebfe62422d8afdaffab17a"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"prettyplease",
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"syn 2.0.114",
|
||||
"wit-bindgen-core",
|
||||
"wit-bindgen-rust",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wit-component"
|
||||
version = "0.244.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9d66ea20e9553b30172b5e831994e35fbde2d165325bec84fc43dbf6f4eb9cb2"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"bitflags 2.9.4",
|
||||
"indexmap 2.11.1",
|
||||
"log",
|
||||
"serde",
|
||||
"serde_derive",
|
||||
"serde_json",
|
||||
"wasm-encoder",
|
||||
"wasm-metadata",
|
||||
"wasmparser",
|
||||
"wit-parser",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "wit-parser"
|
||||
version = "0.244.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "ecc8ac4bc1dc3381b7f59c34f00b67e18f910c2c0f50015669dde7def656a736"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"id-arena",
|
||||
"indexmap 2.11.1",
|
||||
"log",
|
||||
"semver 1.0.27",
|
||||
"serde",
|
||||
"serde_derive",
|
||||
"serde_json",
|
||||
"unicode-xid",
|
||||
"wasmparser",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "writeable"
|
||||
|
||||
+9
-3
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "windmill"
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
authors.workspace = true
|
||||
edition.workspace = true
|
||||
|
||||
@@ -59,7 +59,7 @@ members = [
|
||||
exclude = ["./windmill-duckdb-ffi-internal"]
|
||||
|
||||
[workspace.package]
|
||||
version = "1.628.3"
|
||||
version = "1.629.1"
|
||||
authors = ["Ruben Fiszel <ruben@windmill.dev>"]
|
||||
edition = "2021"
|
||||
|
||||
@@ -229,6 +229,12 @@ axum.workspace = true
|
||||
serde.workspace = true
|
||||
windmill-api-client.workspace = true
|
||||
tempfile.workspace = true
|
||||
rumqttc.workspace = true
|
||||
rdkafka.workspace = true
|
||||
async-nats.workspace = true
|
||||
aws-sdk-sqs.workspace = true
|
||||
aws-config.workspace = true
|
||||
aws-credential-types.workspace = true
|
||||
|
||||
|
||||
[workspace.dependencies]
|
||||
@@ -413,7 +419,7 @@ const_format = { version = "0.2.35", features = ["rust_1_64", "rust_1_51"] }
|
||||
const-str = "0.5"
|
||||
constant_time_eq = "0.3.1"
|
||||
rsa = "^0"
|
||||
aes-gcm = "^0"
|
||||
aes-gcm = "0.10.3"
|
||||
async_zip = { version = "0.0.17", features = ["tokio", "tokio-fs", "deflate", "chrono"] }
|
||||
once_cell = "1.17.1"
|
||||
dashmap = "6.1.0"
|
||||
|
||||
@@ -1 +1 @@
|
||||
cd2713671bc7ac5762302704ee4ecb81e023301f
|
||||
7596cefdba81482c0b0c0b61be26369f112d8009
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
-- No-op: cannot restore deleted telemetry data
|
||||
@@ -0,0 +1,2 @@
|
||||
-- Delete all saved telemetry data from metrics table
|
||||
DELETE FROM metrics WHERE id = 'telemetry';
|
||||
@@ -381,5 +381,6 @@ pub static SHORT_IMPORTS_MAP: PyMap = phf_map! {
|
||||
"docx" => "python-docx",
|
||||
"vt" => "vt-py",
|
||||
"grpc" => "grpcio",
|
||||
"serpapi" => "google-search-results",
|
||||
// Add new entry here ^
|
||||
};
|
||||
|
||||
@@ -422,11 +422,15 @@ fn parse_expr(
|
||||
match e.as_ref() {
|
||||
Expr::Name(ExprName { id, .. }) => (parse_typ(id.as_ref(), enums, module), false),
|
||||
Expr::Attribute(x) => {
|
||||
if x.value
|
||||
.as_name_expr()
|
||||
.is_some_and(|x| x.id.as_str() == "wmill")
|
||||
{
|
||||
(parse_typ(x.attr.as_str(), enums, module), false)
|
||||
if let Some(name) = x.value.as_name_expr() {
|
||||
match name.id.as_str() {
|
||||
"wmill" => (parse_typ(x.attr.as_str(), enums, module), false),
|
||||
"datetime" => {
|
||||
let full_name = format!("datetime.{}", x.attr.as_str());
|
||||
(parse_typ(&full_name, enums, module), false)
|
||||
}
|
||||
_ => (Typ::Unknown, false),
|
||||
}
|
||||
} else {
|
||||
(Typ::Unknown, false)
|
||||
}
|
||||
@@ -493,6 +497,8 @@ fn parse_typ(id: &str, enums: &HashMap<String, EnumInfo>, module: Option<&[Stmt]
|
||||
"bytes" => Typ::Bytes,
|
||||
"datetime" => Typ::Datetime,
|
||||
"datetime.datetime" => Typ::Datetime,
|
||||
"date" => Typ::Date,
|
||||
"datetime.date" => Typ::Date,
|
||||
"Sql" | "sql" => Typ::Sql,
|
||||
x @ _ if x.starts_with("DynSelect_") => {
|
||||
Typ::DynSelect(x.strip_prefix("DynSelect_").unwrap().to_string())
|
||||
@@ -620,7 +626,7 @@ def main(test1: str, name: datetime.datetime = datetime.now(), byte: bytes = byt
|
||||
Arg {
|
||||
otyp: None,
|
||||
name: "name".to_string(),
|
||||
typ: Typ::Unknown,
|
||||
typ: Typ::Datetime,
|
||||
default: Some(json!("<function call>")),
|
||||
has_default: true,
|
||||
oidx: None
|
||||
@@ -709,7 +715,7 @@ def main(test1: str,
|
||||
Arg {
|
||||
otyp: None,
|
||||
name: "name".to_string(),
|
||||
typ: Typ::Unknown,
|
||||
typ: Typ::Datetime,
|
||||
default: Some(json!("<function call>")),
|
||||
has_default: true,
|
||||
oidx: None
|
||||
|
||||
@@ -295,6 +295,7 @@ fn extract_field_type(
|
||||
"bool" => Typ::Bool,
|
||||
"bytes" => Typ::Bytes,
|
||||
"datetime" => Typ::Datetime,
|
||||
"date" => Typ::Date,
|
||||
"Any" => Typ::Unknown, // typing.Any maps to Unknown
|
||||
// Custom class - check if it's a model
|
||||
custom_type => {
|
||||
|
||||
@@ -65,6 +65,7 @@ pub enum Typ {
|
||||
List(Box<Typ>),
|
||||
Bytes,
|
||||
Datetime,
|
||||
Date,
|
||||
Resource(String),
|
||||
Email,
|
||||
Sql,
|
||||
|
||||
@@ -10,6 +10,7 @@ fn to_str(typ: Typ) -> String {
|
||||
Typ::List(t) => format!("list-{}", to_str(*t)),
|
||||
Typ::Bytes => "bytes".to_string(),
|
||||
Typ::Datetime => "datetime".to_string(),
|
||||
Typ::Date => "date".to_string(),
|
||||
_ => "unknown".to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
+46
-2
@@ -1103,7 +1103,6 @@ async fn delete_expired_jobs_batch(
|
||||
job_retention_secs
|
||||
);
|
||||
|
||||
// Delete related records for this batch
|
||||
if let Err(e) = sqlx::query!(
|
||||
"DELETE FROM job_stats WHERE job_id = ANY($1)",
|
||||
&deleted_jobs
|
||||
@@ -1783,6 +1782,19 @@ pub async fn monitor_db(
|
||||
}
|
||||
};
|
||||
|
||||
let cleanup_job_live_rows_f = async {
|
||||
if server_mode && !initial_load {
|
||||
if let Some(db) = conn.as_sql() {
|
||||
if let Err(e) = cleanup_job_perms_orphaned(&db).await {
|
||||
tracing::error!("Error cleaning up orphaned job_perms: {:?}", e);
|
||||
}
|
||||
if let Err(e) = cleanup_job_result_stream_orphaned_jobs(&db).await {
|
||||
tracing::error!("Error cleaning up orphaned job_result_stream_v2: {:?}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
// run every hour (60 minutes / 30 seconds = 120)
|
||||
let cleanup_worker_group_stats_f = async {
|
||||
if server_mode && iteration.is_some() && iteration.as_ref().unwrap().should_run(120) {
|
||||
@@ -1956,6 +1968,7 @@ pub async fn monitor_db(
|
||||
cleanup_debounce_keys_f,
|
||||
cleanup_debounce_keys_completed_f,
|
||||
cleanup_flow_iterator_data_f,
|
||||
cleanup_job_live_rows_f,
|
||||
cleanup_worker_group_stats_f,
|
||||
native_triggers_sync_f,
|
||||
cleanup_notify_events_f,
|
||||
@@ -1963,7 +1976,7 @@ pub async fn monitor_db(
|
||||
}
|
||||
|
||||
async fn vacuuming_tables(db: &Pool<Postgres>) -> error::Result<()> {
|
||||
sqlx::query!("VACUUM v2_job, v2_job_completed, job_result_stream_v2, job_stats, job_logs, concurrency_key, log_file, metrics")
|
||||
sqlx::query!("VACUUM v2_job, v2_job_completed, job_result_stream_v2, job_stats, job_logs, job_perms, concurrency_key, log_file, metrics")
|
||||
.execute(db)
|
||||
.await?;
|
||||
Ok(())
|
||||
@@ -3062,6 +3075,37 @@ RETURNING key,job_id
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn cleanup_job_perms_orphaned(db: &DB) -> error::Result<()> {
|
||||
let result = sqlx::query_scalar!(
|
||||
"DELETE FROM job_perms
|
||||
WHERE job_id NOT IN (SELECT id FROM v2_job_queue)
|
||||
RETURNING job_id"
|
||||
)
|
||||
.fetch_all(db)
|
||||
.await?;
|
||||
|
||||
if !result.is_empty() {
|
||||
tracing::info!("Cleaned up {} orphaned job_perms rows", result.len());
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn cleanup_job_result_stream_orphaned_jobs(db: &DB) -> error::Result<()> {
|
||||
let result = sqlx::query!(
|
||||
"DELETE FROM job_result_stream_v2 WHERE job_id NOT IN (SELECT id FROM v2_job_queue) RETURNING job_id",
|
||||
)
|
||||
.fetch_all(db)
|
||||
.await?;
|
||||
|
||||
if result.len() > 0 {
|
||||
tracing::info!(
|
||||
"Cleaned up {} orphaned job_result_stream_v2 rows",
|
||||
result.len()
|
||||
);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn cleanup_flow_iterator_data_orphaned_jobs(db: &DB) -> error::Result<()> {
|
||||
let result = sqlx::query!(
|
||||
"
|
||||
|
||||
@@ -87,6 +87,12 @@ impl ApiServer {
|
||||
Self::start_inner(db, true).await
|
||||
}
|
||||
|
||||
/// Start the API server with server_mode=true so trigger listeners are active.
|
||||
/// Alias for `start_agent_mode` with a clearer name for trigger e2e tests.
|
||||
pub async fn start_with_listeners(db: Pool<Postgres>) -> anyhow::Result<Self> {
|
||||
Self::start_inner(db, true).await
|
||||
}
|
||||
|
||||
async fn start_inner(db: Pool<Postgres>, agent_mode: bool) -> anyhow::Result<Self> {
|
||||
let (tx, rx) = tokio::sync::broadcast::channel::<()>(1);
|
||||
|
||||
@@ -143,11 +149,17 @@ pub struct RunJob {
|
||||
pub payload: JobPayload,
|
||||
pub args: serde_json::Map<String, serde_json::Value>,
|
||||
pub scheduled_for_o: Option<chrono::DateTime<chrono::Utc>>,
|
||||
pub email: String,
|
||||
}
|
||||
|
||||
impl From<JobPayload> for RunJob {
|
||||
fn from(payload: JobPayload) -> Self {
|
||||
Self { payload, args: Default::default(), scheduled_for_o: None }
|
||||
Self {
|
||||
payload,
|
||||
args: Default::default(),
|
||||
scheduled_for_o: None,
|
||||
email: "test@windmill.dev".to_string(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -165,8 +177,13 @@ impl RunJob {
|
||||
self
|
||||
}
|
||||
|
||||
pub fn email(mut self, email: impl Into<String>) -> Self {
|
||||
self.email = email.into();
|
||||
self
|
||||
}
|
||||
|
||||
pub async fn push(self, db: &Pool<Postgres>) -> Uuid {
|
||||
let RunJob { payload, args, scheduled_for_o } = self;
|
||||
let RunJob { payload, args, scheduled_for_o, email } = self;
|
||||
let mut hm_args = std::collections::HashMap::new();
|
||||
for (k, v) in args {
|
||||
hm_args.insert(k, windmill_common::worker::to_raw_value(&v));
|
||||
@@ -180,7 +197,7 @@ impl RunJob {
|
||||
payload,
|
||||
windmill_queue::PushArgs::from(&hm_args),
|
||||
/* user */ "test-user",
|
||||
/* email */ "test@windmill.dev",
|
||||
/* email */ &email,
|
||||
/* permissioned_as */ "u/test-user".to_string(),
|
||||
/* token_prefix */ None,
|
||||
scheduled_for_o,
|
||||
|
||||
+6
-6
@@ -14,18 +14,18 @@ INSERT INTO workspace_key(workspace_id, kind, key)
|
||||
VALUES ('test-workspace-2', 'cloud', 'test-key-2')
|
||||
ON CONFLICT DO NOTHING;
|
||||
|
||||
-- Insert test secrets for workspace 1
|
||||
-- Note: The 'value' column stores encrypted values in production,
|
||||
-- but for tests we'll use plain text that the migration will handle
|
||||
-- Insert test variables with placeholder values.
|
||||
-- Secret values are encrypted by the test setup using build_crypt + encrypt
|
||||
-- with the workspace key, matching production behavior.
|
||||
INSERT INTO variable (workspace_id, path, value, is_secret, description, extra_perms)
|
||||
VALUES
|
||||
('test-workspace', 'u/test-user/db_password', 'encrypted-db-pass-123', true, 'Database password', '{}'),
|
||||
('test-workspace', 'u/test-user/api_key', 'encrypted-api-key-abc', true, 'API key for external service', '{}'),
|
||||
('test-workspace', 'u/test-user/db_password', 'PLACEHOLDER', true, 'Database password', '{}'),
|
||||
('test-workspace', 'u/test-user/api_key', 'PLACEHOLDER', true, 'API key for external service', '{}'),
|
||||
('test-workspace', 'u/test-user/public_var', 'not-a-secret', false, 'A non-secret variable', '{}')
|
||||
ON CONFLICT DO NOTHING;
|
||||
|
||||
-- Insert test secrets for workspace 2 (to test isolation)
|
||||
INSERT INTO variable (workspace_id, path, value, is_secret, description, extra_perms)
|
||||
VALUES
|
||||
('test-workspace-2', 'u/test-user/other_secret', 'encrypted-other-secret', true, 'Secret in workspace 2', '{}')
|
||||
('test-workspace-2', 'u/test-user/other_secret', 'PLACEHOLDER', true, 'Secret in workspace 2', '{}')
|
||||
ON CONFLICT DO NOTHING;
|
||||
|
||||
+58
@@ -0,0 +1,58 @@
|
||||
#!/usr/bin/env bash
|
||||
# Starts (or stops) all external services needed for the trigger e2e tests.
|
||||
#
|
||||
# Usage:
|
||||
# ./tests/fixtures/start_all_triggers.sh # start everything
|
||||
# ./tests/fixtures/start_all_triggers.sh stop # stop everything
|
||||
# ./tests/fixtures/start_all_triggers.sh oss # start only OSS services
|
||||
set -euo pipefail
|
||||
|
||||
DIR="$(cd "$(dirname "$0")" && pwd)"
|
||||
ACTION="${1:-start}"
|
||||
|
||||
SCRIPTS_OSS=(
|
||||
"$DIR/start_mqtt.sh"
|
||||
"$DIR/start_websocket.sh"
|
||||
"$DIR/start_postgres_replication.sh"
|
||||
)
|
||||
|
||||
SCRIPTS_EE=(
|
||||
"$DIR/start_kafka.sh"
|
||||
"$DIR/start_nats.sh"
|
||||
"$DIR/start_sqs.sh"
|
||||
"$DIR/start_gcp_pubsub.sh"
|
||||
)
|
||||
|
||||
if [[ "$ACTION" == "stop" ]]; then
|
||||
for s in "${SCRIPTS_OSS[@]}" "${SCRIPTS_EE[@]}"; do
|
||||
echo "--- $(basename "$s" .sh) stop ---"
|
||||
bash "$s" stop
|
||||
done
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if [[ "$ACTION" == "oss" ]]; then
|
||||
SCRIPTS=("${SCRIPTS_OSS[@]}")
|
||||
else
|
||||
SCRIPTS=("${SCRIPTS_OSS[@]}" "${SCRIPTS_EE[@]}")
|
||||
fi
|
||||
|
||||
for s in "${SCRIPTS[@]}"; do
|
||||
echo "--- $(basename "$s" .sh) ---"
|
||||
bash "$s"
|
||||
echo ""
|
||||
done
|
||||
|
||||
echo "============================================"
|
||||
echo "All services ready. Run the e2e tests with:"
|
||||
echo ""
|
||||
|
||||
if [[ "$ACTION" == "oss" ]]; then
|
||||
echo " cargo test --test trigger_e2e --features mqtt_trigger,websocket,postgres_trigger -- --ignored --nocapture"
|
||||
else
|
||||
echo " # OSS triggers"
|
||||
echo " cargo test --test trigger_e2e --features mqtt_trigger,websocket,postgres_trigger -- --ignored --nocapture"
|
||||
echo ""
|
||||
echo " # Enterprise triggers"
|
||||
echo " AWS_ENDPOINT_URL=http://localhost:4566 PUBSUB_EMULATOR_HOST=localhost:8085 cargo test --test trigger_e2e --features kafka,nats,sqs_trigger,gcp_trigger,enterprise,private -- --ignored --nocapture"
|
||||
fi
|
||||
+47
@@ -0,0 +1,47 @@
|
||||
#!/usr/bin/env bash
|
||||
# Starts the GCP Pub/Sub emulator for trigger_e2e::test_gcp_e2e (Enterprise)
|
||||
#
|
||||
# Usage:
|
||||
# ./tests/fixtures/start_gcp_pubsub.sh # start
|
||||
# ./tests/fixtures/start_gcp_pubsub.sh stop # stop & remove
|
||||
set -euo pipefail
|
||||
|
||||
NAME="windmill-test-pubsub"
|
||||
PORT=8085
|
||||
|
||||
if [[ "${1:-}" == "stop" ]]; then
|
||||
docker rm -f "$NAME" 2>/dev/null && echo "stopped $NAME" || echo "$NAME not running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if docker ps --format '{{.Names}}' | grep -q "^${NAME}$"; then
|
||||
echo "$NAME is already running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
docker rm -f "$NAME" 2>/dev/null || true
|
||||
|
||||
docker run -d --name "$NAME" -p "${PORT}:8085" \
|
||||
gcr.io/google.com/cloudsdktool/google-cloud-cli:emulators \
|
||||
gcloud beta emulators pubsub start --host-port="0.0.0.0:${PORT}"
|
||||
|
||||
echo "Waiting for Pub/Sub emulator to become ready..."
|
||||
for i in $(seq 1 30); do
|
||||
if curl -sf "http://localhost:${PORT}" &>/dev/null; then
|
||||
break
|
||||
fi
|
||||
sleep 1
|
||||
done
|
||||
|
||||
# Create the test topic and subscription
|
||||
curl -sX PUT "http://localhost:${PORT}/v1/projects/local-project/topics/windmill-e2e-test" >/dev/null
|
||||
curl -sX PUT "http://localhost:${PORT}/v1/projects/local-project/subscriptions/windmill-e2e-sub" \
|
||||
-H "Content-Type: application/json" \
|
||||
-d '{"topic": "projects/local-project/topics/windmill-e2e-test"}' >/dev/null
|
||||
|
||||
echo "GCP Pub/Sub emulator listening on localhost:${PORT}"
|
||||
echo " topic: windmill-e2e-test"
|
||||
echo " subscription: windmill-e2e-sub"
|
||||
echo ""
|
||||
echo "Run the test:"
|
||||
echo " PUBSUB_EMULATOR_HOST=localhost:${PORT} cargo test --test trigger_e2e test_gcp_e2e --features gcp_trigger,enterprise,private -- --ignored --nocapture"
|
||||
+52
@@ -0,0 +1,52 @@
|
||||
#!/usr/bin/env bash
|
||||
# Starts a Kafka broker for trigger_e2e::test_kafka_e2e (Enterprise)
|
||||
#
|
||||
# Usage:
|
||||
# ./tests/fixtures/start_kafka.sh # start
|
||||
# ./tests/fixtures/start_kafka.sh stop # stop & remove
|
||||
set -euo pipefail
|
||||
|
||||
NAME="windmill-test-kafka"
|
||||
PORT=9092
|
||||
|
||||
if [[ "${1:-}" == "stop" ]]; then
|
||||
docker rm -f "$NAME" 2>/dev/null && echo "stopped $NAME" || echo "$NAME not running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if docker ps --format '{{.Names}}' | grep -q "^${NAME}$"; then
|
||||
echo "$NAME is already running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
docker rm -f "$NAME" 2>/dev/null || true
|
||||
|
||||
docker run -d --name "$NAME" -p "${PORT}:9092" \
|
||||
-e KAFKA_NODE_ID=1 \
|
||||
-e KAFKA_PROCESS_ROLES=broker,controller \
|
||||
-e KAFKA_LISTENERS="PLAINTEXT://0.0.0.0:${PORT},CONTROLLER://0.0.0.0:9093" \
|
||||
-e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT \
|
||||
-e KAFKA_CONTROLLER_QUORUM_VOTERS=1@localhost:9093 \
|
||||
-e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
|
||||
-e KAFKA_ADVERTISED_LISTENERS="PLAINTEXT://localhost:${PORT}" \
|
||||
-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
|
||||
-e KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS=0 \
|
||||
apache/kafka:latest
|
||||
|
||||
echo "Waiting for Kafka to become ready..."
|
||||
for i in $(seq 1 60); do
|
||||
if docker exec "$NAME" /opt/kafka/bin/kafka-topics.sh --list --bootstrap-server "localhost:${PORT}" &>/dev/null; then
|
||||
break
|
||||
fi
|
||||
sleep 1
|
||||
done
|
||||
|
||||
docker exec "$NAME" /opt/kafka/bin/kafka-topics.sh --create \
|
||||
--topic windmill-e2e-test \
|
||||
--bootstrap-server "localhost:${PORT}" \
|
||||
--partitions 1 --replication-factor 1 2>/dev/null || true
|
||||
|
||||
echo "Kafka broker listening on localhost:${PORT} with topic 'windmill-e2e-test'"
|
||||
echo ""
|
||||
echo "Run the test:"
|
||||
echo " cargo test --test trigger_e2e test_kafka_e2e --features kafka,enterprise,private -- --ignored --nocapture"
|
||||
+31
@@ -0,0 +1,31 @@
|
||||
#!/usr/bin/env bash
|
||||
# Starts a Mosquitto MQTT broker for trigger_e2e::test_mqtt_e2e
|
||||
#
|
||||
# Usage:
|
||||
# ./tests/fixtures/start_mqtt.sh # start
|
||||
# ./tests/fixtures/start_mqtt.sh stop # stop & remove
|
||||
set -euo pipefail
|
||||
|
||||
NAME="windmill-test-mqtt"
|
||||
PORT=1883
|
||||
|
||||
if [[ "${1:-}" == "stop" ]]; then
|
||||
docker rm -f "$NAME" 2>/dev/null && echo "stopped $NAME" || echo "$NAME not running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if docker ps --format '{{.Names}}' | grep -q "^${NAME}$"; then
|
||||
echo "$NAME is already running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
docker rm -f "$NAME" 2>/dev/null || true
|
||||
|
||||
docker run -d --name "$NAME" -p "${PORT}:1883" \
|
||||
eclipse-mosquitto:latest \
|
||||
mosquitto -c /mosquitto-no-auth.conf
|
||||
|
||||
echo "MQTT broker listening on localhost:${PORT}"
|
||||
echo ""
|
||||
echo "Run the test:"
|
||||
echo " cargo test --test trigger_e2e test_mqtt_e2e --features mqtt_trigger -- --ignored --nocapture"
|
||||
+29
@@ -0,0 +1,29 @@
|
||||
#!/usr/bin/env bash
|
||||
# Starts a NATS server for trigger_e2e::test_nats_e2e (Enterprise)
|
||||
#
|
||||
# Usage:
|
||||
# ./tests/fixtures/start_nats.sh # start
|
||||
# ./tests/fixtures/start_nats.sh stop # stop & remove
|
||||
set -euo pipefail
|
||||
|
||||
NAME="windmill-test-nats"
|
||||
PORT=4222
|
||||
|
||||
if [[ "${1:-}" == "stop" ]]; then
|
||||
docker rm -f "$NAME" 2>/dev/null && echo "stopped $NAME" || echo "$NAME not running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if docker ps --format '{{.Names}}' | grep -q "^${NAME}$"; then
|
||||
echo "$NAME is already running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
docker rm -f "$NAME" 2>/dev/null || true
|
||||
|
||||
docker run -d --name "$NAME" -p "${PORT}:4222" nats:latest
|
||||
|
||||
echo "NATS server listening on localhost:${PORT}"
|
||||
echo ""
|
||||
echo "Run the test:"
|
||||
echo " cargo test --test trigger_e2e test_nats_e2e --features nats,enterprise,private -- --ignored --nocapture"
|
||||
+53
@@ -0,0 +1,53 @@
|
||||
#!/usr/bin/env bash
|
||||
# Configures the local PostgreSQL for logical replication trigger tests.
|
||||
#
|
||||
# Prerequisites: wal_level=logical must be set (requires PG restart).
|
||||
# Check with: SHOW wal_level;
|
||||
#
|
||||
# Usage:
|
||||
# ./tests/fixtures/start_postgres_replication.sh # setup
|
||||
# ./tests/fixtures/start_postgres_replication.sh stop # teardown
|
||||
set -euo pipefail
|
||||
|
||||
PGURL="${DATABASE_URL:-postgres://postgres:changeme@localhost:5432/windmill}"
|
||||
|
||||
if [[ "${1:-}" == "stop" ]]; then
|
||||
psql "$PGURL" <<'SQL'
|
||||
SELECT pg_drop_replication_slot('test_e2e_slot')
|
||||
WHERE EXISTS (SELECT 1 FROM pg_replication_slots WHERE slot_name = 'test_e2e_slot');
|
||||
DROP PUBLICATION IF EXISTS test_e2e_pub;
|
||||
DROP TABLE IF EXISTS test_trigger_table;
|
||||
SQL
|
||||
echo "Postgres replication teardown complete"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
# Check wal_level
|
||||
WAL_LEVEL=$(psql "$PGURL" -tAc "SHOW wal_level;")
|
||||
if [[ "$WAL_LEVEL" != "logical" ]]; then
|
||||
echo "ERROR: wal_level is '$WAL_LEVEL', must be 'logical'"
|
||||
echo ""
|
||||
echo "Fix with:"
|
||||
echo " psql \"$PGURL\" -c \"ALTER SYSTEM SET wal_level = logical;\""
|
||||
echo " # then restart PostgreSQL"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
psql "$PGURL" <<'SQL'
|
||||
CREATE TABLE IF NOT EXISTS test_trigger_table (id serial PRIMARY KEY, data text);
|
||||
|
||||
DO $$
|
||||
BEGIN
|
||||
IF NOT EXISTS (SELECT 1 FROM pg_publication WHERE pubname = 'test_e2e_pub') THEN
|
||||
CREATE PUBLICATION test_e2e_pub FOR TABLE test_trigger_table;
|
||||
END IF;
|
||||
END $$;
|
||||
|
||||
SELECT pg_create_logical_replication_slot('test_e2e_slot', 'pgoutput')
|
||||
WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots WHERE slot_name = 'test_e2e_slot');
|
||||
SQL
|
||||
|
||||
echo "Postgres logical replication ready (publication=test_e2e_pub, slot=test_e2e_slot)"
|
||||
echo ""
|
||||
echo "Run the test:"
|
||||
echo " cargo test --test trigger_e2e test_postgres_e2e --features postgres_trigger -- --ignored --nocapture"
|
||||
+45
@@ -0,0 +1,45 @@
|
||||
#!/usr/bin/env bash
|
||||
# Starts LocalStack for SQS trigger_e2e::test_sqs_e2e (Enterprise)
|
||||
#
|
||||
# Usage:
|
||||
# ./tests/fixtures/start_sqs.sh # start
|
||||
# ./tests/fixtures/start_sqs.sh stop # stop & remove
|
||||
set -euo pipefail
|
||||
|
||||
NAME="windmill-test-localstack"
|
||||
PORT=4566
|
||||
|
||||
if [[ "${1:-}" == "stop" ]]; then
|
||||
docker rm -f "$NAME" 2>/dev/null && echo "stopped $NAME" || echo "$NAME not running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if docker ps --format '{{.Names}}' | grep -q "^${NAME}$"; then
|
||||
echo "$NAME is already running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
docker rm -f "$NAME" 2>/dev/null || true
|
||||
|
||||
docker run -d --name "$NAME" -p "${PORT}:4566" \
|
||||
-e SERVICES=sqs \
|
||||
localstack/localstack
|
||||
|
||||
echo "Waiting for LocalStack to become ready..."
|
||||
for i in $(seq 1 30); do
|
||||
if curl -sf "http://localhost:${PORT}/_localstack/health" &>/dev/null; then
|
||||
break
|
||||
fi
|
||||
sleep 1
|
||||
done
|
||||
|
||||
# Create the test queue
|
||||
aws --endpoint-url="http://localhost:${PORT}" \
|
||||
--region us-east-1 \
|
||||
--no-sign-request \
|
||||
sqs create-queue --queue-name windmill-e2e-test 2>/dev/null || true
|
||||
|
||||
echo "LocalStack SQS listening on localhost:${PORT} with queue 'windmill-e2e-test'"
|
||||
echo ""
|
||||
echo "Run the test:"
|
||||
echo " AWS_ENDPOINT_URL=http://localhost:${PORT} cargo test --test trigger_e2e test_sqs_e2e --features sqs_trigger,enterprise,private -- --ignored --nocapture"
|
||||
+31
@@ -0,0 +1,31 @@
|
||||
#!/usr/bin/env bash
|
||||
# Starts a WebSocket echo server for trigger_e2e::test_websocket_e2e
|
||||
#
|
||||
# Usage:
|
||||
# ./tests/fixtures/start_websocket.sh # start
|
||||
# ./tests/fixtures/start_websocket.sh stop # stop & remove
|
||||
set -euo pipefail
|
||||
|
||||
NAME="windmill-test-ws-echo"
|
||||
PORT=8765
|
||||
|
||||
if [[ "${1:-}" == "stop" ]]; then
|
||||
docker rm -f "$NAME" 2>/dev/null && echo "stopped $NAME" || echo "$NAME not running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if docker ps --format '{{.Names}}' | grep -q "^${NAME}$"; then
|
||||
echo "$NAME is already running"
|
||||
exit 0
|
||||
fi
|
||||
|
||||
docker rm -f "$NAME" 2>/dev/null || true
|
||||
|
||||
docker run -d --name "$NAME" -p "${PORT}:8080" \
|
||||
-e PORT=8080 \
|
||||
jmalloc/echo-server
|
||||
|
||||
echo "WebSocket echo server listening on localhost:${PORT}"
|
||||
echo ""
|
||||
echo "Run the test:"
|
||||
echo " cargo test --test trigger_e2e test_websocket_e2e --features websocket -- --ignored --nocapture"
|
||||
@@ -26,28 +26,56 @@
|
||||
#[cfg(all(feature = "private", feature = "enterprise"))]
|
||||
mod tests {
|
||||
use sqlx::{Pool, Postgres};
|
||||
use std::collections::HashMap;
|
||||
use windmill_common::secret_backend::{
|
||||
migrate_secrets_to_database, migrate_secrets_to_vault, test_vault_connection,
|
||||
SecretBackend, VaultBackend, VaultSettings,
|
||||
};
|
||||
use windmill_common::variables::{build_crypt, decrypt, encrypt};
|
||||
|
||||
/// Plaintext values for test secrets. The fixture inserts PLACEHOLDERs that
|
||||
/// `encrypt_fixture_secrets` replaces with properly encrypted values.
|
||||
const TEST_SECRETS: &[(&str, &str, &str)] = &[
|
||||
("test-workspace", "u/test-user/db_password", "db-pass-123"),
|
||||
("test-workspace", "u/test-user/api_key", "api-key-abc"),
|
||||
(
|
||||
"test-workspace-2",
|
||||
"u/test-user/other_secret",
|
||||
"other-secret-value",
|
||||
),
|
||||
];
|
||||
|
||||
/// Encrypt the PLACEHOLDER values inserted by the fixture using the real
|
||||
/// workspace encryption keys, so migration tests can decrypt them correctly.
|
||||
async fn encrypt_fixture_secrets(db: &Pool<Postgres>) {
|
||||
for &(workspace_id, path, plaintext) in TEST_SECRETS {
|
||||
let mc = build_crypt(db, workspace_id).await.unwrap();
|
||||
let encrypted = encrypt(&mc, plaintext);
|
||||
sqlx::query!(
|
||||
"UPDATE variable SET value = $1 WHERE workspace_id = $2 AND path = $3 AND is_secret = true",
|
||||
encrypted,
|
||||
workspace_id,
|
||||
path,
|
||||
)
|
||||
.execute(db)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
/// Check if vault tests should run (requires RUN_VAULT_TESTS=1 env var)
|
||||
fn should_run_vault_tests() -> bool {
|
||||
std::env::var("RUN_VAULT_TESTS")
|
||||
.map(|v| v == "1" || v.to_lowercase() == "true")
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
/// Set up BASE_URL for JWT tests (required for OIDC issuer URL generation)
|
||||
#[cfg(feature = "openidconnect")]
|
||||
async fn setup_base_url() {
|
||||
let base_url = std::env::var("BASE_URL")
|
||||
.unwrap_or_else(|_| "http://localhost:8000".to_string());
|
||||
let base_url =
|
||||
std::env::var("BASE_URL").unwrap_or_else(|_| "http://localhost:8000".to_string());
|
||||
let mut url = windmill_common::BASE_URL.write().await;
|
||||
*url = base_url;
|
||||
}
|
||||
|
||||
/// Skip test if RUN_VAULT_TESTS is not set
|
||||
macro_rules! skip_if_no_vault {
|
||||
() => {
|
||||
if !should_run_vault_tests() {
|
||||
@@ -71,6 +99,7 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "openidconnect")]
|
||||
fn vault_settings_jwt() -> VaultSettings {
|
||||
VaultSettings {
|
||||
address: std::env::var("VAULT_ADDR")
|
||||
@@ -163,6 +192,7 @@ mod tests {
|
||||
// ==================== JWT Auth Tests ====================
|
||||
|
||||
/// Test Vault connection with JWT authentication
|
||||
#[cfg(feature = "openidconnect")]
|
||||
#[sqlx::test(fixtures("base", "secret_backend"))]
|
||||
async fn test_vault_connection_jwt(db: Pool<Postgres>) {
|
||||
skip_if_no_vault!();
|
||||
@@ -236,70 +266,51 @@ mod tests {
|
||||
|
||||
let settings = vault_settings_static_token();
|
||||
|
||||
// Verify Vault connection
|
||||
test_vault_connection(&settings, Some(&db))
|
||||
.await
|
||||
.expect("Failed to connect to Vault");
|
||||
|
||||
// Check initial state
|
||||
let secrets_before = sqlx::query!(
|
||||
"SELECT workspace_id, path, value FROM variable WHERE is_secret = true ORDER BY workspace_id, path"
|
||||
)
|
||||
.fetch_all(&db)
|
||||
.await
|
||||
.expect("Failed to query secrets");
|
||||
// Encrypt fixture placeholders with real workspace keys
|
||||
encrypt_fixture_secrets(&db).await;
|
||||
|
||||
println!(
|
||||
"Found {} secrets in database before migration:",
|
||||
secrets_before.len()
|
||||
);
|
||||
for s in &secrets_before {
|
||||
println!(" - {}/{}: {} chars", s.workspace_id, s.path, s.value.len());
|
||||
}
|
||||
let secret_count = sqlx::query_scalar!(
|
||||
"SELECT COUNT(*) FROM variable WHERE is_secret = true"
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await
|
||||
.expect("Failed to count secrets");
|
||||
println!("Found {} secrets in database before migration", secret_count.unwrap_or(0));
|
||||
|
||||
// Run migration
|
||||
println!("\nMigrating secrets to Vault...");
|
||||
println!("Migrating secrets to Vault...");
|
||||
let report = migrate_secrets_to_vault(&db, &settings)
|
||||
.await
|
||||
.expect("Migration to Vault failed");
|
||||
|
||||
println!("Migration report:");
|
||||
println!(" Total secrets: {}", report.total_secrets);
|
||||
println!(" Migrated: {}", report.migrated_count);
|
||||
println!(" Failed: {}", report.failed_count);
|
||||
println!("Migration report: total={}, migrated={}, failed={}",
|
||||
report.total_secrets, report.migrated_count, report.failed_count);
|
||||
|
||||
if !report.failures.is_empty() {
|
||||
println!(" Failures:");
|
||||
for f in &report.failures {
|
||||
println!(" - {}/{}: {}", f.workspace_id, f.path, f.error);
|
||||
println!(" FAIL: {}/{}: {}", f.workspace_id, f.path, f.error);
|
||||
}
|
||||
}
|
||||
|
||||
assert_eq!(report.failed_count, 0, "Migration had failures");
|
||||
assert!(report.migrated_count > 0, "No secrets were migrated");
|
||||
|
||||
// Verify secrets in Vault
|
||||
println!("\nVerifying secrets in Vault...");
|
||||
// Verify decrypted values in Vault match original plaintexts
|
||||
let vault_backend = VaultBackend::new(settings.clone());
|
||||
|
||||
for secret in &secrets_before {
|
||||
let result = vault_backend
|
||||
.get_secret(&secret.workspace_id, &secret.path)
|
||||
.await;
|
||||
assert!(
|
||||
result.is_ok(),
|
||||
"Failed to read secret {}/{} from Vault: {:?}",
|
||||
secret.workspace_id,
|
||||
secret.path,
|
||||
result.err()
|
||||
);
|
||||
println!(
|
||||
" ✓ {}/{} exists in Vault",
|
||||
secret.workspace_id, secret.path
|
||||
);
|
||||
for &(ws, path, expected_plaintext) in TEST_SECRETS {
|
||||
let value = vault_backend
|
||||
.get_secret(ws, path)
|
||||
.await
|
||||
.unwrap_or_else(|e| panic!("Failed to read {}/{} from Vault: {:?}", ws, path, e));
|
||||
assert_eq!(value, expected_plaintext, "Vault value mismatch for {}/{}", ws, path);
|
||||
println!(" ✓ {}/{} correct in Vault", ws, path);
|
||||
}
|
||||
|
||||
println!("\n✓ Migration to Vault completed successfully");
|
||||
println!("✓ Migration to Vault completed successfully");
|
||||
}
|
||||
|
||||
/// Test migration from Vault back to database
|
||||
@@ -313,49 +324,50 @@ mod tests {
|
||||
.await
|
||||
.expect("Failed to connect to Vault");
|
||||
|
||||
encrypt_fixture_secrets(&db).await;
|
||||
|
||||
// First migrate TO Vault
|
||||
println!("Setting up: migrating secrets to Vault first...");
|
||||
let to_vault = migrate_secrets_to_vault(&db, &settings)
|
||||
.await
|
||||
.expect("Initial migration to Vault failed");
|
||||
assert!(to_vault.migrated_count > 0, "No secrets to test with");
|
||||
println!(" Migrated {} secrets to Vault", to_vault.migrated_count);
|
||||
|
||||
// Clear database values
|
||||
println!("\nClearing database secret values...");
|
||||
sqlx::query!("UPDATE variable SET value = 'CLEARED' WHERE is_secret = true")
|
||||
.execute(&db)
|
||||
.await
|
||||
.expect("Failed to clear values");
|
||||
|
||||
// Migrate back from Vault
|
||||
println!("\nMigrating secrets from Vault to database...");
|
||||
println!("Migrating secrets from Vault to database...");
|
||||
let report = migrate_secrets_to_database(&db, &settings)
|
||||
.await
|
||||
.expect("Migration to database failed");
|
||||
|
||||
println!("Migration report:");
|
||||
println!(" Total secrets: {}", report.total_secrets);
|
||||
println!(" Migrated: {}", report.migrated_count);
|
||||
println!(" Failed: {}", report.failed_count);
|
||||
println!("Migration report: total={}, migrated={}, failed={}",
|
||||
report.total_secrets, report.migrated_count, report.failed_count);
|
||||
|
||||
assert_eq!(report.failed_count, 0, "Migration had failures");
|
||||
assert!(report.migrated_count > 0, "No secrets were migrated");
|
||||
|
||||
// Verify restored
|
||||
let restored = sqlx::query!(
|
||||
"SELECT COUNT(*) as count FROM variable WHERE is_secret = true AND value != 'CLEARED'"
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await
|
||||
.expect("Failed to count restored");
|
||||
// Verify restored values decrypt to the original plaintexts
|
||||
for &(ws, path, expected_plaintext) in TEST_SECRETS {
|
||||
let row = sqlx::query_scalar!(
|
||||
"SELECT value FROM variable WHERE workspace_id = $1 AND path = $2 AND is_secret = true",
|
||||
ws, path
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await
|
||||
.unwrap_or_else(|_| panic!("Secret {}/{} not found in DB after migration", ws, path));
|
||||
|
||||
assert!(
|
||||
restored.count.unwrap_or(0) > 0,
|
||||
"No secrets were restored in database"
|
||||
);
|
||||
let mc = build_crypt(&db, ws).await.unwrap();
|
||||
let decrypted = decrypt(&mc, row).expect("Failed to decrypt restored value");
|
||||
assert_eq!(decrypted, expected_plaintext, "Restored value mismatch for {}/{}", ws, path);
|
||||
println!(" ✓ {}/{} correctly restored in DB", ws, path);
|
||||
}
|
||||
|
||||
println!("\n✓ Migration to database completed successfully");
|
||||
println!("✓ Migration to database completed successfully");
|
||||
}
|
||||
|
||||
// ==================== Variable Rename Tests ====================
|
||||
@@ -443,68 +455,47 @@ mod tests {
|
||||
.await
|
||||
.expect("Failed to connect to Vault");
|
||||
|
||||
// Get original secrets
|
||||
let original: HashMap<(String, String), String> = sqlx::query!(
|
||||
"SELECT workspace_id, path, value FROM variable WHERE is_secret = true"
|
||||
)
|
||||
.fetch_all(&db)
|
||||
.await
|
||||
.expect("Failed to query")
|
||||
.into_iter()
|
||||
.map(|r| ((r.workspace_id, r.path), r.value))
|
||||
.collect();
|
||||
encrypt_fixture_secrets(&db).await;
|
||||
|
||||
println!("Original secrets: {} entries", original.len());
|
||||
|
||||
// Step 1: DB -> Vault
|
||||
println!("\n=== Step 1: Migrate DB -> Vault ===");
|
||||
println!("=== Step 1: Migrate DB -> Vault ===");
|
||||
let to_vault = migrate_secrets_to_vault(&db, &settings)
|
||||
.await
|
||||
.expect("Migration to Vault failed");
|
||||
println!("Migrated {} secrets to Vault", to_vault.migrated_count);
|
||||
assert_eq!(to_vault.failed_count, 0);
|
||||
println!("Migrated {} secrets to Vault", to_vault.migrated_count);
|
||||
|
||||
// Step 2: Clear DB
|
||||
println!("\n=== Step 2: Clear database values ===");
|
||||
println!("=== Step 2: Clear database values ===");
|
||||
sqlx::query!("UPDATE variable SET value = 'ROUND_TRIP_CLEARED' WHERE is_secret = true")
|
||||
.execute(&db)
|
||||
.await
|
||||
.expect("Failed to clear");
|
||||
|
||||
// Step 3: Vault -> DB
|
||||
println!("\n=== Step 3: Migrate Vault -> DB ===");
|
||||
println!("=== Step 3: Migrate Vault -> DB ===");
|
||||
let to_db = migrate_secrets_to_database(&db, &settings)
|
||||
.await
|
||||
.expect("Migration to database failed");
|
||||
println!("Migrated {} secrets to database", to_db.migrated_count);
|
||||
assert_eq!(to_db.failed_count, 0);
|
||||
println!("Migrated {} secrets to database", to_db.migrated_count);
|
||||
|
||||
// Step 4: Verify
|
||||
println!("\n=== Step 4: Verify round-trip integrity ===");
|
||||
let restored: HashMap<(String, String), String> = sqlx::query!(
|
||||
"SELECT workspace_id, path, value FROM variable WHERE is_secret = true"
|
||||
)
|
||||
.fetch_all(&db)
|
||||
.await
|
||||
.expect("Failed to query")
|
||||
.into_iter()
|
||||
.map(|r| ((r.workspace_id, r.path), r.value))
|
||||
.collect();
|
||||
|
||||
for ((ws, path), _) in &original {
|
||||
let restored_value = restored
|
||||
.get(&(ws.clone(), path.clone()))
|
||||
.expect(&format!("Secret {}/{} not found after round-trip", ws, path));
|
||||
|
||||
assert_ne!(
|
||||
restored_value, "ROUND_TRIP_CLEARED",
|
||||
"Secret {}/{} was not restored",
|
||||
println!("=== Step 4: Verify round-trip integrity ===");
|
||||
for &(ws, path, expected_plaintext) in TEST_SECRETS {
|
||||
let encrypted = sqlx::query_scalar!(
|
||||
"SELECT value FROM variable WHERE workspace_id = $1 AND path = $2 AND is_secret = true",
|
||||
ws, path
|
||||
);
|
||||
println!(" ✓ {}/{}: restored", ws, path);
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await
|
||||
.unwrap_or_else(|_| panic!("Secret {}/{} not found after round-trip", ws, path));
|
||||
|
||||
assert_ne!(encrypted, "ROUND_TRIP_CLEARED", "Secret {}/{} was not restored", ws, path);
|
||||
|
||||
let mc = build_crypt(&db, ws).await.unwrap();
|
||||
let decrypted = decrypt(&mc, encrypted).expect("Failed to decrypt");
|
||||
assert_eq!(decrypted, expected_plaintext, "Round-trip value mismatch for {}/{}", ws, path);
|
||||
println!(" ✓ {}/{}: round-trip OK", ws, path);
|
||||
}
|
||||
|
||||
println!("\n✓ Full round-trip completed successfully!");
|
||||
println!("✓ Full round-trip completed successfully!");
|
||||
}
|
||||
|
||||
// ==================== Workspace Isolation Test ====================
|
||||
@@ -517,7 +508,8 @@ mod tests {
|
||||
let settings = vault_settings_static_token();
|
||||
let backend = VaultBackend::new(settings.clone());
|
||||
|
||||
// First migrate secrets to Vault
|
||||
encrypt_fixture_secrets(&db).await;
|
||||
|
||||
migrate_secrets_to_vault(&db, &settings)
|
||||
.await
|
||||
.expect("Migration failed");
|
||||
@@ -542,7 +534,7 @@ mod tests {
|
||||
assert!(ws1.is_ok(), "Same-workspace access should work");
|
||||
println!("✓ Same-workspace access works");
|
||||
|
||||
println!("\n✓ Workspace isolation verified!");
|
||||
println!("✓ Workspace isolation verified!");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,726 @@
|
||||
/*!
|
||||
* End-to-end integration tests for Windmill trigger listeners.
|
||||
*
|
||||
* Each test is `#[ignore]` because it requires a running external service
|
||||
* (MQTT broker, NATS server, Kafka broker, etc.). See individual test doc
|
||||
* comments for setup instructions.
|
||||
*
|
||||
* Quick start — use the helper scripts in `tests/fixtures/`:
|
||||
* ```bash
|
||||
* ./tests/fixtures/start_all_triggers.sh # start all services
|
||||
* ./tests/fixtures/start_all_triggers.sh oss # start OSS services only
|
||||
* ./tests/fixtures/start_all_triggers.sh stop # tear down everything
|
||||
* ```
|
||||
*
|
||||
* The general pattern:
|
||||
* 1. Insert a test script + trigger row + resource into the DB
|
||||
* 2. Start the API server with listeners enabled (server_mode=true)
|
||||
* 3. Connect to the external service and send a test message
|
||||
* 4. Poll `v2_job` for a job matching the trigger path + trigger_kind
|
||||
* 5. Verify the args shape/content
|
||||
*/
|
||||
|
||||
use serde_json::json;
|
||||
use sqlx::{Pool, Postgres};
|
||||
use std::time::Duration;
|
||||
|
||||
mod common;
|
||||
use common::*;
|
||||
|
||||
// ============================================================================
|
||||
// Helpers
|
||||
// ============================================================================
|
||||
|
||||
/// Row shape for polling v2_job.
|
||||
#[derive(Debug)]
|
||||
#[allow(dead_code)]
|
||||
struct TriggerJobRow {
|
||||
id: uuid::Uuid,
|
||||
runnable_path: Option<String>,
|
||||
trigger_kind: Option<String>,
|
||||
args: Option<sqlx::types::Json<serde_json::Value>>,
|
||||
}
|
||||
|
||||
/// Poll `v2_job` every 500ms for up to `timeout` for a job whose
|
||||
/// `runnable_path` and `trigger_kind` match the expected values.
|
||||
async fn poll_for_trigger_job(
|
||||
db: &Pool<Postgres>,
|
||||
script_path: &str,
|
||||
trigger_kind: &str,
|
||||
timeout: Duration,
|
||||
) -> anyhow::Result<TriggerJobRow> {
|
||||
let deadline = tokio::time::Instant::now() + timeout;
|
||||
loop {
|
||||
let row = sqlx::query_as!(
|
||||
TriggerJobRow,
|
||||
r#"
|
||||
SELECT id, runnable_path, trigger_kind AS "trigger_kind: String",
|
||||
args AS "args: sqlx::types::Json<serde_json::Value>"
|
||||
FROM v2_job
|
||||
WHERE runnable_path = $1
|
||||
AND trigger_kind = $2::job_trigger_kind
|
||||
ORDER BY created_at DESC
|
||||
LIMIT 1
|
||||
"#,
|
||||
script_path,
|
||||
trigger_kind as _,
|
||||
)
|
||||
.fetch_optional(db)
|
||||
.await?;
|
||||
|
||||
if let Some(job) = row {
|
||||
return Ok(job);
|
||||
}
|
||||
|
||||
if tokio::time::Instant::now() >= deadline {
|
||||
anyhow::bail!(
|
||||
"timed out waiting for trigger job (script_path={}, trigger_kind={})",
|
||||
script_path,
|
||||
trigger_kind
|
||||
);
|
||||
}
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(500)).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Insert a minimal test script row that trigger listeners can reference.
|
||||
async fn insert_test_script(db: &Pool<Postgres>, path: &str) -> anyhow::Result<i64> {
|
||||
let hash: i64 = rand::random::<i64>().unsigned_abs() as i64;
|
||||
sqlx::query(
|
||||
"INSERT INTO script (workspace_id, hash, path, summary, description, content,
|
||||
created_by, language, kind, lock)
|
||||
VALUES ('test-workspace', $1, $2, '', '', 'def main(): pass',
|
||||
'test-user', 'python3', 'script', '')",
|
||||
)
|
||||
.bind(hash)
|
||||
.bind(path)
|
||||
.execute(db)
|
||||
.await?;
|
||||
Ok(hash)
|
||||
}
|
||||
|
||||
/// Insert a resource row for triggers that resolve connection details from the
|
||||
/// `resource` table.
|
||||
async fn insert_resource(
|
||||
db: &Pool<Postgres>,
|
||||
path: &str,
|
||||
resource_type: &str,
|
||||
value: serde_json::Value,
|
||||
) -> anyhow::Result<()> {
|
||||
sqlx::query(
|
||||
"INSERT INTO resource (workspace_id, path, value, resource_type, extra_perms, created_by)
|
||||
VALUES ('test-workspace', $1, $2::jsonb, $3, '{}'::jsonb, 'test-user')",
|
||||
)
|
||||
.bind(path)
|
||||
.bind(value)
|
||||
.bind(resource_type)
|
||||
.execute(db)
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// MQTT Trigger E2E
|
||||
// ============================================================================
|
||||
|
||||
/// End-to-end test for MQTT trigger.
|
||||
///
|
||||
/// Requires a running MQTT broker. Setup:
|
||||
/// ```bash
|
||||
/// ./tests/fixtures/start_mqtt.sh
|
||||
/// ```
|
||||
///
|
||||
/// Run:
|
||||
/// ```bash
|
||||
/// cargo test --test trigger_e2e test_mqtt_e2e --features mqtt_trigger \
|
||||
/// -- --ignored --nocapture
|
||||
/// ```
|
||||
#[ignore = "requires running MQTT broker on localhost:1883"]
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_mqtt_e2e(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
initialize_tracing().await;
|
||||
|
||||
let script_path = "f/test/mqtt_e2e_handler";
|
||||
insert_test_script(&db, script_path).await?;
|
||||
|
||||
insert_resource(
|
||||
&db,
|
||||
"u/test-user/mqtt_res",
|
||||
"mqtt",
|
||||
json!({
|
||||
"broker": "localhost",
|
||||
"port": 1883
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO mqtt_trigger (
|
||||
path, mqtt_resource_path, subscribe_topics, client_version,
|
||||
script_path, is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, ARRAY[$3::jsonb], $4::mqtt_client_version, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
)
|
||||
.bind("f/test/mqtt_e2e_trigger")
|
||||
.bind("u/test-user/mqtt_res")
|
||||
.bind(json!({"topic": "windmill/test/e2e", "qos": "qos0"}))
|
||||
.bind("v5")
|
||||
.bind(script_path)
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let _server = ApiServer::start_with_listeners(db.clone()).await?;
|
||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||
|
||||
// Publish a message using rumqttc
|
||||
let mut mqtt_opts = rumqttc::MqttOptions::new("windmill-e2e-test", "localhost", 1883);
|
||||
mqtt_opts.set_keep_alive(Duration::from_secs(5));
|
||||
let (client, mut eventloop) = rumqttc::AsyncClient::new(mqtt_opts, 10);
|
||||
|
||||
// Drive the event loop in the background
|
||||
let el_handle = tokio::spawn(async move {
|
||||
loop {
|
||||
match eventloop.poll().await {
|
||||
Ok(_) => {}
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(500)).await;
|
||||
client
|
||||
.publish(
|
||||
"windmill/test/e2e",
|
||||
rumqttc::QoS::AtLeastOnce,
|
||||
false,
|
||||
b"hello from e2e test".to_vec(),
|
||||
)
|
||||
.await?;
|
||||
|
||||
let job = poll_for_trigger_job(&db, script_path, "mqtt", Duration::from_secs(30)).await?;
|
||||
assert!(job.args.is_some(), "job should have args");
|
||||
|
||||
client.disconnect().await.ok();
|
||||
el_handle.abort();
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// WebSocket Trigger E2E
|
||||
// ============================================================================
|
||||
|
||||
/// End-to-end test for WebSocket trigger.
|
||||
///
|
||||
/// Requires a WebSocket echo server. Setup:
|
||||
/// ```bash
|
||||
/// ./tests/fixtures/start_websocket.sh
|
||||
/// ```
|
||||
///
|
||||
/// Run:
|
||||
/// ```bash
|
||||
/// cargo test --test trigger_e2e test_websocket_e2e --features websocket \
|
||||
/// -- --ignored --nocapture
|
||||
/// ```
|
||||
#[ignore = "requires running WebSocket echo server on localhost:8765"]
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_websocket_e2e(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
initialize_tracing().await;
|
||||
|
||||
let script_path = "f/test/ws_e2e_handler";
|
||||
insert_test_script(&db, script_path).await?;
|
||||
|
||||
sqlx::query!(
|
||||
r#"
|
||||
INSERT INTO websocket_trigger (
|
||||
path, url, script_path, is_flow, workspace_id,
|
||||
edited_by, email, initial_messages
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
|
||||
"#,
|
||||
"f/test/ws_e2e_trigger",
|
||||
"ws://localhost:8765",
|
||||
script_path,
|
||||
false,
|
||||
"test-workspace",
|
||||
"test-user",
|
||||
"test@windmill.dev",
|
||||
&[json!({"type": "RawMessage", "content": "hello from e2e test"})] as &[serde_json::Value],
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let _server = ApiServer::start_with_listeners(db.clone()).await?;
|
||||
|
||||
// The WebSocket trigger connects to the server and sends initial_messages,
|
||||
// and each received message triggers a job.
|
||||
let job = poll_for_trigger_job(&db, script_path, "websocket", Duration::from_secs(30)).await?;
|
||||
assert!(job.args.is_some(), "job should have args");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Postgres Trigger E2E
|
||||
// ============================================================================
|
||||
|
||||
/// End-to-end test for Postgres trigger (logical replication).
|
||||
///
|
||||
/// Requires PostgreSQL with `wal_level=logical`. Setup:
|
||||
/// ```bash
|
||||
/// ./tests/fixtures/start_postgres_replication.sh
|
||||
/// ```
|
||||
/// (The script checks wal_level and creates the table/publication/slot in the
|
||||
/// main DB. This test re-creates them in its isolated sqlx::test database.)
|
||||
///
|
||||
/// Run:
|
||||
/// ```bash
|
||||
/// cargo test --test trigger_e2e test_postgres_e2e --features postgres_trigger \
|
||||
/// -- --ignored --nocapture
|
||||
/// ```
|
||||
#[ignore = "requires PostgreSQL with wal_level=logical"]
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_postgres_e2e(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
initialize_tracing().await;
|
||||
|
||||
let script_path = "f/test/pg_e2e_handler";
|
||||
insert_test_script(&db, script_path).await?;
|
||||
|
||||
// Create the tracked table + publication + replication slot inside the
|
||||
// isolated test database (sqlx::test gives us a fresh DB each run).
|
||||
// Replication slots are server-wide so we use a random suffix.
|
||||
let suffix: u32 = rand::random();
|
||||
let slot_name = format!("test_e2e_slot_{suffix}");
|
||||
let pub_name = format!("test_e2e_pub_{suffix}");
|
||||
|
||||
sqlx::query("CREATE TABLE test_trigger_table (id serial PRIMARY KEY, data text)")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
sqlx::query(&format!("CREATE PUBLICATION {pub_name} FOR TABLE test_trigger_table"))
|
||||
.execute(&db)
|
||||
.await?;
|
||||
sqlx::query(&format!(
|
||||
"SELECT pg_create_logical_replication_slot('{slot_name}', 'pgoutput')"
|
||||
))
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
// Extract the test DB name from the pool so the resource points here,
|
||||
// not at the main windmill database.
|
||||
let test_db_name: String =
|
||||
sqlx::query_scalar("SELECT current_database()")
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
insert_resource(
|
||||
&db,
|
||||
"u/test-user/pg_res",
|
||||
"postgresql",
|
||||
json!({
|
||||
"user": "postgres",
|
||||
"password": "changeme",
|
||||
"host": "localhost",
|
||||
"port": 5432,
|
||||
"dbname": test_db_name,
|
||||
"sslmode": "disable"
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO postgres_trigger (
|
||||
path, script_path, is_flow, workspace_id, edited_by, email,
|
||||
postgres_resource_path, replication_slot_name, publication_name
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
)
|
||||
.bind("f/test/pg_e2e_trigger")
|
||||
.bind(script_path)
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.bind("u/test-user/pg_res")
|
||||
.bind(&slot_name)
|
||||
.bind(&pub_name)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let _server = ApiServer::start_with_listeners(db.clone()).await?;
|
||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||
|
||||
// Insert a row into the tracked table to trigger a change event
|
||||
sqlx::query("INSERT INTO test_trigger_table (data) VALUES ('e2e test data')")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let job = poll_for_trigger_job(&db, script_path, "postgres", Duration::from_secs(30)).await?;
|
||||
assert!(job.args.is_some(), "job should have args");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Kafka Trigger E2E (Enterprise)
|
||||
// ============================================================================
|
||||
|
||||
/// End-to-end test for Kafka trigger (Enterprise only).
|
||||
///
|
||||
/// Requires a running Kafka broker with the test topic. Setup:
|
||||
/// ```bash
|
||||
/// ./tests/fixtures/start_kafka.sh
|
||||
/// ```
|
||||
///
|
||||
/// Run:
|
||||
/// ```bash
|
||||
/// cargo test --test trigger_e2e test_kafka_e2e \
|
||||
/// --features kafka,enterprise,private -- --ignored --nocapture
|
||||
/// ```
|
||||
#[cfg(all(feature = "enterprise", feature = "private"))]
|
||||
#[ignore = "requires running Kafka broker on localhost:9092"]
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_kafka_e2e(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
initialize_tracing().await;
|
||||
|
||||
let script_path = "f/test/kafka_e2e_handler";
|
||||
insert_test_script(&db, script_path).await?;
|
||||
|
||||
insert_resource(
|
||||
&db,
|
||||
"u/test-user/kafka_res",
|
||||
"kafka",
|
||||
json!({
|
||||
"brokers": ["localhost:9092"],
|
||||
"security": { "label": "PLAINTEXT" }
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
r#"
|
||||
INSERT INTO kafka_trigger (
|
||||
path, kafka_resource_path, topics, group_id,
|
||||
script_path, is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
"f/test/kafka_e2e_trigger",
|
||||
"u/test-user/kafka_res",
|
||||
&["windmill-e2e-test"] as &[&str],
|
||||
"windmill-e2e-test-group",
|
||||
script_path,
|
||||
false,
|
||||
"test-workspace",
|
||||
"test-user",
|
||||
"test@windmill.dev",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let _server = ApiServer::start_with_listeners(db.clone()).await?;
|
||||
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||
|
||||
// Produce messages using rdkafka. The consumer starts with auto.offset.reset=latest
|
||||
// and needs time for group rebalance, so we send repeatedly until a job appears.
|
||||
use rdkafka::config::ClientConfig;
|
||||
use rdkafka::producer::{FutureProducer, FutureRecord};
|
||||
|
||||
let producer: FutureProducer = ClientConfig::new()
|
||||
.set("bootstrap.servers", "localhost:9092")
|
||||
.create()?;
|
||||
|
||||
let db2 = db.clone();
|
||||
let produce_handle = tokio::spawn(async move {
|
||||
for _ in 0..30 {
|
||||
let _ = producer
|
||||
.send(
|
||||
FutureRecord::to("windmill-e2e-test")
|
||||
.payload("hello from kafka e2e test")
|
||||
.key("test-key"),
|
||||
Duration::from_secs(5),
|
||||
)
|
||||
.await;
|
||||
tokio::time::sleep(Duration::from_secs(1)).await;
|
||||
}
|
||||
});
|
||||
|
||||
let job = poll_for_trigger_job(&db2, script_path, "kafka", Duration::from_secs(30)).await?;
|
||||
produce_handle.abort();
|
||||
assert!(job.args.is_some(), "job should have args");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// NATS Trigger E2E (Enterprise)
|
||||
// ============================================================================
|
||||
|
||||
/// End-to-end test for NATS trigger (Enterprise only).
|
||||
///
|
||||
/// Requires a running NATS server. Setup:
|
||||
/// ```bash
|
||||
/// ./tests/fixtures/start_nats.sh
|
||||
/// ```
|
||||
///
|
||||
/// Run:
|
||||
/// ```bash
|
||||
/// cargo test --test trigger_e2e test_nats_e2e \
|
||||
/// --features nats,enterprise,private -- --ignored --nocapture
|
||||
/// ```
|
||||
#[cfg(all(feature = "enterprise", feature = "private"))]
|
||||
#[ignore = "requires running NATS server on localhost:4222"]
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_nats_e2e(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
initialize_tracing().await;
|
||||
|
||||
let script_path = "f/test/nats_e2e_handler";
|
||||
insert_test_script(&db, script_path).await?;
|
||||
|
||||
insert_resource(
|
||||
&db,
|
||||
"u/test-user/nats_res",
|
||||
"nats",
|
||||
json!({
|
||||
"servers": ["nats://localhost:4222"],
|
||||
"auth": { "label": "NO_AUTH" },
|
||||
"require_tls": false
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
r#"
|
||||
INSERT INTO nats_trigger (
|
||||
path, nats_resource_path, subjects, script_path,
|
||||
is_flow, workspace_id, edited_by, email, use_jetstream
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
"f/test/nats_e2e_trigger",
|
||||
"u/test-user/nats_res",
|
||||
&["windmill.e2e.test"] as &[&str],
|
||||
script_path,
|
||||
false,
|
||||
"test-workspace",
|
||||
"test-user",
|
||||
"test@windmill.dev",
|
||||
false,
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let _server = ApiServer::start_with_listeners(db.clone()).await?;
|
||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||
|
||||
// Publish a message using async-nats
|
||||
let nats_client = async_nats::connect("localhost:4222").await?;
|
||||
nats_client
|
||||
.publish("windmill.e2e.test", "hello from nats e2e test".into())
|
||||
.await?;
|
||||
nats_client.flush().await?;
|
||||
|
||||
let job = poll_for_trigger_job(&db, script_path, "nats", Duration::from_secs(30)).await?;
|
||||
assert!(job.args.is_some(), "job should have args");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// SQS Trigger E2E (Enterprise)
|
||||
// ============================================================================
|
||||
|
||||
/// End-to-end test for SQS trigger (Enterprise only).
|
||||
///
|
||||
/// Requires LocalStack with the test queue. Setup:
|
||||
/// ```bash
|
||||
/// ./tests/fixtures/start_sqs.sh
|
||||
/// ```
|
||||
///
|
||||
/// Run:
|
||||
/// ```bash
|
||||
/// AWS_ENDPOINT_URL=http://localhost:4566 \
|
||||
/// cargo test --test trigger_e2e test_sqs_e2e \
|
||||
/// --features sqs_trigger,enterprise,private -- --ignored --nocapture
|
||||
/// ```
|
||||
#[cfg(all(feature = "enterprise", feature = "private"))]
|
||||
#[ignore = "requires LocalStack SQS on localhost:4566"]
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_sqs_e2e(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
initialize_tracing().await;
|
||||
|
||||
// The SQS listener uses aws_config which respects AWS_ENDPOINT_URL for LocalStack.
|
||||
std::env::set_var("AWS_ENDPOINT_URL", "http://localhost:4566");
|
||||
|
||||
let script_path = "f/test/sqs_e2e_handler";
|
||||
insert_test_script(&db, script_path).await?;
|
||||
|
||||
insert_resource(
|
||||
&db,
|
||||
"u/test-user/aws_res",
|
||||
"aws",
|
||||
json!({
|
||||
"awsAccessKeyId": "test",
|
||||
"awsSecretAccessKey": "test",
|
||||
"region": "us-east-1"
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
r#"
|
||||
INSERT INTO sqs_trigger (
|
||||
path, queue_url, aws_resource_path, script_path,
|
||||
is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
|
||||
"#,
|
||||
"f/test/sqs_e2e_trigger",
|
||||
"http://localhost:4566/000000000000/windmill-e2e-test",
|
||||
"u/test-user/aws_res",
|
||||
script_path,
|
||||
false,
|
||||
"test-workspace",
|
||||
"test-user",
|
||||
"test@windmill.dev",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let _server = ApiServer::start_with_listeners(db.clone()).await?;
|
||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||
|
||||
// Send a message using aws-sdk-sqs
|
||||
let config = aws_config::defaults(aws_config::BehaviorVersion::latest())
|
||||
.endpoint_url("http://localhost:4566")
|
||||
.region(aws_config::Region::new("us-east-1"))
|
||||
.credentials_provider(aws_credential_types::Credentials::new(
|
||||
"test", "test", None, None, "test",
|
||||
))
|
||||
.load()
|
||||
.await;
|
||||
let sqs_client = aws_sdk_sqs::Client::new(&config);
|
||||
|
||||
sqs_client
|
||||
.send_message()
|
||||
.queue_url("http://localhost:4566/000000000000/windmill-e2e-test")
|
||||
.message_body("hello from sqs e2e test")
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
let job = poll_for_trigger_job(&db, script_path, "sqs", Duration::from_secs(30)).await?;
|
||||
assert!(job.args.is_some(), "job should have args");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// GCP Pub/Sub Trigger E2E (Enterprise)
|
||||
// ============================================================================
|
||||
|
||||
/// End-to-end test for GCP Pub/Sub trigger (Enterprise only).
|
||||
///
|
||||
/// Requires the GCP Pub/Sub emulator with test topic/subscription. Setup:
|
||||
/// ```bash
|
||||
/// ./tests/fixtures/start_gcp_pubsub.sh
|
||||
/// ```
|
||||
///
|
||||
/// Run:
|
||||
/// ```bash
|
||||
/// PUBSUB_EMULATOR_HOST=localhost:8085 \
|
||||
/// cargo test --test trigger_e2e test_gcp_e2e \
|
||||
/// --features gcp_trigger,enterprise,private -- --ignored --nocapture
|
||||
/// ```
|
||||
#[cfg(all(feature = "enterprise", feature = "private"))]
|
||||
#[ignore = "requires GCP Pub/Sub emulator on localhost:8085"]
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_gcp_e2e(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
initialize_tracing().await;
|
||||
|
||||
let script_path = "f/test/gcp_e2e_handler";
|
||||
insert_test_script(&db, script_path).await?;
|
||||
|
||||
// The GCP emulator doesn't require real credentials, but the resource
|
||||
// row must still exist for the listener to resolve it.
|
||||
// The private_key must use literal \n (backslash-n) as in real GCP service
|
||||
// account JSON files. The trigger code re-parses it through serde_json to
|
||||
// convert those escape sequences to actual newlines.
|
||||
insert_resource(
|
||||
&db,
|
||||
"u/test-user/gcp_res",
|
||||
"google",
|
||||
json!({
|
||||
"project_id": "test-project",
|
||||
"private_key_id": "test",
|
||||
"private_key": "-----BEGIN RSA PRIVATE KEY-----\\nMIIBogIBAAJBALRiMLAH\\n-----END RSA PRIVATE KEY-----\\n",
|
||||
"client_email": "test@test-project.iam.gserviceaccount.com",
|
||||
"auth_uri": "https://accounts.google.com/o/oauth2/auth",
|
||||
"token_uri": "https://oauth2.googleapis.com/token",
|
||||
"auth_provider_x509_cert_url": "https://www.googleapis.com/oauth2/v1/certs"
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO gcp_trigger (
|
||||
path, gcp_resource_path, topic_id, subscription_id,
|
||||
delivery_type, subscription_mode, script_path, is_flow,
|
||||
workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5::delivery_mode, $6::gcp_subscription_mode, $7, $8, $9, $10, $11)
|
||||
"#,
|
||||
)
|
||||
.bind("f/test/gcp_e2e_trigger")
|
||||
.bind("u/test-user/gcp_res")
|
||||
.bind("windmill-e2e-test")
|
||||
.bind("windmill-e2e-sub")
|
||||
.bind("pull")
|
||||
.bind("existing")
|
||||
.bind(script_path)
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let _server = ApiServer::start_with_listeners(db.clone()).await?;
|
||||
tokio::time::sleep(Duration::from_secs(3)).await;
|
||||
|
||||
// Publish a message to the emulator via HTTP
|
||||
let client = reqwest::Client::new();
|
||||
let emulator_host =
|
||||
std::env::var("PUBSUB_EMULATOR_HOST").unwrap_or_else(|_| "localhost:8085".to_string());
|
||||
// The google-cloud-pubsub crate uses "local-project" as the default project ID
|
||||
// when PUBSUB_EMULATOR_HOST is set, so we must publish to that project's topic.
|
||||
let publish_url = format!(
|
||||
"http://{}/v1/projects/local-project/topics/windmill-e2e-test:publish",
|
||||
emulator_host
|
||||
);
|
||||
|
||||
let message_data = base64::Engine::encode(
|
||||
&base64::engine::general_purpose::STANDARD,
|
||||
"hello from gcp e2e test",
|
||||
);
|
||||
client
|
||||
.post(&publish_url)
|
||||
.json(&json!({
|
||||
"messages": [{ "data": message_data }]
|
||||
}))
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
let job = poll_for_trigger_job(&db, script_path, "gcp", Duration::from_secs(30)).await?;
|
||||
assert!(job.args.is_some(), "job should have args");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -1272,3 +1272,431 @@ async fn test_schedule_insert_and_query(db: Pool<Postgres>) -> anyhow::Result<()
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// MQTT Trigger Tests (DB-level)
|
||||
// ============================================================================
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_mqtt_trigger_insert(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO mqtt_trigger (
|
||||
path, mqtt_resource_path, subscribe_topics, client_version,
|
||||
script_path, is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, ARRAY[$3::jsonb], $4::mqtt_client_version, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
)
|
||||
.bind("f/test/mqtt_trigger")
|
||||
.bind("u/admin/mqtt_resource")
|
||||
.bind(json!({"topic": "test/+", "qos": "qos1"}))
|
||||
.bind("v5")
|
||||
.bind("f/test/mqtt_handler")
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let trigger = sqlx::query!(
|
||||
r#"
|
||||
SELECT mqtt_resource_path, client_version AS "client_version: String",
|
||||
script_path, mode AS "mode: String"
|
||||
FROM mqtt_trigger
|
||||
WHERE workspace_id = $1 AND path = $2
|
||||
"#,
|
||||
"test-workspace",
|
||||
"f/test/mqtt_trigger",
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
assert_eq!(trigger.mqtt_resource_path, "u/admin/mqtt_resource");
|
||||
assert_eq!(trigger.client_version, "v5");
|
||||
assert_eq!(trigger.script_path, "f/test/mqtt_handler");
|
||||
assert_eq!(trigger.mode, "enabled");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_mqtt_trigger_update(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO mqtt_trigger (
|
||||
path, mqtt_resource_path, subscribe_topics, client_version,
|
||||
script_path, is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, ARRAY[$3::jsonb], $4::mqtt_client_version, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
)
|
||||
.bind("f/test/mqtt_trigger")
|
||||
.bind("u/admin/mqtt_resource")
|
||||
.bind(json!({"topic": "test/+", "qos": "qos1"}))
|
||||
.bind("v5")
|
||||
.bind("f/test/old_handler")
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"UPDATE mqtt_trigger SET script_path = $1 WHERE workspace_id = $2 AND path = $3",
|
||||
"f/test/new_handler",
|
||||
"test-workspace",
|
||||
"f/test/mqtt_trigger",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let trigger = sqlx::query!(
|
||||
"SELECT script_path FROM mqtt_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"test-workspace",
|
||||
"f/test/mqtt_trigger",
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
assert_eq!(trigger.script_path, "f/test/new_handler");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_mqtt_trigger_delete(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO mqtt_trigger (
|
||||
path, mqtt_resource_path, subscribe_topics, client_version,
|
||||
script_path, is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, ARRAY[$3::jsonb], $4::mqtt_client_version, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
)
|
||||
.bind("f/test/mqtt_trigger")
|
||||
.bind("u/admin/mqtt_resource")
|
||||
.bind(json!({"topic": "test/+", "qos": "qos1"}))
|
||||
.bind("v5")
|
||||
.bind("f/test/mqtt_handler")
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"DELETE FROM mqtt_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"test-workspace",
|
||||
"f/test/mqtt_trigger",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let count = sqlx::query_scalar!(
|
||||
"SELECT COUNT(*) FROM mqtt_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"test-workspace",
|
||||
"f/test/mqtt_trigger",
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
assert_eq!(count, Some(0));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// GCP Trigger Tests (DB-level)
|
||||
// ============================================================================
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_gcp_trigger_insert_pull(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO gcp_trigger (
|
||||
path, gcp_resource_path, topic_id, subscription_id,
|
||||
delivery_type, subscription_mode, script_path, is_flow,
|
||||
workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5::delivery_mode, $6::gcp_subscription_mode, $7, $8, $9, $10, $11)
|
||||
"#,
|
||||
)
|
||||
.bind("f/test/gcp_trigger_pull")
|
||||
.bind("u/admin/gcp_resource")
|
||||
.bind("my-topic")
|
||||
.bind("my-subscription")
|
||||
.bind("pull")
|
||||
.bind("create_update")
|
||||
.bind("f/test/gcp_handler")
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let trigger = sqlx::query!(
|
||||
r#"
|
||||
SELECT gcp_resource_path, topic_id, subscription_id,
|
||||
delivery_type AS "delivery_type: String",
|
||||
subscription_mode AS "subscription_mode: String",
|
||||
mode AS "mode: String"
|
||||
FROM gcp_trigger
|
||||
WHERE workspace_id = $1 AND path = $2
|
||||
"#,
|
||||
"test-workspace",
|
||||
"f/test/gcp_trigger_pull",
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
assert_eq!(trigger.gcp_resource_path, "u/admin/gcp_resource");
|
||||
assert_eq!(trigger.topic_id, "my-topic");
|
||||
assert_eq!(trigger.subscription_id, "my-subscription");
|
||||
assert_eq!(trigger.delivery_type, "pull");
|
||||
assert_eq!(trigger.subscription_mode, "create_update");
|
||||
assert_eq!(trigger.mode, "enabled");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_gcp_trigger_insert_push(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO gcp_trigger (
|
||||
path, gcp_resource_path, topic_id, subscription_id,
|
||||
delivery_type, delivery_config, subscription_mode,
|
||||
script_path, is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5::delivery_mode, $6::jsonb, $7::gcp_subscription_mode, $8, $9, $10, $11, $12)
|
||||
"#,
|
||||
)
|
||||
.bind("f/test/gcp_trigger_push")
|
||||
.bind("u/admin/gcp_resource")
|
||||
.bind("my-topic")
|
||||
.bind("my-push-subscription")
|
||||
.bind("push")
|
||||
.bind(json!({"endpoint": "https://example.com/push"}))
|
||||
.bind("create_update")
|
||||
.bind("f/test/gcp_handler")
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let trigger = sqlx::query!(
|
||||
r#"
|
||||
SELECT delivery_type AS "delivery_type: String",
|
||||
delivery_config
|
||||
FROM gcp_trigger
|
||||
WHERE workspace_id = $1 AND path = $2
|
||||
"#,
|
||||
"test-workspace",
|
||||
"f/test/gcp_trigger_push",
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
assert_eq!(trigger.delivery_type, "push");
|
||||
assert!(trigger.delivery_config.is_some());
|
||||
assert_eq!(
|
||||
trigger.delivery_config.unwrap()["endpoint"],
|
||||
"https://example.com/push"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_gcp_trigger_unique_constraint(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
let insert_query = r#"
|
||||
INSERT INTO gcp_trigger (
|
||||
path, gcp_resource_path, topic_id, subscription_id,
|
||||
delivery_type, subscription_mode, script_path, is_flow,
|
||||
workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5::delivery_mode, $6::gcp_subscription_mode, $7, $8, $9, $10, $11)
|
||||
"#;
|
||||
|
||||
sqlx::query(insert_query)
|
||||
.bind("f/test/gcp_trigger_1")
|
||||
.bind("u/admin/gcp_resource")
|
||||
.bind("my-topic")
|
||||
.bind("shared-subscription")
|
||||
.bind("pull")
|
||||
.bind("create_update")
|
||||
.bind("f/test/gcp_handler")
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
// Inserting a second trigger with same (subscription_id, gcp_resource_path, workspace_id) should fail
|
||||
let result = sqlx::query(insert_query)
|
||||
.bind("f/test/gcp_trigger_2")
|
||||
.bind("u/admin/gcp_resource")
|
||||
.bind("my-topic")
|
||||
.bind("shared-subscription")
|
||||
.bind("pull")
|
||||
.bind("create_update")
|
||||
.bind("f/test/gcp_handler_2")
|
||||
.bind(false)
|
||||
.bind("test-workspace")
|
||||
.bind("test-user")
|
||||
.bind("test@windmill.dev")
|
||||
.execute(&db)
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
result.is_err(),
|
||||
"should fail due to unique constraint on (subscription_id, gcp_resource_path, workspace_id)"
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Email Trigger Tests (DB-level)
|
||||
// ============================================================================
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_email_trigger_insert(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
sqlx::query!(
|
||||
r#"
|
||||
INSERT INTO email_trigger (
|
||||
path, local_part, workspaced_local_part, script_path,
|
||||
is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
|
||||
"#,
|
||||
"f/test/email_trigger",
|
||||
"support",
|
||||
true,
|
||||
"f/test/email_handler",
|
||||
false,
|
||||
"test-workspace",
|
||||
"test-user",
|
||||
"test@windmill.dev",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let trigger = sqlx::query!(
|
||||
r#"
|
||||
SELECT local_part, workspaced_local_part, script_path,
|
||||
mode AS "mode: String"
|
||||
FROM email_trigger
|
||||
WHERE workspace_id = $1 AND path = $2
|
||||
"#,
|
||||
"test-workspace",
|
||||
"f/test/email_trigger",
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
assert_eq!(trigger.local_part, "support");
|
||||
assert_eq!(trigger.workspaced_local_part, true);
|
||||
assert_eq!(trigger.script_path, "f/test/email_handler");
|
||||
assert_eq!(trigger.mode, "enabled");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_email_trigger_update(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
sqlx::query!(
|
||||
r#"
|
||||
INSERT INTO email_trigger (
|
||||
path, local_part, workspaced_local_part, script_path,
|
||||
is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
|
||||
"#,
|
||||
"f/test/email_trigger",
|
||||
"support",
|
||||
true,
|
||||
"f/test/old_handler",
|
||||
false,
|
||||
"test-workspace",
|
||||
"test-user",
|
||||
"test@windmill.dev",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"UPDATE email_trigger SET script_path = $1, local_part = $2 WHERE workspace_id = $3 AND path = $4",
|
||||
"f/test/new_handler",
|
||||
"billing",
|
||||
"test-workspace",
|
||||
"f/test/email_trigger",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let trigger = sqlx::query!(
|
||||
"SELECT script_path, local_part FROM email_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"test-workspace",
|
||||
"f/test/email_trigger",
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
assert_eq!(trigger.script_path, "f/test/new_handler");
|
||||
assert_eq!(trigger.local_part, "billing");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_email_trigger_delete(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
sqlx::query!(
|
||||
r#"
|
||||
INSERT INTO email_trigger (
|
||||
path, local_part, workspaced_local_part, script_path,
|
||||
is_flow, workspace_id, edited_by, email
|
||||
)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
|
||||
"#,
|
||||
"f/test/email_trigger",
|
||||
"support",
|
||||
true,
|
||||
"f/test/email_handler",
|
||||
false,
|
||||
"test-workspace",
|
||||
"test-user",
|
||||
"test@windmill.dev",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"DELETE FROM email_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"test-workspace",
|
||||
"f/test/email_trigger",
|
||||
)
|
||||
.execute(&db)
|
||||
.await?;
|
||||
|
||||
let count = sqlx::query_scalar!(
|
||||
"SELECT COUNT(*) FROM email_trigger WHERE workspace_id = $1 AND path = $2",
|
||||
"test-workspace",
|
||||
"f/test/email_trigger",
|
||||
)
|
||||
.fetch_one(&db)
|
||||
.await?;
|
||||
|
||||
assert_eq!(count, Some(0));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -259,7 +259,7 @@ async fn test_user_endpoints(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
let resp = authed(client().post(format!(
|
||||
"{global_base}/rename/newglobaluser@windmill.dev"
|
||||
)))
|
||||
.json(&json!({"new_username": "renamed-user"}))
|
||||
.json(&json!({"new_username": "renamed_user"}))
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
@@ -296,7 +296,7 @@ async fn test_user_endpoints(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
resp.status() == 401 || resp.status() == 500,
|
||||
resp.status() == 400 || resp.status() == 401 || resp.status() == 500,
|
||||
"login: unexpected status {}",
|
||||
resp.status()
|
||||
);
|
||||
|
||||
@@ -3343,7 +3343,6 @@ async fn test_duckdb_ffi(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
/// Test that flow substeps with tags that are not available for the workspace fail.
|
||||
/// This validates that `check_tag_available_for_workspace_internal` is properly called
|
||||
/// when pushing jobs from worker_flow.
|
||||
#[cfg(feature = "deno_core")]
|
||||
#[sqlx::test(fixtures("base"))]
|
||||
async fn test_flow_substep_tag_availability_check(db: Pool<Postgres>) -> anyhow::Result<()> {
|
||||
use windmill_common::worker::{
|
||||
@@ -3385,6 +3384,7 @@ async fn test_flow_substep_tag_availability_check(db: Pool<Postgres>) -> anyhow:
|
||||
|
||||
let result =
|
||||
RunJob::from(JobPayload::RawFlow { value: flow.clone(), path: None, restarted_from: None })
|
||||
.email("test2@windmill.dev")
|
||||
.run_until_complete(&db, false, server.addr.port())
|
||||
.await;
|
||||
|
||||
|
||||
@@ -447,8 +447,9 @@ pub async fn fetch_api_authed_from_permissioned_as(
|
||||
}
|
||||
_ => {
|
||||
tracing::debug!("API authed cache miss for user {}", email);
|
||||
|
||||
let authed =
|
||||
fetch_authed_from_permissioned_as(permissioned_as, email.clone(), w_id, db).await?;
|
||||
fetch_authed_from_permissioned_as(&permissioned_as, &email, w_id, db).await?;
|
||||
|
||||
let api_authed = ApiAuthed {
|
||||
username: authed.username,
|
||||
|
||||
@@ -1496,7 +1496,7 @@ async fn toggle_workspace_error_handler(
|
||||
}
|
||||
None => {
|
||||
tx.commit().await?;
|
||||
Err(Error::ExecutionErr(
|
||||
Err(Error::BadRequest(
|
||||
"Workspace error handler needs to be defined".to_string(),
|
||||
))
|
||||
}
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
openapi: "3.0.3"
|
||||
|
||||
info:
|
||||
version: 1.628.3
|
||||
version: 1.629.1
|
||||
title: Windmill API
|
||||
|
||||
contact:
|
||||
@@ -1310,6 +1310,20 @@ paths:
|
||||
schema:
|
||||
type: string
|
||||
|
||||
/settings/get_stats:
|
||||
get:
|
||||
summary: get encrypted telemetry stats (EE only)
|
||||
operationId: getStats
|
||||
tags:
|
||||
- setting
|
||||
responses:
|
||||
"200":
|
||||
description: base64-encoded encrypted telemetry blob
|
||||
content:
|
||||
text/plain:
|
||||
schema:
|
||||
type: string
|
||||
|
||||
/settings/latest_key_renewal_attempt:
|
||||
get:
|
||||
summary: get latest key renewal attempt
|
||||
|
||||
@@ -68,7 +68,7 @@ lazy_static::lazy_static! {
|
||||
).replace("CREATE INDEX", "CREATE INDEX CONCURRENTLY").replace("DROP INDEX", "DROP INDEX CONCURRENTLY")),
|
||||
(20260207000002, include_str!(
|
||||
"../../migrations/20260207000002_concurrent_indexes_v2_job_completed.up.sql"
|
||||
).replace("CREATE INDEX", "CREATE INDEX CONCURRENTLY").replace("DROP INDEX", "DROP INDEX CONCURRENTLY")),
|
||||
).replace("CREATE INDEX", "CREATE INDEX CONCURRENTLY").replace("DROP INDEX", "DROP INDEX CONCURRENTLY").replace("DROP INDEX CONCURRENTLY IF EXISTS labeled_jobs_on_jobs;", "")),
|
||||
(20260207000003, include_str!(
|
||||
"../../migrations/20260207000003_concurrent_indexes_v2_job_queue.up.sql"
|
||||
).replace("CREATE INDEX", "CREATE INDEX CONCURRENTLY").replace("DROP INDEX", "DROP INDEX CONCURRENTLY")),
|
||||
|
||||
@@ -58,6 +58,7 @@ pub fn global_service() -> Router {
|
||||
.route("/test_smtp", post(test_email))
|
||||
.route("/test_license_key", post(test_license_key))
|
||||
.route("/send_stats", post(send_stats))
|
||||
.route("/get_stats", get(get_stats))
|
||||
.route(
|
||||
"/latest_key_renewal_attempt",
|
||||
get(get_latest_key_renewal_attempt),
|
||||
@@ -434,6 +435,25 @@ pub async fn send_stats(Extension(db): Extension<DB>, authed: ApiAuthed) -> Resu
|
||||
Ok("Sent stats".to_string())
|
||||
}
|
||||
|
||||
#[cfg(feature = "enterprise")]
|
||||
pub async fn get_stats(Extension(db): Extension<DB>, authed: ApiAuthed) -> Result<String> {
|
||||
require_super_admin(&db, &authed.email).await?;
|
||||
let stats = windmill_common::stats_oss::get_stats_payload(
|
||||
&db,
|
||||
&windmill_common::stats_oss::SendStatsReason::Manual,
|
||||
)
|
||||
.await?;
|
||||
let encrypted = windmill_common::stats_oss::encrypt_stats(&stats)?;
|
||||
Ok(encrypted)
|
||||
}
|
||||
|
||||
#[cfg(not(feature = "enterprise"))]
|
||||
pub async fn get_stats() -> Result<String> {
|
||||
Err(error::Error::BadRequest(
|
||||
"Downloading telemetry is only available on enterprise edition".to_string(),
|
||||
))
|
||||
}
|
||||
|
||||
#[derive(serde::Serialize)]
|
||||
pub struct KeyRenewalAttempt {
|
||||
result: String,
|
||||
|
||||
@@ -244,20 +244,26 @@ pub fn permissioned_as_to_username(permissioned_as: &str) -> String {
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn fetch_authed_from_permissioned_as(
|
||||
permissioned_as: String,
|
||||
email: String,
|
||||
w_id: &str,
|
||||
db: &DB,
|
||||
) -> Result<Authed> {
|
||||
let mut conn = db
|
||||
.acquire()
|
||||
.await
|
||||
.map_err(|e| Error::internal_err(format!("acquiring connection: {e:#}")))?;
|
||||
fetch_authed_from_permissioned_as_conn(&permissioned_as, &email, w_id, &mut conn).await
|
||||
pub fn fetch_authed_from_permissioned_as<'a, A>(
|
||||
permissioned_as: &'a str,
|
||||
email: &'a str,
|
||||
w_id: &'a str,
|
||||
db: A,
|
||||
) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<Authed>> + Send + 'a>>
|
||||
where
|
||||
A: sqlx::Acquire<'a, Database = sqlx::Postgres> + Send + 'a,
|
||||
{
|
||||
Box::pin(async move {
|
||||
let mut conn = db
|
||||
.acquire()
|
||||
.await
|
||||
.map_err(|e| Error::internal_err(format!("acquiring connection: {e:#}")))?;
|
||||
|
||||
fetch_authed_from_permissioned_as_inner(permissioned_as, email, w_id, &mut *conn).await
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn fetch_authed_from_permissioned_as_conn(
|
||||
async fn fetch_authed_from_permissioned_as_inner(
|
||||
permissioned_as: &str,
|
||||
email: &str,
|
||||
w_id: &str,
|
||||
@@ -418,7 +424,7 @@ pub async fn create_token_for_owner(
|
||||
Ok(Some(jp)) => jp.into(),
|
||||
_ => {
|
||||
tracing::warn!("Could not get permissions for job {job_id} from job_perms table, getting permissions directly...");
|
||||
fetch_authed_from_permissioned_as(owner.to_string(), email.to_string(), w_id, db)
|
||||
fetch_authed_from_permissioned_as(owner, email, w_id, db)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Error::internal_err(format!(
|
||||
|
||||
@@ -5,6 +5,31 @@ use crate::{
|
||||
use serde::Serialize;
|
||||
use tokio::time::Instant;
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct PoolStats {
|
||||
pub peak_active_conns: u32,
|
||||
pub pool_saturation_histogram: Vec<u32>,
|
||||
}
|
||||
|
||||
impl PoolStats {
|
||||
pub fn new(pool_size: u32) -> Self {
|
||||
PoolStats {
|
||||
peak_active_conns: 0,
|
||||
pool_saturation_histogram: vec![0; pool_size as usize + 1],
|
||||
}
|
||||
}
|
||||
|
||||
pub fn sample(&mut self, pool_size: u32, num_idle: u32) {
|
||||
let active = pool_size.saturating_sub(num_idle);
|
||||
if active > self.peak_active_conns {
|
||||
self.peak_active_conns = active;
|
||||
}
|
||||
if let Some(bucket) = self.pool_saturation_histogram.get_mut(active as usize) {
|
||||
*bucket += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct BenchmarkInfo {
|
||||
#[serde(skip)]
|
||||
@@ -14,6 +39,7 @@ pub struct BenchmarkInfo {
|
||||
timings: Vec<BenchmarkIter>,
|
||||
pub iter_durations: Vec<u64>,
|
||||
pub total_duration: Option<u64>,
|
||||
pub pool_stats: Option<PoolStats>,
|
||||
}
|
||||
|
||||
impl BenchmarkInfo {
|
||||
@@ -24,6 +50,17 @@ impl BenchmarkInfo {
|
||||
start: Instant::now(),
|
||||
iter_durations: vec![],
|
||||
total_duration: None,
|
||||
pool_stats: None,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn init_pool_stats(&mut self, pool_size: u32) {
|
||||
self.pool_stats = Some(PoolStats::new(pool_size));
|
||||
}
|
||||
|
||||
pub fn sample_pool(&mut self, pool_size: u32, num_idle: u32) {
|
||||
if let Some(stats) = self.pool_stats.as_mut() {
|
||||
stats.sample(pool_size, num_idle);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -40,8 +77,11 @@ impl BenchmarkInfo {
|
||||
let total_duration = self.start.elapsed().as_millis() as u64;
|
||||
self.total_duration = Some(total_duration as u64);
|
||||
|
||||
let pool_info = self.pool_stats.as_ref().map_or(String::new(), |ps| {
|
||||
format!(", peak active conns: {}", ps.peak_active_conns)
|
||||
});
|
||||
println!(
|
||||
"Writing benchmark {path}, duration of benchmark: {total_duration}ms and RPS: {}",
|
||||
"Writing benchmark {path}, duration of benchmark: {total_duration}ms and RPS: {}{pool_info}",
|
||||
self.iters as f64 / total_duration as f64 * 1000.0
|
||||
);
|
||||
write_file(TMP_DIR, path, &serde_json::to_string(&self).unwrap()).expect("write profiling");
|
||||
@@ -220,6 +260,7 @@ pub async fn benchmark_init(benchmark_jobs: i32, db: &DB) {
|
||||
.await
|
||||
.unwrap_or_else(|_e| panic!("failed to insert parallelflow jobs (4)"));
|
||||
}
|
||||
"none" => {}
|
||||
_ => {
|
||||
let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9 FROM generate_series(1, $10)) RETURNING id",
|
||||
None::<i64>,
|
||||
|
||||
@@ -50,3 +50,19 @@ pub async fn get_user_usage<'c, E: sqlx::Executor<'c, Database = Postgres>>(
|
||||
let usage = ActiveUserUsage { author_count: None, operator_count: None };
|
||||
Ok(usage)
|
||||
}
|
||||
|
||||
#[cfg(not(feature = "private"))]
|
||||
#[derive(serde::Serialize)]
|
||||
pub struct Stats {}
|
||||
|
||||
#[cfg(not(feature = "private"))]
|
||||
pub async fn get_stats_payload(_db: &DB, _reason: &SendStatsReason) -> Result<Stats> {
|
||||
// stats details are closed source
|
||||
Ok(Stats {})
|
||||
}
|
||||
|
||||
#[cfg(not(feature = "private"))]
|
||||
pub fn encrypt_stats(_stats: &Stats) -> Result<String> {
|
||||
// stats details are closed source
|
||||
Ok(String::new())
|
||||
}
|
||||
|
||||
@@ -1045,16 +1045,17 @@ pub async fn extract_tar(tar: bytes::Bytes, folder: &str) -> error::Result<()> {
|
||||
}
|
||||
#[cfg(all(feature = "enterprise", feature = "parquet"))]
|
||||
fn write_binary_file(main_path: &str, byts: &mut bytes::Bytes) -> error::Result<()> {
|
||||
use std::fs::{File, Permissions};
|
||||
use std::fs::File;
|
||||
use std::io::Write;
|
||||
|
||||
#[cfg(unix)]
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
|
||||
let mut file = File::create(main_path)?;
|
||||
file.write_all(byts)?;
|
||||
#[cfg(unix)]
|
||||
file.set_permissions(Permissions::from_mode(0o755))?;
|
||||
{
|
||||
use std::fs::Permissions;
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
file.set_permissions(Permissions::from_mode(0o755))?;
|
||||
}
|
||||
file.flush()?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -53,7 +53,7 @@ use windmill_common::utils::{calculate_hash, configure_client, now_from_db};
|
||||
use windmill_common::worker::{Connection, SCRIPT_TOKEN_EXPIRY};
|
||||
|
||||
use windmill_common::{
|
||||
auth::{fetch_authed_from_permissioned_as, permissioned_as_to_username},
|
||||
auth::permissioned_as_to_username,
|
||||
cache::{self, FlowData},
|
||||
db::{Authed, UserDB},
|
||||
error::{self, Error},
|
||||
@@ -782,7 +782,6 @@ pub async fn add_completed_job_error(
|
||||
flow_is_done,
|
||||
duration,
|
||||
false,
|
||||
false,
|
||||
)
|
||||
.warn_after_seconds(10)
|
||||
.await?;
|
||||
@@ -818,7 +817,6 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
canceled_by: Option<CanceledBy>,
|
||||
flow_is_done: bool,
|
||||
duration: Option<i64>,
|
||||
has_stream: bool,
|
||||
from_cache: bool,
|
||||
) -> Result<(Uuid, i64), Error> {
|
||||
// tracing::error!("Start");
|
||||
@@ -844,7 +842,6 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
&canceled_by,
|
||||
flow_is_done,
|
||||
duration,
|
||||
has_stream,
|
||||
from_cache,
|
||||
)
|
||||
.warn_after_seconds(10)
|
||||
@@ -904,7 +901,6 @@ async fn commit_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
canceled_by: &Option<CanceledBy>,
|
||||
flow_is_done: bool,
|
||||
duration: Option<i64>,
|
||||
has_stream: bool,
|
||||
from_cache: bool,
|
||||
) -> windmill_common::error::Result<(Option<Uuid>, i64, bool)> {
|
||||
// let start = std::time::Instant::now();
|
||||
@@ -1184,7 +1180,19 @@ async fn commit_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
.concurrent_limit
|
||||
.is_some()
|
||||
{
|
||||
let concurrency_key = concurrency_key(db, &completed_job.id).await?;
|
||||
let concurrency_key = sqlx::query_scalar!(
|
||||
"SELECT key FROM concurrency_key WHERE job_id = $1",
|
||||
&completed_job.id
|
||||
)
|
||||
.fetch_optional(&mut *tx)
|
||||
.warn_after_seconds(10)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Error::internal_err(format!(
|
||||
"Could not get concurrency key for job {}: {e:#}",
|
||||
completed_job.id
|
||||
))
|
||||
})?;
|
||||
if *DISABLE_CONCURRENCY_LIMIT || concurrency_key.is_none() {
|
||||
tracing::warn!("Concurrency limit is disabled, skipping");
|
||||
} else {
|
||||
@@ -1221,18 +1229,6 @@ async fn commit_completed_job<T: Serialize + Send + Sync + ValidableJson>(
|
||||
tracing::debug!("decremented concurrency counter");
|
||||
}
|
||||
|
||||
sqlx::query!("DELETE FROM job_perms WHERE job_id = $1", job_id)
|
||||
.execute(&mut *tx)
|
||||
.warn_after_seconds(10)
|
||||
.await?;
|
||||
|
||||
if !success || has_stream {
|
||||
sqlx::query!("DELETE FROM job_result_stream_v2 WHERE job_id = $1", job_id)
|
||||
.execute(&mut *tx)
|
||||
.warn_after_seconds(10)
|
||||
.await?;
|
||||
}
|
||||
|
||||
tx.commit().warn_after_seconds(10).await?;
|
||||
|
||||
tracing::info!(
|
||||
@@ -1827,7 +1823,7 @@ pub async fn try_schedule_next_job<'c>(
|
||||
&job.workspace_id
|
||||
);
|
||||
|
||||
let schedule_authed = windmill_common::auth::fetch_authed_from_permissioned_as_conn(
|
||||
let schedule_authed = windmill_common::auth::fetch_authed_from_permissioned_as(
|
||||
&windmill_common::users::username_to_permissioned_as(&schedule.edited_by),
|
||||
&schedule.email,
|
||||
&job.workspace_id,
|
||||
@@ -5386,7 +5382,7 @@ async fn push_inner<'c, 'd>(
|
||||
&runnable_path,
|
||||
job_kind,
|
||||
concurrency_settings.concurrency_key.clone(),
|
||||
&mut tx,
|
||||
&mut *tx,
|
||||
job_id,
|
||||
)
|
||||
.await?;
|
||||
@@ -5418,11 +5414,11 @@ async fn push_inner<'c, 'd>(
|
||||
if authed.is_some() {
|
||||
tracing::warn!("Authed passed to push is not the same as permissioned_as, refetching direclty permissions for job {job_id}...")
|
||||
}
|
||||
fetch_authed_from_permissioned_as(
|
||||
permissioned_as.clone(),
|
||||
email.to_string(),
|
||||
windmill_common::auth::fetch_authed_from_permissioned_as(
|
||||
&permissioned_as,
|
||||
email,
|
||||
workspace_id,
|
||||
_db,
|
||||
&mut *tx,
|
||||
)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
@@ -5709,7 +5705,7 @@ pub async fn insert_concurrency_key<'d, 'c>(
|
||||
script_path: &Option<String>,
|
||||
job_kind: JobKind,
|
||||
custom_concurrency_key: Option<String>,
|
||||
tx: &mut Transaction<'c, Postgres>,
|
||||
db: impl PgExecutor<'c>,
|
||||
job_id: Uuid,
|
||||
) -> Result<(), Error> {
|
||||
let concurrency_key = custom_concurrency_key
|
||||
@@ -5737,7 +5733,7 @@ pub async fn insert_concurrency_key<'d, 'c>(
|
||||
));
|
||||
sqlx::query!(
|
||||
"WITH inserted_concurrency_counter AS (
|
||||
INSERT INTO concurrency_counter (concurrency_id, job_uuids)
|
||||
INSERT INTO concurrency_counter (concurrency_id, job_uuids)
|
||||
VALUES ($1, '{}'::jsonb)
|
||||
ON CONFLICT DO NOTHING
|
||||
)
|
||||
@@ -5745,7 +5741,7 @@ pub async fn insert_concurrency_key<'d, 'c>(
|
||||
concurrency_key,
|
||||
job_id,
|
||||
)
|
||||
.execute(&mut **tx)
|
||||
.execute(db)
|
||||
.warn_after_seconds(3)
|
||||
.await
|
||||
.map_err(|e| Error::internal_err(format!("Could not insert concurrency_key={concurrency_key} for job_id={job_id} script_path={script_path:?} workspace_id={workspace_id}: {e:#}")))?;
|
||||
|
||||
@@ -468,7 +468,7 @@ pub async fn push_scheduled_job<'c>(
|
||||
let push_authed = match push_authed {
|
||||
Some(a) => Some(a),
|
||||
None => {
|
||||
obo_authed = windmill_common::auth::fetch_authed_from_permissioned_as_conn(
|
||||
obo_authed = windmill_common::auth::fetch_authed_from_permissioned_as(
|
||||
&permissioned_as,
|
||||
email,
|
||||
&schedule.workspace_id,
|
||||
|
||||
@@ -35,3 +35,4 @@ anyhow.workspace = true
|
||||
chrono.workspace = true
|
||||
backon.workspace = true
|
||||
thiserror.workspace = true
|
||||
aws-smithy-types.workspace = true
|
||||
@@ -460,6 +460,11 @@ impl OpenAPISchema {
|
||||
Typ::Bool => Self::from_str("boolean"),
|
||||
Typ::Bytes => Self::from_str("string"),
|
||||
Typ::Datetime => Self::datetime(),
|
||||
Typ::Date => Self {
|
||||
r#type: Some(SchemaType::Single("string".to_string())),
|
||||
format: Some("date".to_string()),
|
||||
..Default::default()
|
||||
},
|
||||
Typ::Resource(_) => Self::from_str("string"),
|
||||
Typ::Email => Self::from_str("string"),
|
||||
Typ::Sql => Self::from_str("string"),
|
||||
|
||||
@@ -752,10 +752,12 @@ pub async fn get_git_ssh_cmd(
|
||||
})?;
|
||||
content.push_str("\n");
|
||||
|
||||
let file = write_file(job_dir, &id_file_name, &content)?;
|
||||
#[cfg(not(unix))]
|
||||
let _ = write_file(job_dir, &id_file_name, &content)?;
|
||||
|
||||
#[cfg(unix)]
|
||||
{
|
||||
let file = write_file(job_dir, &id_file_name, &content)?;
|
||||
let perm = std::os::unix::fs::PermissionsExt::from_mode(0o600);
|
||||
file.set_permissions(perm)?;
|
||||
}
|
||||
@@ -1218,10 +1220,10 @@ fi
|
||||
ANSIBLE_PLAYBOOK_PATH.as_str()
|
||||
);
|
||||
|
||||
let file = write_file(job_dir, "wrapper.sh", &wrapper)?;
|
||||
let _file = write_file(job_dir, "wrapper.sh", &wrapper)?;
|
||||
|
||||
#[cfg(unix)]
|
||||
file.metadata()?.permissions().set_mode(0o777);
|
||||
_file.metadata()?.permissions().set_mode(0o777);
|
||||
// let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str());
|
||||
let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str());
|
||||
nsjail_cmd
|
||||
|
||||
@@ -1016,6 +1016,21 @@ async fn prepare_wrapper(
|
||||
kwargs[\"{name}\"] = datetime.fromisoformat(kwargs[\"{name}\"])\n",
|
||||
)
|
||||
}
|
||||
windmill_parser::Typ::Date => {
|
||||
let name = &x.name;
|
||||
format!(
|
||||
"if \"{name}\" in kwargs and kwargs[\"{name}\"] is not None:\n \
|
||||
try:\n \
|
||||
kwargs[\"{name}\"] = date.fromisoformat(kwargs[\"{name}\"])\n \
|
||||
except ValueError:\n \
|
||||
for _fmt in (\"%d-%m-%Y\", \"%m/%d/%Y\", \"%d/%m/%Y\", \"%Y/%m/%d\"):\n \
|
||||
try:\n \
|
||||
kwargs[\"{name}\"] = datetime.strptime(kwargs[\"{name}\"], _fmt).date()\n \
|
||||
break\n \
|
||||
except ValueError:\n \
|
||||
continue\n",
|
||||
)
|
||||
}
|
||||
_ => "".to_string(),
|
||||
})
|
||||
.collect::<Vec<String>>()
|
||||
@@ -1035,14 +1050,19 @@ async fn prepare_wrapper(
|
||||
} else {
|
||||
""
|
||||
};
|
||||
let import_datetime = if init_sig
|
||||
let has_datetime = init_sig
|
||||
.args
|
||||
.iter()
|
||||
.any(|x| x.typ == windmill_parser::Typ::Datetime)
|
||||
{
|
||||
"from datetime import datetime"
|
||||
} else {
|
||||
""
|
||||
.any(|x| x.typ == windmill_parser::Typ::Datetime);
|
||||
let has_date = init_sig
|
||||
.args
|
||||
.iter()
|
||||
.any(|x| x.typ == windmill_parser::Typ::Date);
|
||||
let import_datetime = match (has_datetime, has_date) {
|
||||
(true, true) => "from datetime import datetime, date",
|
||||
(true, false) => "from datetime import datetime",
|
||||
(false, true) => "from datetime import datetime, date",
|
||||
(false, false) => "",
|
||||
};
|
||||
let spread = if sig.star_kwargs {
|
||||
"args = kwargs".to_string()
|
||||
|
||||
@@ -21,10 +21,12 @@ use windmill_queue::append_logs;
|
||||
use crate::{
|
||||
common::{start_child_process, OccupancyMetrics},
|
||||
handle_child::handle_child,
|
||||
python_executor::{INDEX_CERT, NATIVE_CERT, PYTHON_PATH, UV_PATH},
|
||||
python_executor::{INDEX_CERT, NATIVE_CERT, PYTHON_PATH},
|
||||
HOME_ENV, INSTANCE_PYTHON_VERSION, PATH_ENV, PROXY_ENVS, PY_INSTALL_DIR, UV_CACHE_DIR,
|
||||
WIN_ENVS,
|
||||
};
|
||||
#[cfg(unix)]
|
||||
use crate::python_executor::UV_PATH;
|
||||
|
||||
impl From<PyV> for PyVAlias {
|
||||
fn from(value: PyV) -> Self {
|
||||
|
||||
@@ -586,7 +586,6 @@ pub async fn process_completed_job(
|
||||
duration,
|
||||
result_columns,
|
||||
preprocessed_args,
|
||||
has_stream,
|
||||
from_cache,
|
||||
flow_runners,
|
||||
done_tx,
|
||||
@@ -655,7 +654,6 @@ pub async fn process_completed_job(
|
||||
canceled_by.clone(),
|
||||
false,
|
||||
duration,
|
||||
has_stream.unwrap_or(false),
|
||||
from_cache.unwrap_or(false),
|
||||
)
|
||||
.await?;
|
||||
|
||||
@@ -34,7 +34,7 @@ fn make_rules_for_arg_typ(typ: &Typ) -> Vec<SchemaValidationRule> {
|
||||
rules.push(SchemaValidationRule::IsString);
|
||||
rules.push(SchemaValidationRule::IsBytes);
|
||||
}
|
||||
Typ::Datetime => {
|
||||
Typ::Datetime | Typ::Date => {
|
||||
rules.push(SchemaValidationRule::IsString);
|
||||
rules.push(SchemaValidationRule::IsDatetime);
|
||||
}
|
||||
|
||||
@@ -12,7 +12,6 @@
|
||||
use anyhow::anyhow;
|
||||
use futures::TryFutureExt;
|
||||
use tokio::sync::Mutex;
|
||||
use tokio::time::sleep;
|
||||
use tokio::time::timeout;
|
||||
use windmill_common::client::AuthedClient;
|
||||
use windmill_common::jobs::WorkerInternalServerInlineUtils;
|
||||
@@ -1618,6 +1617,11 @@ pub async fn run_worker(
|
||||
#[cfg(feature = "benchmark")]
|
||||
let mut infos = BenchmarkInfo::new();
|
||||
|
||||
#[cfg(feature = "benchmark")]
|
||||
if let Some(db) = conn.as_sql() {
|
||||
infos.init_pool_stats(db.size());
|
||||
}
|
||||
|
||||
let vacuum_shift = rand::rng().random_range(0..VACUUM_PERIOD);
|
||||
|
||||
IS_READY.store(true, Ordering::Relaxed);
|
||||
@@ -2067,6 +2071,9 @@ pub async fn run_worker(
|
||||
#[cfg(feature = "benchmark")]
|
||||
{
|
||||
add_time!(bench, "sent to dedicated worker");
|
||||
if let Some(db) = conn.as_sql() {
|
||||
infos.sample_pool(db.size(), db.num_idle() as u32);
|
||||
}
|
||||
infos.add_iter(bench, true);
|
||||
}
|
||||
|
||||
@@ -2137,6 +2144,9 @@ pub async fn run_worker(
|
||||
#[cfg(feature = "benchmark")]
|
||||
{
|
||||
add_time!(bench, "sent to flow runner");
|
||||
if let Some(db) = conn.as_sql() {
|
||||
infos.sample_pool(db.size(), db.num_idle() as u32);
|
||||
}
|
||||
infos.add_iter(bench, true);
|
||||
}
|
||||
|
||||
@@ -2437,6 +2447,9 @@ pub async fn run_worker(
|
||||
{
|
||||
if started {
|
||||
add_time!(bench, "job processed");
|
||||
if let Some(db) = conn.as_sql() {
|
||||
infos.sample_pool(db.size(), db.num_idle() as u32);
|
||||
}
|
||||
infos.add_iter(bench, true);
|
||||
}
|
||||
}
|
||||
@@ -2465,6 +2478,9 @@ pub async fn run_worker(
|
||||
#[cfg(feature = "benchmark")]
|
||||
{
|
||||
add_time!(bench, "sleep because empty job queue");
|
||||
if let Some(db) = conn.as_sql() {
|
||||
infos.sample_pool(db.size(), db.num_idle() as u32);
|
||||
}
|
||||
infos.add_iter(bench, false);
|
||||
}
|
||||
#[cfg(feature = "prometheus")]
|
||||
@@ -3047,7 +3063,7 @@ pub async fn handle_queued_job(
|
||||
.flatten()
|
||||
{
|
||||
tracing::debug!("Debug: {} going to sleep for {}", job.id, dbg_djob_sleep);
|
||||
sleep(std::time::Duration::from_secs(dbg_djob_sleep as u64)).await;
|
||||
tokio::time::sleep(std::time::Duration::from_secs(dbg_djob_sleep as u64)).await;
|
||||
}
|
||||
|
||||
tracing::debug!(
|
||||
|
||||
@@ -1365,35 +1365,31 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
let args_hm = args.unwrap_or_default().0;
|
||||
let args = PushArgs::from(&args_hm);
|
||||
if let Some(ck) = concurrency_key {
|
||||
let mut tx = db.begin().await?;
|
||||
insert_concurrency_key(
|
||||
&flow_job.workspace_id,
|
||||
&args,
|
||||
&flow_job.runnable_path,
|
||||
JobKind::Flow,
|
||||
Some(ck),
|
||||
&mut tx,
|
||||
db,
|
||||
flow,
|
||||
)
|
||||
.await?;
|
||||
tx.commit().await?;
|
||||
}
|
||||
if let Some(t) = tag {
|
||||
tag = Some(interpolate_args(t, &args, &flow_job.workspace_id));
|
||||
}
|
||||
} else if concurrent_limit.is_some() {
|
||||
let mut tx = db.begin().await?;
|
||||
insert_concurrency_key(
|
||||
&flow_job.workspace_id,
|
||||
&PushArgs::from(&HashMap::new()),
|
||||
&flow_job.runnable_path,
|
||||
JobKind::Flow,
|
||||
concurrency_key,
|
||||
&mut tx,
|
||||
db,
|
||||
flow,
|
||||
)
|
||||
.await?;
|
||||
tx.commit().await?;
|
||||
}
|
||||
|
||||
// let tag = tag_and_concurrency_key.and_then(|tc| tc.tag.map(|t| interpolate_args(t.clone(), &args, &workspace_id)));
|
||||
@@ -1644,7 +1640,6 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
true,
|
||||
None,
|
||||
false,
|
||||
false,
|
||||
)
|
||||
.await?;
|
||||
duration
|
||||
@@ -1665,7 +1660,6 @@ pub async fn update_flow_status_after_job_completion_internal(
|
||||
true,
|
||||
None,
|
||||
false,
|
||||
false,
|
||||
)
|
||||
.await?;
|
||||
duration
|
||||
@@ -3257,7 +3251,7 @@ async fn push_next_flow_job(
|
||||
"UPDATE v2_job_runtime SET ping = now() WHERE id = $1 AND ping < now()",
|
||||
flow_job.id,
|
||||
)
|
||||
.execute(db)
|
||||
.execute(&mut *tx)
|
||||
.warn_after_seconds(3)
|
||||
.await?;
|
||||
}
|
||||
@@ -5098,6 +5092,10 @@ pub async fn get_previous_job_result(
|
||||
Some(FlowStatusModule::Success { flow_jobs: Some(flow_jobs), .. }) => {
|
||||
Ok(Some(retrieve_flow_jobs_results(db, w_id, flow_jobs).await?))
|
||||
}
|
||||
Some(FlowStatusModule::Success { job, .. }) if *job == Uuid::nil() => {
|
||||
// Empty branch — no real job was executed, return empty object
|
||||
Ok(None)
|
||||
}
|
||||
Some(FlowStatusModule::Success { job, .. }) => Ok(Some(
|
||||
sqlx::query_scalar!(
|
||||
"SELECT result AS \"result!: Json<Box<RawValue>>\"
|
||||
|
||||
+1
-1
@@ -2,7 +2,7 @@ import { sleep } from "https://deno.land/x/sleep@v1.2.1/mod.ts";
|
||||
import * as windmill from "https://deno.land/x/windmill@v1.174.0/mod.ts";
|
||||
import * as api from "https://deno.land/x/windmill@v1.174.0/windmill-api/index.ts";
|
||||
|
||||
export const VERSION = "v1.628.3";
|
||||
export const VERSION = "v1.629.1";
|
||||
|
||||
export async function login(email: string, password: string): Promise<string> {
|
||||
return await windmill.UserService.login({
|
||||
|
||||
+1
-1
@@ -77,7 +77,7 @@ export {
|
||||
// }
|
||||
// });
|
||||
|
||||
export const VERSION = "1.628.3";
|
||||
export const VERSION = "1.629.1";
|
||||
|
||||
// Re-exported from constants.ts to maintain backwards compatibility
|
||||
export { WM_FORK_PREFIX } from "./core/constants.ts";
|
||||
|
||||
+2
-2
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"name": "windmill-utils-internal",
|
||||
"version": "1.3.2",
|
||||
"version": "1.3.4",
|
||||
"lockfileVersion": 3,
|
||||
"requires": true,
|
||||
"packages": {
|
||||
"": {
|
||||
"name": "windmill-utils-internal",
|
||||
"version": "1.3.2",
|
||||
"version": "1.3.4",
|
||||
"license": "Apache 2.0",
|
||||
"devDependencies": {
|
||||
"@types/node": "^24.2.0",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "windmill-utils-internal",
|
||||
"version": "1.3.3",
|
||||
"version": "1.3.4",
|
||||
"description": "Internal utility functions for Windmill",
|
||||
"main": "dist/cjs/index.js",
|
||||
"module": "dist/esm/index.js",
|
||||
|
||||
@@ -108,6 +108,9 @@ export function argSigToJsonSchemaType(
|
||||
} else if (t === "datetime") {
|
||||
newS.type = "string";
|
||||
newS.format = "date-time";
|
||||
} else if (t === "date") {
|
||||
newS.type = "string";
|
||||
newS.format = "date";
|
||||
} else if (typeof t !== "string" && "oneof" in t) {
|
||||
newS.type = "object";
|
||||
if (t.oneof) {
|
||||
|
||||
@@ -64,7 +64,7 @@ RUN --mount=type=secret,id=rh_username \
|
||||
RUN subscription-manager repos --enable codeready-builder-for-rhel-8-$(arch)-rpms
|
||||
|
||||
RUN yum update -y && \
|
||||
yum install -y perl-interpreter perl-IPC-Cmd perl-Time-Piece libxml2-devel xmlsec1-devel xmlsec1-openssl-devel krb5-devel cyrus-sasl-devel cyrus-sasl-gssapi clang llvm-devel cmake libtool-ltdl-devel
|
||||
yum install -y perl-interpreter perl-IPC-Cmd perl-Time-Piece libxml2-devel xmlsec1-devel xmlsec1-openssl-devel krb5-devel cyrus-sasl-devel cyrus-sasl-gssapi libcurl-devel clang llvm-devel cmake libtool-ltdl-devel
|
||||
|
||||
# RUN --mount=type=cache,target=/usr/local/cargo/registry \
|
||||
# CARGO_NET_GIT_FETCH_WITH_CLI=true RUST_BACKTRACE=1 cargo chef cook --release --features "$features" --recipe-path recipe.json
|
||||
|
||||
@@ -64,7 +64,7 @@ RUN --mount=type=secret,id=rh_username \
|
||||
RUN subscription-manager repos --enable codeready-builder-for-rhel-9-$(arch)-rpms
|
||||
|
||||
RUN yum update -y && \
|
||||
yum install -y perl-FindBin perl-IPC-Cmd perl-Time-Piece libxml2-devel xmlsec1-devel xmlsec1-openssl-devel krb5-devel cyrus-sasl-devel cyrus-sasl-gssapi clang llvm-devel cmake libtool-ltdl-devel
|
||||
yum install -y perl-FindBin perl-IPC-Cmd perl-Time-Piece libxml2-devel xmlsec1-devel xmlsec1-openssl-devel krb5-devel cyrus-sasl-devel cyrus-sasl-gssapi libcurl-devel clang llvm-devel cmake libtool-ltdl-devel
|
||||
|
||||
# RUN --mount=type=cache,target=/usr/local/cargo/registry \
|
||||
# CARGO_NET_GIT_FETCH_WITH_CLI=true RUST_BACKTRACE=1 cargo chef cook --release --features "$features" --recipe-path recipe.json
|
||||
|
||||
@@ -55,8 +55,11 @@
|
||||
nodejs
|
||||
postgresql
|
||||
pkg-config
|
||||
clang
|
||||
llvmPackages_18.clang
|
||||
mold
|
||||
cmake
|
||||
cyrus_sasl
|
||||
krb5
|
||||
];
|
||||
coursier = pkgs.fetchFromGitHub {
|
||||
owner = "coursier";
|
||||
@@ -66,7 +69,7 @@
|
||||
};
|
||||
|
||||
PKG_CONFIG_PATH = pkgs.lib.makeSearchPath "lib/pkgconfig"
|
||||
(with pkgs; [ openssl.dev libxml2.dev xmlsec.dev libxslt.dev ]);
|
||||
(with pkgs; [ openssl.dev libxml2.dev xmlsec.dev libxslt.dev cyrus_sasl.dev krb5.dev ]);
|
||||
RUSTY_V8_ARCHIVE = let
|
||||
# NOTE: needs to be same as in Cargo.toml
|
||||
version = "130.0.7";
|
||||
@@ -303,6 +306,14 @@
|
||||
# RUST_LOG = "debug";
|
||||
# RUST_LOG = "kube=debug";
|
||||
|
||||
# Override cargo linker to use clang 18 (stdenv brings clang 21 which causes SIGSEGV with mold)
|
||||
CARGO_TARGET_X86_64_UNKNOWN_LINUX_GNU_LINKER = "${pkgs.llvmPackages_18.clang}/bin/clang";
|
||||
CARGO_TARGET_AARCH64_UNKNOWN_LINUX_GNU_LINKER = "${pkgs.llvmPackages_18.clang}/bin/clang";
|
||||
CARGO_TARGET_X86_64_UNKNOWN_LINUX_GNU_RUSTFLAGS = "-C link-arg=-fuse-ld=mold -C link-arg=-Wl,-rpath,${pkgs.lib.makeLibraryPath [ pkgs.openssl pkgs.libffi pkgs.cyrus_sasl pkgs.krb5 pkgs.libxml2 pkgs.xmlsec pkgs.libxslt stdenv.cc.cc.lib ]}";
|
||||
CARGO_TARGET_AARCH64_UNKNOWN_LINUX_GNU_RUSTFLAGS = "-C link-arg=-fuse-ld=mold -C link-arg=-Wl,-rpath,${pkgs.lib.makeLibraryPath [ pkgs.openssl pkgs.libffi pkgs.cyrus_sasl pkgs.krb5 pkgs.libxml2 pkgs.xmlsec pkgs.libxslt stdenv.cc.cc.lib ]}";
|
||||
# rpath for build scripts and proc macros (host compilation)
|
||||
CARGO_HOST_RUSTFLAGS = "-C link-arg=-Wl,-rpath,${pkgs.lib.makeLibraryPath [ pkgs.openssl pkgs.libffi pkgs.cyrus_sasl pkgs.krb5 pkgs.libxml2 pkgs.xmlsec pkgs.libxslt stdenv.cc.cc.lib ]}";
|
||||
|
||||
# See this issue: https://github.com/NixOS/nixpkgs/issues/370494
|
||||
# Allows to build jemalloc on nixos
|
||||
CFLAGS = "-Wno-error=int-conversion";
|
||||
@@ -310,6 +321,7 @@
|
||||
# Need to tell bindgen where to find libclang
|
||||
LIBCLANG_PATH = "${pkgs.llvmPackages.libclang.lib}/lib";
|
||||
|
||||
# LD_LIBRARY_PATH set in shellHook with a wrapper to avoid leaking into git/ssh
|
||||
# LD_LIBRARY_PATH = "${pkgs.gcc.lib}/lib";
|
||||
|
||||
# Set C flags for Rust's bindgen program. Unlike ordinary C
|
||||
|
||||
Generated
+10
-169
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"name": "windmill-components",
|
||||
"version": "1.628.3",
|
||||
"version": "1.629.1",
|
||||
"lockfileVersion": 3,
|
||||
"requires": true,
|
||||
"packages": {
|
||||
"": {
|
||||
"name": "windmill-components",
|
||||
"version": "1.628.3",
|
||||
"version": "1.629.1",
|
||||
"hasInstallScript": true,
|
||||
"license": "AGPL-3.0",
|
||||
"dependencies": {
|
||||
@@ -79,14 +79,14 @@
|
||||
"windmill-parser-wasm-java": "1.510.1",
|
||||
"windmill-parser-wasm-nu": "1.510.1",
|
||||
"windmill-parser-wasm-php": "1.574.1",
|
||||
"windmill-parser-wasm-py": "1.623.1",
|
||||
"windmill-parser-wasm-py": "^1.628.3",
|
||||
"windmill-parser-wasm-regex": "1.625.0",
|
||||
"windmill-parser-wasm-ruby": "1.526.1",
|
||||
"windmill-parser-wasm-rust": "1.558.1",
|
||||
"windmill-parser-wasm-ts": "1.623.1",
|
||||
"windmill-parser-wasm-yaml": "1.593.0",
|
||||
"windmill-sql-datatype-parser-wasm": "1.512.0",
|
||||
"windmill-utils-internal": "^1.3.2",
|
||||
"windmill-utils-internal": "^1.3.4",
|
||||
"xterm": "^5.3.0",
|
||||
"xterm-readline": "^1.1.2",
|
||||
"y-monaco": "^0.1.4",
|
||||
@@ -2955,115 +2955,6 @@
|
||||
"integrity": "sha512-WmoN8qaIAo7WTYWbAZuG8PYEhn5fkz7dZrqTBZ7dtt//lL2Gwms1IcnQ5yHqjDfX8Ft5j4YzDM23f87zBfDe9g==",
|
||||
"license": "ISC"
|
||||
},
|
||||
"node_modules/@vitest/browser": {
|
||||
"version": "4.0.15",
|
||||
"resolved": "https://registry.npmjs.org/@vitest/browser/-/browser-4.0.15.tgz",
|
||||
"integrity": "sha512-zedtczX688KehaIaAv7m25CeDLb0gBtAOa2Oi1G1cqvSO5aLSVfH6lpZMJLW8BKYuWMxLQc9/5GYoM+jgvGIrw==",
|
||||
"dev": true,
|
||||
"license": "MIT",
|
||||
"optional": true,
|
||||
"peer": true,
|
||||
"dependencies": {
|
||||
"@vitest/mocker": "4.0.15",
|
||||
"@vitest/utils": "4.0.15",
|
||||
"magic-string": "^0.30.21",
|
||||
"pixelmatch": "7.1.0",
|
||||
"pngjs": "^7.0.0",
|
||||
"sirv": "^3.0.2",
|
||||
"tinyrainbow": "^3.0.3",
|
||||
"ws": "^8.18.3"
|
||||
},
|
||||
"funding": {
|
||||
"url": "https://opencollective.com/vitest"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"vitest": "4.0.15"
|
||||
}
|
||||
},
|
||||
"node_modules/@vitest/browser-playwright": {
|
||||
"version": "4.0.15",
|
||||
"resolved": "https://registry.npmjs.org/@vitest/browser-playwright/-/browser-playwright-4.0.15.tgz",
|
||||
"integrity": "sha512-94yVpDbb+ykiT7mK6ToonGnq2GIHEQGBTZTAzGxBGQXcVNCh54YKC2/WkfaDzxy0m6Kgw05kq3FYHKHu+wRdIA==",
|
||||
"dev": true,
|
||||
"license": "MIT",
|
||||
"optional": true,
|
||||
"peer": true,
|
||||
"dependencies": {
|
||||
"@vitest/browser": "4.0.15",
|
||||
"@vitest/mocker": "4.0.15",
|
||||
"tinyrainbow": "^3.0.3"
|
||||
},
|
||||
"funding": {
|
||||
"url": "https://opencollective.com/vitest"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"playwright": "*",
|
||||
"vitest": "4.0.15"
|
||||
},
|
||||
"peerDependenciesMeta": {
|
||||
"playwright": {
|
||||
"optional": false
|
||||
}
|
||||
}
|
||||
},
|
||||
"node_modules/@vitest/browser-playwright/node_modules/@vitest/mocker": {
|
||||
"version": "4.0.15",
|
||||
"resolved": "https://registry.npmjs.org/@vitest/mocker/-/mocker-4.0.15.tgz",
|
||||
"integrity": "sha512-CZ28GLfOEIFkvCFngN8Sfx5h+Se0zN+h4B7yOsPVCcgtiO7t5jt9xQh2E1UkFep+eb9fjyMfuC5gBypwb07fvQ==",
|
||||
"dev": true,
|
||||
"license": "MIT",
|
||||
"optional": true,
|
||||
"peer": true,
|
||||
"dependencies": {
|
||||
"@vitest/spy": "4.0.15",
|
||||
"estree-walker": "^3.0.3",
|
||||
"magic-string": "^0.30.21"
|
||||
},
|
||||
"funding": {
|
||||
"url": "https://opencollective.com/vitest"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"msw": "^2.4.9",
|
||||
"vite": "^6.0.0 || ^7.0.0-0"
|
||||
},
|
||||
"peerDependenciesMeta": {
|
||||
"msw": {
|
||||
"optional": true
|
||||
},
|
||||
"vite": {
|
||||
"optional": true
|
||||
}
|
||||
}
|
||||
},
|
||||
"node_modules/@vitest/browser/node_modules/@vitest/mocker": {
|
||||
"version": "4.0.15",
|
||||
"resolved": "https://registry.npmjs.org/@vitest/mocker/-/mocker-4.0.15.tgz",
|
||||
"integrity": "sha512-CZ28GLfOEIFkvCFngN8Sfx5h+Se0zN+h4B7yOsPVCcgtiO7t5jt9xQh2E1UkFep+eb9fjyMfuC5gBypwb07fvQ==",
|
||||
"dev": true,
|
||||
"license": "MIT",
|
||||
"optional": true,
|
||||
"peer": true,
|
||||
"dependencies": {
|
||||
"@vitest/spy": "4.0.15",
|
||||
"estree-walker": "^3.0.3",
|
||||
"magic-string": "^0.30.21"
|
||||
},
|
||||
"funding": {
|
||||
"url": "https://opencollective.com/vitest"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"msw": "^2.4.9",
|
||||
"vite": "^6.0.0 || ^7.0.0-0"
|
||||
},
|
||||
"peerDependenciesMeta": {
|
||||
"msw": {
|
||||
"optional": true
|
||||
},
|
||||
"vite": {
|
||||
"optional": true
|
||||
}
|
||||
}
|
||||
},
|
||||
"node_modules/@vitest/expect": {
|
||||
"version": "4.0.15",
|
||||
"resolved": "https://registry.npmjs.org/@vitest/expect/-/expect-4.0.15.tgz",
|
||||
@@ -9950,21 +9841,6 @@
|
||||
"node": ">= 6"
|
||||
}
|
||||
},
|
||||
"node_modules/pixelmatch": {
|
||||
"version": "7.1.0",
|
||||
"resolved": "https://registry.npmjs.org/pixelmatch/-/pixelmatch-7.1.0.tgz",
|
||||
"integrity": "sha512-1wrVzJ2STrpmONHKBy228LM1b84msXDUoAzVEl0R8Mz4Ce6EPr+IVtxm8+yvrqLYMHswREkjYFaMxnyGnaY3Ng==",
|
||||
"dev": true,
|
||||
"license": "ISC",
|
||||
"optional": true,
|
||||
"peer": true,
|
||||
"dependencies": {
|
||||
"pngjs": "^7.0.0"
|
||||
},
|
||||
"bin": {
|
||||
"pixelmatch": "bin/pixelmatch"
|
||||
}
|
||||
},
|
||||
"node_modules/pkg-types": {
|
||||
"version": "1.3.1",
|
||||
"resolved": "https://registry.npmjs.org/pkg-types/-/pkg-types-1.3.1.tgz",
|
||||
@@ -10040,18 +9916,6 @@
|
||||
"node": ">=4"
|
||||
}
|
||||
},
|
||||
"node_modules/pngjs": {
|
||||
"version": "7.0.0",
|
||||
"resolved": "https://registry.npmjs.org/pngjs/-/pngjs-7.0.0.tgz",
|
||||
"integrity": "sha512-LKWqWJRhstyYo9pGvgor/ivk2w94eSjE3RGVuzLGlr3NmD8bf7RcYGze1mNdEHRP6TRP6rMuDHk5t44hnTRyow==",
|
||||
"dev": true,
|
||||
"license": "MIT",
|
||||
"optional": true,
|
||||
"peer": true,
|
||||
"engines": {
|
||||
"node": ">=14.19.0"
|
||||
}
|
||||
},
|
||||
"node_modules/postcss": {
|
||||
"version": "8.5.6",
|
||||
"resolved": "https://registry.npmjs.org/postcss/-/postcss-8.5.6.tgz",
|
||||
@@ -14603,9 +14467,9 @@
|
||||
"integrity": "sha512-COyid6B1RYs+bpzUCInsA4HY/WZkpDLfkQ90+AqU/TVTpzYSbAC2JCbIwy0cRElBvlhI4bQ+9Wg6hSQKMpEkpA=="
|
||||
},
|
||||
"node_modules/windmill-parser-wasm-py": {
|
||||
"version": "1.623.1",
|
||||
"resolved": "https://registry.npmjs.org/windmill-parser-wasm-py/-/windmill-parser-wasm-py-1.623.1.tgz",
|
||||
"integrity": "sha512-lFBlZg6hvhHzsU5oPJq0478UyMTZ9UKVd8Hc8ggmxPIHZaJBeJ+56NR75hmGwg0VJcRff8ed+zEm1PgmjVhD+w=="
|
||||
"version": "1.628.3",
|
||||
"resolved": "https://registry.npmjs.org/windmill-parser-wasm-py/-/windmill-parser-wasm-py-1.628.3.tgz",
|
||||
"integrity": "sha512-TlluqknZpg8cZ+A3m6JFLPseY2PpKtDsxdj26fAnCUzKPtse8TxQR+n0dwC80rfW5TwdWSulvNGRDgcNuf7CTw=="
|
||||
},
|
||||
"node_modules/windmill-parser-wasm-regex": {
|
||||
"version": "1.625.0",
|
||||
@@ -14638,9 +14502,9 @@
|
||||
"integrity": "sha512-uHNL8F72/Tf96xF3hOHnPDjkEyqXw7fNjcPJiUhth9sTQkcwUIoJMOdwm8/cs+j9kKVRJ4tgNYMHEBLylazp6g=="
|
||||
},
|
||||
"node_modules/windmill-utils-internal": {
|
||||
"version": "1.3.2",
|
||||
"resolved": "https://registry.npmjs.org/windmill-utils-internal/-/windmill-utils-internal-1.3.2.tgz",
|
||||
"integrity": "sha512-sLE7sfmN3vBwPReqKorBrusnn83Qnm8vVFqh/7+kz8mf98ojMVDMgBginjlh9KBtVfWu9NhyjGS+WXRcnWalzA==",
|
||||
"version": "1.3.4",
|
||||
"resolved": "https://registry.npmjs.org/windmill-utils-internal/-/windmill-utils-internal-1.3.4.tgz",
|
||||
"integrity": "sha512-XVypDKIZ6P4fwIjZwvuvq1m+j0rtAA7BDp1rI2F7hQ+VBKZUHsLskP+jgstXs+kN1LqGGsQJj4ecMYDImpIZ6A==",
|
||||
"license": "Apache 2.0"
|
||||
},
|
||||
"node_modules/word-wrap": {
|
||||
@@ -14785,29 +14649,6 @@
|
||||
"node": "^14.17.0 || ^16.13.0 || >=18.0.0"
|
||||
}
|
||||
},
|
||||
"node_modules/ws": {
|
||||
"version": "8.19.0",
|
||||
"resolved": "https://registry.npmjs.org/ws/-/ws-8.19.0.tgz",
|
||||
"integrity": "sha512-blAT2mjOEIi0ZzruJfIhb3nps74PRWTCz1IjglWEEpQl5XS/UNama6u2/rjFkDDouqr4L67ry+1aGIALViWjDg==",
|
||||
"license": "MIT",
|
||||
"optional": true,
|
||||
"peer": true,
|
||||
"engines": {
|
||||
"node": ">=10.0.0"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"bufferutil": "^4.0.1",
|
||||
"utf-8-validate": ">=5.0.2"
|
||||
},
|
||||
"peerDependenciesMeta": {
|
||||
"bufferutil": {
|
||||
"optional": true
|
||||
},
|
||||
"utf-8-validate": {
|
||||
"optional": true
|
||||
}
|
||||
}
|
||||
},
|
||||
"node_modules/xml-utils": {
|
||||
"version": "1.10.2",
|
||||
"resolved": "https://registry.npmjs.org/xml-utils/-/xml-utils-1.10.2.tgz",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "windmill-components",
|
||||
"version": "1.628.3",
|
||||
"version": "1.629.1",
|
||||
"scripts": {
|
||||
"dev": "vite dev",
|
||||
"build": "vite build",
|
||||
@@ -149,14 +149,14 @@
|
||||
"windmill-parser-wasm-java": "1.510.1",
|
||||
"windmill-parser-wasm-nu": "1.510.1",
|
||||
"windmill-parser-wasm-php": "1.574.1",
|
||||
"windmill-parser-wasm-py": "1.623.1",
|
||||
"windmill-parser-wasm-py": "^1.628.3",
|
||||
"windmill-parser-wasm-regex": "1.625.0",
|
||||
"windmill-parser-wasm-ruby": "1.526.1",
|
||||
"windmill-parser-wasm-rust": "1.558.1",
|
||||
"windmill-parser-wasm-ts": "1.623.1",
|
||||
"windmill-parser-wasm-yaml": "1.593.0",
|
||||
"windmill-sql-datatype-parser-wasm": "1.512.0",
|
||||
"windmill-utils-internal": "^1.3.2",
|
||||
"windmill-utils-internal": "^1.3.4",
|
||||
"xterm": "^5.3.0",
|
||||
"xterm-readline": "^1.1.2",
|
||||
"y-monaco": "^0.1.4",
|
||||
|
||||
@@ -258,7 +258,11 @@
|
||||
nvalue = structuredClone($state.snapshot(defaultValue))
|
||||
if (defaultValue === undefined || defaultValue === null) {
|
||||
if (inputCat === 'string') {
|
||||
nvalue = nullable ? null : format === 'uuid' && extra?.['x-auto-generate'] ? randomUUID() : ''
|
||||
nvalue = nullable
|
||||
? null
|
||||
: format === 'uuid' && extra?.['x-auto-generate']
|
||||
? randomUUID()
|
||||
: ''
|
||||
} else if (inputCat == 'enum' && required) {
|
||||
let firstV = enum_?.[0]
|
||||
if (typeof firstV === 'string') {
|
||||
@@ -1177,7 +1181,6 @@
|
||||
titleClass="text-2xs"
|
||||
/>
|
||||
{/if}
|
||||
|
||||
{:else if disabled}
|
||||
<textarea disabled></textarea>
|
||||
{:else}
|
||||
|
||||
@@ -8,10 +8,13 @@
|
||||
import { Button } from './common'
|
||||
import Toggle from './Toggle.svelte'
|
||||
import { emptyString } from '$lib/utils'
|
||||
import { validateDeployPathFilters } from '$lib/validators/workspaceSettings'
|
||||
import Alert from './common/alert/Alert.svelte'
|
||||
import SettingsFooter from './workspaceSettings/SettingsFooter.svelte'
|
||||
|
||||
$: deployableWorkspaces = $usersWorkspaceStore?.workspaces
|
||||
.map((w) => w.id)
|
||||
.filter((w) => w != $workspaceStore)
|
||||
let deployableWorkspaces = $derived(
|
||||
$usersWorkspaceStore?.workspaces.map((w) => w.id).filter((w) => w != $workspaceStore)
|
||||
)
|
||||
|
||||
type DeployUITypeMap = {
|
||||
scripts: boolean
|
||||
@@ -34,14 +37,39 @@
|
||||
triggers: true
|
||||
}
|
||||
|
||||
export let workspaceToDeployTo: string | undefined
|
||||
export let deployUiSettings: {
|
||||
include_path: string[]
|
||||
include_type: DeployUITypeMap
|
||||
} = {
|
||||
include_path: [],
|
||||
include_type: all_ok
|
||||
}
|
||||
let {
|
||||
workspaceToDeployTo = $bindable(),
|
||||
deployUiSettings = $bindable({
|
||||
include_path: [],
|
||||
include_type: all_ok
|
||||
}),
|
||||
hasUnsavedChanges = false,
|
||||
onSave,
|
||||
onDiscard,
|
||||
onWorkspaceToDeployToSave
|
||||
}: {
|
||||
workspaceToDeployTo: string | undefined
|
||||
deployUiSettings: {
|
||||
include_path: string[]
|
||||
include_type: DeployUITypeMap
|
||||
}
|
||||
hasUnsavedChanges?: boolean
|
||||
onSave?: () => void
|
||||
onDiscard: () => void
|
||||
onWorkspaceToDeployToSave?: (workspaceToDeployTo: string | undefined) => void
|
||||
} = $props()
|
||||
|
||||
// Validation state
|
||||
let pathValidationErrors: Record<number, string> = $state({})
|
||||
let hasValidationErrors = $derived(Object.keys(pathValidationErrors).length > 0)
|
||||
|
||||
// Validate path filters whenever they change
|
||||
$effect(() => {
|
||||
if (deployUiSettings?.include_path) {
|
||||
const validationResult = validateDeployPathFilters(deployUiSettings.include_path)
|
||||
pathValidationErrors = validationResult.errors
|
||||
}
|
||||
})
|
||||
function deployUITypeMapToArray(
|
||||
typesMap: DeployUITypeMap,
|
||||
expectedValue: boolean
|
||||
@@ -71,39 +99,61 @@
|
||||
return result
|
||||
}
|
||||
|
||||
async function editWorkspaceToDeployTo() {
|
||||
try {
|
||||
await WorkspaceService.editDeployTo({
|
||||
workspace: $workspaceStore ?? '',
|
||||
requestBody: { deploy_to: workspaceToDeployTo === '' ? undefined : workspaceToDeployTo }
|
||||
})
|
||||
|
||||
if (workspaceToDeployTo === '' || workspaceToDeployTo === undefined) {
|
||||
sendUserToast('Disabled setting deployable workspace')
|
||||
onWorkspaceToDeployToSave?.(undefined)
|
||||
} else {
|
||||
sendUserToast('Set deployable workspace to ' + workspaceToDeployTo)
|
||||
onWorkspaceToDeployToSave?.(workspaceToDeployTo)
|
||||
}
|
||||
} catch (error) {
|
||||
sendUserToast(`Failed to save workspace deployment setting: ${error}`, true)
|
||||
}
|
||||
}
|
||||
|
||||
async function editWindmillDeploymentUISettings() {
|
||||
// Validate before saving
|
||||
const validationResult = validateDeployPathFilters(deployUiSettings.include_path)
|
||||
if (!validationResult.isValid) {
|
||||
sendUserToast('Please fix validation errors before saving', true)
|
||||
return
|
||||
}
|
||||
|
||||
let include_path = deployUiSettings.include_path.filter((elmt) => !emptyString(elmt))
|
||||
let include_type = deployUITypeMapToArray(deployUiSettings.include_type, true)
|
||||
await WorkspaceService.editWorkspaceDeployUiSettings({
|
||||
workspace: $workspaceStore!,
|
||||
requestBody: {
|
||||
deploy_ui_settings: {
|
||||
include_path: include_path,
|
||||
include_type: include_type
|
||||
|
||||
try {
|
||||
// Save workspace to deploy to first
|
||||
await editWorkspaceToDeployTo()
|
||||
|
||||
// Then save deployment UI settings
|
||||
await WorkspaceService.editWorkspaceDeployUiSettings({
|
||||
workspace: $workspaceStore!,
|
||||
requestBody: {
|
||||
deploy_ui_settings: {
|
||||
include_path: include_path,
|
||||
include_type: include_type
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
sendUserToast('Workspace Deployment UI settings updated')
|
||||
})
|
||||
sendUserToast('Workspace Deployment UI settings updated')
|
||||
onSave?.()
|
||||
} catch (error) {
|
||||
sendUserToast(`Failed to save deployment settings: ${error}`, true)
|
||||
}
|
||||
}
|
||||
</script>
|
||||
|
||||
<h3 class="mt-6 text-xs font-semibold text-emphasis">Workspace to link to</h3>
|
||||
<div class="flex min-w-0 mt-1">
|
||||
<select
|
||||
bind:value={workspaceToDeployTo}
|
||||
on:change={async (e) => {
|
||||
await WorkspaceService.editDeployTo({
|
||||
workspace: $workspaceStore ?? '',
|
||||
requestBody: { deploy_to: workspaceToDeployTo == '' ? undefined : workspaceToDeployTo }
|
||||
})
|
||||
if (workspaceToDeployTo == '') {
|
||||
workspaceToDeployTo = undefined
|
||||
sendUserToast('Disabled setting deployable workspace')
|
||||
} else {
|
||||
sendUserToast('Set deployable workspace to ' + workspaceToDeployTo)
|
||||
}
|
||||
}}
|
||||
>
|
||||
<select bind:value={workspaceToDeployTo}>
|
||||
{#if deployableWorkspaces?.length == 0}
|
||||
<option disabled>No workspace deployable to</option>
|
||||
{/if}
|
||||
@@ -113,8 +163,12 @@
|
||||
{/each}
|
||||
</select>
|
||||
</div>
|
||||
<h3 class="mt-6 mb-3 text-sm font-semibold text-emphasis">Deployable items</h3>
|
||||
<div class="flex flex-wrap gap-20">
|
||||
<h3 class="mt-6 mb-1 text-xs font-semibold text-emphasis">Deployable items</h3>
|
||||
<div class="text-xs text-secondary mb-1">
|
||||
You can filter which items can be deployed to the production workspace. By default everything is
|
||||
deployable.
|
||||
</div>
|
||||
<div class="flex flex-wrap gap-6 p-4 rounded-md border">
|
||||
<div class="max-w-md w-full">
|
||||
{#if Array.isArray(deployUiSettings?.include_path)}
|
||||
<h4 class="flex gap-2 mb-2 text-xs font-semibold text-emphasis"
|
||||
@@ -125,20 +179,36 @@
|
||||
anything including slashes.
|
||||
</Tooltip></h4
|
||||
>
|
||||
{#each deployUiSettings.include_path ?? [] as regexpPath, idx}
|
||||
<div class="flex mt-1 items-center">
|
||||
<input type="text" bind:value={regexpPath} id="arg-input-array" />
|
||||
<button
|
||||
transition:fade|local={{ duration: 100 }}
|
||||
class="rounded-full p-1 bg-surface-secondary duration-200 hover:bg-surface-hover ml-2"
|
||||
aria-label="Clear"
|
||||
on:click={() => {
|
||||
deployUiSettings.include_path.splice(idx, 1)
|
||||
deployUiSettings.include_path = [...deployUiSettings.include_path]
|
||||
}}
|
||||
>
|
||||
<X size={14} />
|
||||
</button>
|
||||
{#each deployUiSettings.include_path ?? [] as _, idx}
|
||||
<div class="flex flex-col mt-1">
|
||||
<div class="flex items-center">
|
||||
<input
|
||||
type="text"
|
||||
bind:value={deployUiSettings.include_path[idx]}
|
||||
id="arg-input-array-{idx}"
|
||||
class="flex-1 {pathValidationErrors[idx] ? 'border-red-500' : ''}"
|
||||
placeholder="e.g., f/*, u/admin/**"
|
||||
/>
|
||||
<button
|
||||
transition:fade|local={{ duration: 100 }}
|
||||
class="rounded-full p-1 bg-surface-secondary duration-200 hover:bg-surface-hover ml-2"
|
||||
aria-label="Clear"
|
||||
onclick={() => {
|
||||
deployUiSettings.include_path.splice(idx, 1)
|
||||
deployUiSettings.include_path = [...deployUiSettings.include_path]
|
||||
// Clear validation error for this index
|
||||
delete pathValidationErrors[idx]
|
||||
pathValidationErrors = { ...pathValidationErrors }
|
||||
}}
|
||||
>
|
||||
<X size={14} />
|
||||
</button>
|
||||
</div>
|
||||
{#if pathValidationErrors[idx]}
|
||||
<div class="text-xs text-red-600 dark:text-red-400 mt-1"
|
||||
>{pathValidationErrors[idx]}</div
|
||||
>
|
||||
{/if}
|
||||
</div>
|
||||
{/each}
|
||||
{/if}
|
||||
@@ -197,14 +267,18 @@
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
{#if hasValidationErrors}
|
||||
<Alert type="error" title="Validation Errors" class="mt-4">
|
||||
Please fix the validation errors in the path filters before saving.
|
||||
</Alert>
|
||||
{/if}
|
||||
{#if $enterpriseLicense}
|
||||
<div class="flex mt-5 mb-5 gap-1">
|
||||
<Button
|
||||
variant="accent"
|
||||
disabled={workspaceToDeployTo == undefined}
|
||||
on:click={() => {
|
||||
editWindmillDeploymentUISettings()
|
||||
}}>Save Deployment UI settings</Button
|
||||
>
|
||||
</div>
|
||||
<SettingsFooter
|
||||
{hasUnsavedChanges}
|
||||
onSave={editWindmillDeploymentUISettings}
|
||||
{onDiscard}
|
||||
saveLabel="Save deployment UI"
|
||||
disabled={workspaceToDeployTo == undefined || hasValidationErrors}
|
||||
class="border-none"
|
||||
/>
|
||||
{/if}
|
||||
|
||||
@@ -1,27 +1,16 @@
|
||||
<script lang="ts">
|
||||
import { Building } from 'lucide-svelte'
|
||||
import { twMerge } from 'tailwind-merge'
|
||||
import { Tooltip } from './meltComponents'
|
||||
import Tooltip from '$lib/components/meltComponents/Tooltip.svelte'
|
||||
import Badge from './common/badge/Badge.svelte'
|
||||
|
||||
interface Props {
|
||||
class?: string
|
||||
children?: import('svelte').Snippet
|
||||
}
|
||||
|
||||
let { class: className = '', children = undefined }: Props = $props()
|
||||
let { children = undefined }: Props = $props()
|
||||
</script>
|
||||
|
||||
<Tooltip>
|
||||
<div
|
||||
class={twMerge(
|
||||
'flex text-xs items-center gap-1 text-yellow-500 whitespace-nowrap px-1',
|
||||
className
|
||||
)}
|
||||
aria-label="Enterprise Edition only feature"
|
||||
role="tooltip"
|
||||
>
|
||||
EE only <Building size={16} />
|
||||
</div>
|
||||
<Badge verySmall color="blue" class="px-2">EE only</Badge>
|
||||
{#snippet text()}
|
||||
{#if children}
|
||||
{@render children()}
|
||||
|
||||
@@ -33,7 +33,7 @@
|
||||
import { base } from '$lib/base'
|
||||
import { enterpriseLicense, workspaceStore } from '$lib/stores'
|
||||
import MsTeamsIcon from '$lib/components/icons/MSTeamsIcon.svelte'
|
||||
import { emptySchema, emptyString, sendUserToast, tryEvery } from '$lib/utils'
|
||||
import { classNames, emptySchema, emptyString, sendUserToast, tryEvery } from '$lib/utils'
|
||||
import MultiSelect from '$lib/components/select/MultiSelect.svelte'
|
||||
import {
|
||||
FlowService,
|
||||
@@ -80,6 +80,7 @@
|
||||
customScriptTemplate: string
|
||||
customHandlerKind?: 'flow' | 'script'
|
||||
customTabTooltip?: import('svelte').Snippet
|
||||
noMargin?: boolean
|
||||
}
|
||||
|
||||
let {
|
||||
@@ -92,7 +93,8 @@
|
||||
handlerExtraArgs = $bindable(),
|
||||
customScriptTemplate,
|
||||
customHandlerKind = $bindable('script'),
|
||||
customTabTooltip
|
||||
customTabTooltip,
|
||||
noMargin = false
|
||||
}: Props = $props()
|
||||
|
||||
let customHandlerSchema: Schema | undefined = $state()
|
||||
@@ -297,21 +299,26 @@
|
||||
teams: undefined as string | undefined,
|
||||
email: undefined as string[] | undefined
|
||||
})
|
||||
let handlerPathCache: Partial<Record<ErrorHandler, string | undefined>> = $state({})
|
||||
$effect(() => {
|
||||
if (lastHandlerSelected !== handlerSelected && lastHandlerSelected !== undefined) {
|
||||
if (lastHandlerSelected != 'custom') {
|
||||
const key = lastHandlerSelected === 'email' ? EMAIL_RECIPIENTS_KEY : CHANNEL_KEY
|
||||
handlerCache[lastHandlerSelected] = handlerExtraArgs[key]
|
||||
}
|
||||
handlerPathCache[lastHandlerSelected] = handlerPath
|
||||
|
||||
if (handlerSelected === 'custom') {
|
||||
handlerExtraArgs[CHANNEL_KEY] = ''
|
||||
handlerExtraArgs[EMAIL_RECIPIENTS_KEY] = []
|
||||
handlerPath = undefined
|
||||
delete handlerExtraArgs[CHANNEL_KEY]
|
||||
delete handlerExtraArgs[EMAIL_RECIPIENTS_KEY]
|
||||
handlerPath = handlerPathCache['custom']
|
||||
} else if (handlerSelected === 'email') {
|
||||
handlerExtraArgs[EMAIL_RECIPIENTS_KEY] = handlerCache[handlerSelected] ?? []
|
||||
delete handlerExtraArgs[CHANNEL_KEY]
|
||||
} else {
|
||||
handlerExtraArgs[CHANNEL_KEY] = handlerCache[handlerSelected] ?? ''
|
||||
delete handlerExtraArgs[EMAIL_RECIPIENTS_KEY]
|
||||
handlerPath = handlerPathCache[handlerSelected]
|
||||
}
|
||||
}
|
||||
|
||||
@@ -343,7 +350,7 @@
|
||||
})
|
||||
</script>
|
||||
|
||||
<div class="mt-2 space-y-2">
|
||||
<div class={classNames('space-y-2', noMargin ? '' : 'mt-2')}>
|
||||
<ToggleButtonGroup bind:selected={handlerSelected} disabled={!isEditable}>
|
||||
{#snippet children({ item })}
|
||||
<ToggleButton label="Slack" value="slack" {item} disabled={!isEditable} />
|
||||
@@ -361,56 +368,60 @@
|
||||
|
||||
<div class="flex flex-col gap-6 p-4 rounded-md border">
|
||||
{#if handlerSelected === 'custom'}
|
||||
<div class="flex flex-row mb-6">
|
||||
<ScriptPicker
|
||||
disabled={!isEditable || !$enterpriseLicense}
|
||||
kinds={['script', 'failure']}
|
||||
allowFlow={true}
|
||||
bind:scriptPath={handlerPath}
|
||||
bind:itemKind={customHandlerKind}
|
||||
allowRefresh={isEditable}
|
||||
clearable
|
||||
/>
|
||||
<div class="flex flex-col gap-1">
|
||||
<div class="flex flex-row">
|
||||
<ScriptPicker
|
||||
disabled={!isEditable || !$enterpriseLicense}
|
||||
kinds={['script', 'failure']}
|
||||
allowFlow={true}
|
||||
bind:scriptPath={handlerPath}
|
||||
bind:itemKind={customHandlerKind}
|
||||
allowRefresh={isEditable}
|
||||
clearable
|
||||
/>
|
||||
|
||||
{#if !handlerPath}
|
||||
<Button
|
||||
btnClasses="ml-4 whitespace-nowrap"
|
||||
variant="default"
|
||||
size="xs"
|
||||
href={customScriptTemplate}
|
||||
disabled={!isEditable}
|
||||
target="_blank"
|
||||
>
|
||||
Create from template
|
||||
</Button>
|
||||
{#if !handlerPath}
|
||||
<Button
|
||||
btnClasses="ml-4 whitespace-nowrap"
|
||||
variant="default"
|
||||
size="xs"
|
||||
href={customScriptTemplate}
|
||||
disabled={!isEditable}
|
||||
target="_blank"
|
||||
>
|
||||
Create from template
|
||||
</Button>
|
||||
{/if}
|
||||
</div>
|
||||
{#if showScriptHelpText}
|
||||
<div class="text-2xs text-secondary">
|
||||
Example of error handler scripts can be found on <a
|
||||
target="_blank"
|
||||
href="{$hubBaseUrlStore}/failures"
|
||||
>
|
||||
Windmill Hub</a
|
||||
>
|
||||
</div>
|
||||
{/if}
|
||||
</div>
|
||||
{#if showScriptHelpText}
|
||||
<div class="text-2xs text-secondary">
|
||||
Example of error handler scripts can be found on <a
|
||||
target="_blank"
|
||||
href="{$hubBaseUrlStore}/failures"
|
||||
>
|
||||
Windmill Hub</a
|
||||
>
|
||||
</div>
|
||||
{/if}
|
||||
{#if handlerPath}
|
||||
<p class="font-semibold text-xs mt-6 mb-1">Extra arguments</p>
|
||||
{#await import('$lib/components/SchemaForm.svelte')}
|
||||
<Loader2 class="animate-spin" />
|
||||
{:then Module}
|
||||
<Module.default
|
||||
disabled={!isEditable}
|
||||
schema={customHandlerSchema}
|
||||
bind:args={handlerExtraArgs}
|
||||
shouldHideNoInputs
|
||||
className="text-xs"
|
||||
/>
|
||||
{/await}
|
||||
{#if customHandlerSchema && customHandlerSchema.properties && Object.keys(customHandlerSchema.properties).length === 0}
|
||||
<div class="text-xs texg-gray-700">This error handler takes no extra arguments</div>
|
||||
{/if}
|
||||
<div>
|
||||
<p class="font-semibold text-xs mb-1">Extra arguments</p>
|
||||
{#await import('$lib/components/SchemaForm.svelte')}
|
||||
<Loader2 class="animate-spin" />
|
||||
{:then Module}
|
||||
<Module.default
|
||||
disabled={!isEditable}
|
||||
schema={customHandlerSchema}
|
||||
bind:args={handlerExtraArgs}
|
||||
shouldHideNoInputs
|
||||
className="text-xs"
|
||||
/>
|
||||
{/await}
|
||||
{#if customHandlerSchema && customHandlerSchema.properties && Object.keys(customHandlerSchema.properties).length === 0}
|
||||
<div class="text-xs text-secondary">This error handler takes no extra arguments</div>
|
||||
{/if}
|
||||
</div>
|
||||
{/if}
|
||||
{:else if handlerSelected === 'slack'}
|
||||
<!-- Slack Connection Status -->
|
||||
|
||||
@@ -12,6 +12,8 @@
|
||||
selectedJobStepType: 'single' | 'forloop' | 'branchall'
|
||||
restartBranchNames?: [number, string][]
|
||||
flowPath?: string
|
||||
/** The flow version ID used in this run (from job.script_hash) */
|
||||
flowVersionId?: number
|
||||
disabled?: boolean
|
||||
enterpriseOnly?: boolean
|
||||
variant?: 'default' | 'accent'
|
||||
@@ -28,6 +30,7 @@
|
||||
selectedJobStepType,
|
||||
restartBranchNames = [],
|
||||
flowPath = undefined,
|
||||
flowVersionId = undefined,
|
||||
disabled = false,
|
||||
enterpriseOnly = false,
|
||||
variant = 'default',
|
||||
@@ -36,11 +39,15 @@
|
||||
onRestartComplete
|
||||
}: Props = $props()
|
||||
|
||||
// Sentinel value meaning "use the same version as the original run" (backend receives undefined)
|
||||
const RUN_VERSION_SENTINEL = -1
|
||||
|
||||
let branchOrIterationN = $state(0)
|
||||
let selectedVersionMode: 'run' | 'custom' = $state('run')
|
||||
let customFlowVersion: number | undefined = $state(undefined)
|
||||
let selectedFlowVersion: number = $state(RUN_VERSION_SENTINEL)
|
||||
let flowVersions: Array<FlowVersion> = $state([])
|
||||
let loadingVersions: boolean = $state(false)
|
||||
let versionsLoaded: boolean = $state(false)
|
||||
let runVersionInList: boolean = $state(false)
|
||||
|
||||
async function restartFlow(stepId: string, branchOrIterationN: number, flowVersion?: number) {
|
||||
let run = await JobService.restartFlowAtStep({
|
||||
@@ -56,16 +63,20 @@
|
||||
}
|
||||
|
||||
async function loadFlowVersions() {
|
||||
if (!flowPath || loadingVersions) return
|
||||
if (!flowPath || loadingVersions || versionsLoaded) return
|
||||
loadingVersions = true
|
||||
try {
|
||||
flowVersions = await FlowService.getFlowHistory({
|
||||
workspace: $workspaceStore!,
|
||||
path: flowPath
|
||||
})
|
||||
if (flowVersions.length > 0 && customFlowVersion === undefined) {
|
||||
customFlowVersion = flowVersions[0].id
|
||||
if (flowVersions.length > 0) {
|
||||
const match = flowVersionId ? flowVersions.find((v) => v.id === flowVersionId) : undefined
|
||||
runVersionInList = match !== undefined
|
||||
selectedFlowVersion =
|
||||
match?.id ?? (flowVersionId ? RUN_VERSION_SENTINEL : flowVersions[0].id)
|
||||
}
|
||||
versionsLoaded = true
|
||||
} catch (e) {
|
||||
sendUserToast('Failed to load flow versions', true)
|
||||
} finally {
|
||||
@@ -74,12 +85,12 @@
|
||||
}
|
||||
|
||||
function getFlowVersionForRestart(): number | undefined {
|
||||
if (selectedVersionMode === 'run') {
|
||||
return undefined // use run version
|
||||
} else if (selectedVersionMode === 'custom') {
|
||||
return customFlowVersion
|
||||
}
|
||||
return undefined
|
||||
return selectedFlowVersion === RUN_VERSION_SENTINEL ? undefined : selectedFlowVersion
|
||||
}
|
||||
|
||||
function formatVersionLabel(version: FlowVersion): string {
|
||||
const name = emptyString(version.deployment_msg) ? `v${version.id}` : version.deployment_msg!
|
||||
return `${name} - ${new Date(version.created_at).toLocaleString()}`
|
||||
}
|
||||
|
||||
function handleRestart() {
|
||||
@@ -95,37 +106,28 @@
|
||||
{#snippet flowVersionSelector()}
|
||||
<label>
|
||||
<div class="pb-1 text-xs font-semibold text-emphasis">Flow version</div>
|
||||
<div class="flex flex-col gap-2">
|
||||
<select
|
||||
bind:value={selectedVersionMode}
|
||||
class="grow"
|
||||
onchange={() => {
|
||||
if (selectedVersionMode === 'custom' && flowVersions.length === 0) {
|
||||
loadFlowVersions()
|
||||
}
|
||||
}}
|
||||
>
|
||||
<option value="run">Run version</option>
|
||||
<option value="custom">Specific version</option>
|
||||
</select>
|
||||
|
||||
{#if selectedVersionMode === 'custom'}
|
||||
{#if loadingVersions}
|
||||
<div class="text-xs text-secondary">Loading versions...</div>
|
||||
{:else if flowVersions.length > 0}
|
||||
<select bind:value={customFlowVersion} class="grow text-xs">
|
||||
{#each flowVersions as version}
|
||||
<option value={version.id}>
|
||||
{#if emptyString(version.deployment_msg)}Version {version.id}{:else}{version.deployment_msg}{/if}
|
||||
- {new Date(version.created_at).toLocaleString()}
|
||||
</option>
|
||||
{/each}
|
||||
</select>
|
||||
{:else}
|
||||
<div class="text-xs text-tertiary">No versions available</div>
|
||||
{#if loadingVersions}
|
||||
<div class="text-xs text-secondary">Loading versions...</div>
|
||||
{:else if flowVersions.length > 0}
|
||||
<select bind:value={selectedFlowVersion} class="w-full text-xs">
|
||||
{#if flowVersionId && !runVersionInList}
|
||||
<option value={RUN_VERSION_SENTINEL}>Same as run (v{flowVersionId})</option>
|
||||
{/if}
|
||||
{/if}
|
||||
</div>
|
||||
{#each flowVersions as version, i (version.id)}
|
||||
{@const isLatest = i === 0}
|
||||
{@const isSameAsRun = flowVersionId !== undefined && version.id === flowVersionId}
|
||||
<option value={version.id}>
|
||||
{formatVersionLabel(version)}{isSameAsRun
|
||||
? ' (Same as run)'
|
||||
: isLatest
|
||||
? ' (Latest)'
|
||||
: ''}
|
||||
</option>
|
||||
{/each}
|
||||
</select>
|
||||
{:else}
|
||||
<div class="text-xs text-tertiary">No versions available</div>
|
||||
{/if}
|
||||
</label>
|
||||
{/snippet}
|
||||
{#snippet singleRestartButton()}
|
||||
@@ -159,12 +161,13 @@
|
||||
<Popover
|
||||
floatingConfig={{ strategy: 'absolute', placement: 'bottom-start' }}
|
||||
disablePopup={!flowPath}
|
||||
on:openChange={(e) => { if (e.detail) loadFlowVersions() }}
|
||||
>
|
||||
{#snippet trigger()}
|
||||
{@render singleRestartButton()}
|
||||
{/snippet}
|
||||
{#snippet content()}
|
||||
<div class="flex flex-col gap-4 text-primary p-4 min-w-64">
|
||||
<div class="flex flex-col gap-4 text-primary p-4 w-80">
|
||||
{@render flowVersionSelector()}
|
||||
|
||||
<Button variant="accent" onClick={handleRestart}>Restart</Button>
|
||||
@@ -173,7 +176,10 @@
|
||||
</Popover>
|
||||
{/if}
|
||||
{:else}
|
||||
<Popover floatingConfig={{ strategy: 'absolute', placement: 'bottom-start' }}>
|
||||
<Popover
|
||||
floatingConfig={{ strategy: 'absolute', placement: 'bottom-start' }}
|
||||
on:openChange={(e) => { if (e.detail) loadFlowVersions() }}
|
||||
>
|
||||
{#snippet trigger()}
|
||||
<Button
|
||||
title={`Re-start this flow from step ${selectedJobStep} (included).${enterpriseOnly ? ' This is a feature only available in enterprise edition.' : ''}`}
|
||||
@@ -193,7 +199,7 @@
|
||||
</Button>
|
||||
{/snippet}
|
||||
{#snippet content()}
|
||||
<div class="flex flex-col gap-4 text-primary p-4 min-w-64">
|
||||
<div class="flex flex-col gap-4 text-primary p-4 w-80">
|
||||
<label>
|
||||
<div class="pb-1 text-xs font-semibold text-emphasis"
|
||||
>{selectedJobStepType == 'forloop' ? 'From iteration #' : 'From branch'}</div
|
||||
@@ -203,7 +209,7 @@
|
||||
<input type="number" min="0" bind:value={branchOrIterationN} class="!w-32 grow" />
|
||||
{:else}
|
||||
<select bind:value={branchOrIterationN} class="!w-32 grow">
|
||||
{#each restartBranchNames as [branchIdx, branchName]}
|
||||
{#each restartBranchNames as [branchIdx, branchName] (branchIdx)}
|
||||
<option value={branchIdx}>{branchName}</option>
|
||||
{/each}
|
||||
</select>
|
||||
|
||||
@@ -177,7 +177,7 @@
|
||||
{/snippet}
|
||||
|
||||
<!-- {JSON.stringify($values, null, 2)} -->
|
||||
{#if (!setting.cloudonly || isCloudHosted()) && showSetting(setting.key, $values) && !(setting.hiddenIfNull && $values[setting.key] == null) && !(setting.hiddenIfEmpty && !$values[setting.key])}
|
||||
{#if (!setting.cloudonly || isCloudHosted()) && showSetting(setting.key, $values) && !(setting.hiddenIfNull && $values[setting.key] == null) && !(setting.hiddenIfEmpty && !$values[setting.key]) && !(setting.hiddenInEe && $enterpriseLicense)}
|
||||
{#if setting.fieldType == 'select'}
|
||||
<div>
|
||||
{@render LabelSnippet()}
|
||||
@@ -255,7 +255,7 @@
|
||||
<div class="text-emphasis font-semibold text-xs flex flex-col gap-1 w-full">
|
||||
<div class="flex items-center justify-between gap-2 w-full">
|
||||
{#if setting.fieldType != 'smtp_connect'}
|
||||
<div class="flex gap-1">
|
||||
<div class="flex gap-1 items-baseline">
|
||||
<span class="text-emphasis font-semibold text-xs pb-1">{setting.label}</span>
|
||||
{#if setting.ee_only != undefined && !$enterpriseLicense}
|
||||
{#if setting.ee_only != ''}
|
||||
@@ -834,7 +834,10 @@
|
||||
{/if}
|
||||
</div>
|
||||
{:else if setting.fieldType == 'otel_tracing_proxy'}
|
||||
{@const tracingProxyVal = $values[setting.key] ?? { enabled: false, enabled_languages: [...OTEL_TRACING_PROXY_LANGUAGES] }}
|
||||
{@const tracingProxyVal = $values[setting.key] ?? {
|
||||
enabled: false,
|
||||
enabled_languages: [...OTEL_TRACING_PROXY_LANGUAGES]
|
||||
}}
|
||||
<div class="flex flex-col gap-4">
|
||||
<Toggle
|
||||
id="otel_tracing_proxy_enabled"
|
||||
|
||||
@@ -231,6 +231,28 @@
|
||||
}
|
||||
}
|
||||
|
||||
let downloadingStats = $state(false)
|
||||
async function downloadStats() {
|
||||
try {
|
||||
downloadingStats = true
|
||||
const encryptedData = await SettingService.getStats()
|
||||
const blob = new Blob([encryptedData], { type: 'application/octet-stream' })
|
||||
const url = URL.createObjectURL(blob)
|
||||
const a = document.createElement('a')
|
||||
a.href = url
|
||||
a.download = `windmill-telemetry-${new Date().toISOString().split('T')[0]}.enc`
|
||||
document.body.appendChild(a)
|
||||
a.click()
|
||||
document.body.removeChild(a)
|
||||
URL.revokeObjectURL(url)
|
||||
sendUserToast('Telemetry data downloaded')
|
||||
} catch (err) {
|
||||
throw err
|
||||
} finally {
|
||||
downloadingStats = false
|
||||
}
|
||||
}
|
||||
|
||||
function isValidTeamsChannel(value: any): value is TeamsChannel {
|
||||
return (
|
||||
typeof value === 'object' &&
|
||||
@@ -341,18 +363,29 @@
|
||||
<div class="text-primary pb-4 text-xs">
|
||||
On Enterprise Edition, you must send data to check that usage is in line with the
|
||||
terms of the subscription. You can either enable telemetry or regularly send usage
|
||||
data by clicking the button below.
|
||||
data by clicking the button below. For air-gapped instances, you can download the
|
||||
telemetry data and send it manually.
|
||||
</div>
|
||||
<div class="flex gap-2 mb-4">
|
||||
<Button
|
||||
on:click={sendStats}
|
||||
variant="default"
|
||||
btnClasses="w-auto"
|
||||
loading={sendingStats}
|
||||
size="xs"
|
||||
>
|
||||
Send usage
|
||||
</Button>
|
||||
<Button
|
||||
on:click={downloadStats}
|
||||
variant="default"
|
||||
btnClasses="w-auto"
|
||||
loading={downloadingStats}
|
||||
size="xs"
|
||||
>
|
||||
Download usage
|
||||
</Button>
|
||||
</div>
|
||||
<Button
|
||||
on:click={sendStats}
|
||||
variant="default"
|
||||
btnClasses="w-auto"
|
||||
wrapperClasses="mb-4"
|
||||
loading={sendingStats}
|
||||
size="xs"
|
||||
>
|
||||
Send usage
|
||||
</Button>
|
||||
{/if}
|
||||
{:else if category == 'Auth/OAuth/SAML'}
|
||||
<AuthSettings
|
||||
|
||||
@@ -102,9 +102,9 @@
|
||||
transition:slide={animate || collapsable ? { duration: 200 } : { duration: 0 }}
|
||||
>
|
||||
{#if description}
|
||||
<div class="text-xs text-primary mt-1">{@html description}</div>
|
||||
<div class="text-xs text-primary mt-1 mb-2">{@html description}</div>
|
||||
{/if}
|
||||
<div class="flex flex-col gap-6 grow min-h-0 mt-6">
|
||||
<div class="flex flex-col gap-6 grow min-h-0 mt-4">
|
||||
<div class={twMerge('grow min-h-0', clazz)}>
|
||||
{@render children?.()}
|
||||
</div>
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
<script lang="ts">
|
||||
import type { ComponentType } from 'svelte'
|
||||
import { twMerge } from 'tailwind-merge'
|
||||
import Button from '$lib/components/common/button/Button.svelte'
|
||||
import EEOnly from '$lib/components/EEOnly.svelte'
|
||||
import { enterpriseLicense } from '$lib/stores'
|
||||
|
||||
interface NavigationItem {
|
||||
id: string
|
||||
label: string
|
||||
icon?: ComponentType
|
||||
disabled?: boolean
|
||||
count?: number
|
||||
aiId?: string
|
||||
aiDescription?: string
|
||||
showIf?: boolean
|
||||
isEE?: boolean
|
||||
}
|
||||
|
||||
interface NavigationGroup {
|
||||
title?: string
|
||||
items: NavigationItem[]
|
||||
}
|
||||
|
||||
interface Props {
|
||||
groups: NavigationGroup[]
|
||||
selectedId: string
|
||||
onNavigate: (id: string) => void
|
||||
class?: string
|
||||
}
|
||||
|
||||
let { groups, selectedId, onNavigate, class: className = '' }: Props = $props()
|
||||
</script>
|
||||
|
||||
<div class={twMerge('flex flex-col gap-6', className)}>
|
||||
{#each groups as group (group.title)}
|
||||
<div class="flex flex-col gap-1">
|
||||
{#if group.title}
|
||||
<div class="text-sm font-semibold text-emphasis px-2 mb-1">
|
||||
{group.title}
|
||||
</div>
|
||||
{/if}
|
||||
<nav class="flex flex-col gap-0.5">
|
||||
{#each group.items as item (item.id)}
|
||||
{#if item.showIf !== false}
|
||||
{@const isSelected = selectedId === item.id}
|
||||
<Button
|
||||
variant="subtle"
|
||||
unifiedSize="sm"
|
||||
selected={isSelected}
|
||||
disabled={item.disabled}
|
||||
aiId={item.aiId}
|
||||
aiDescription={item.aiDescription}
|
||||
startIcon={item.icon ? { icon: item.icon } : undefined}
|
||||
btnClasses={'!justify-start text-left !w-full'}
|
||||
onClick={() => onNavigate(item.id)}
|
||||
>
|
||||
<span class="truncate">{item.label}</span>
|
||||
<div class="ml-auto flex items-center gap-1">
|
||||
{#if item.isEE && !$enterpriseLicense}
|
||||
<EEOnly />
|
||||
{/if}
|
||||
{#if item.count !== undefined}
|
||||
<span
|
||||
class="text-2xs text-secondary bg-surface-secondary px-1.5 py-0.5 rounded-full"
|
||||
>
|
||||
{item.count}
|
||||
</span>
|
||||
{/if}
|
||||
</div>
|
||||
</Button>
|
||||
{/if}
|
||||
{/each}
|
||||
</nav>
|
||||
</div>
|
||||
{/each}
|
||||
</div>
|
||||
@@ -899,11 +899,6 @@ class AIChatManager {
|
||||
}
|
||||
|
||||
saveAndClear = async () => {
|
||||
console.log('saveAndClear called', {
|
||||
hasAbortController: !!this.abortController,
|
||||
isLoading: this.loading,
|
||||
stack: new Error().stack
|
||||
})
|
||||
this.cancel('saveAndClear')
|
||||
await this.historyManager.save(this.displayMessages, this.messages)
|
||||
this.displayMessages = []
|
||||
|
||||
@@ -1,6 +1,12 @@
|
||||
import { ResourceService, JobService } from '$lib/gen/services.gen'
|
||||
import type { AIProvider, AIProviderModel, ResourceType, ScriptLang } from '$lib/gen/types.gen'
|
||||
import { capitalize, isObject, toCamel } from '$lib/utils'
|
||||
import { ResourceService, JobService, FlowService } from '$lib/gen/services.gen'
|
||||
import type {
|
||||
AIProvider,
|
||||
AIProviderModel,
|
||||
Flow,
|
||||
ResourceType,
|
||||
ScriptLang
|
||||
} from '$lib/gen/types.gen'
|
||||
import { capitalize, emptyString, isObject, toCamel } from '$lib/utils'
|
||||
import { get } from 'svelte/store'
|
||||
import { compile, phpCompile, pythonCompile } from '../../utils'
|
||||
import type {
|
||||
@@ -16,9 +22,11 @@ import {
|
||||
executeTestRun,
|
||||
buildTestRunArgs,
|
||||
buildContextString,
|
||||
extractAllModules,
|
||||
type ScriptLintResult,
|
||||
formatScriptLintResult
|
||||
} from '../shared'
|
||||
import uFuzzy from '@leeoniya/ufuzzy'
|
||||
import { setupTypeAcquisition, type DepsToGet } from '$lib/ata'
|
||||
import { getModelContextWindow } from '../../lib'
|
||||
import type { ReviewChangesOpts } from '../monaco-adapter'
|
||||
@@ -178,6 +186,7 @@ function buildChatSystemPrompt(currentModel: AIProviderModel) {
|
||||
- You can also receive a \`DIFF\` of the changes that have been made to the code. You should use this diff to give better answers.
|
||||
- Before giving your answer, check again that you carefully followed these instructions.
|
||||
- When asked to create a script that communicates with an external service, you can use the \`search_hub_scripts\` tool to search for relevant scripts in the hub. Make sure the language is the same as what the user is coding in. If you do not find any relevant scripts, you can use the \`search_npm_packages\` tool to search for relevant packages and their documentation. Always give a link to the documentation in your answer if possible.
|
||||
- If the user mentions a flow, you can use the \`search_flows\` tool to find it, and then \`get_flow_details\` to read its details and possibly reuse any inline script module inside it.
|
||||
- After applying code changes with the \`${editToolName}\` tool, ALWAYS use the \`get_lint_errors\` tool to check for lint errors. If there are errors, fix them before proceeding. Then use the \`test_run_script\` tool to test the code, and iterate on the code until it works as expected (MAX 3 times). If the user cancels the test run, do not try again and wait for the next user instruction.
|
||||
|
||||
Important:
|
||||
@@ -329,6 +338,8 @@ export function prepareScriptTools(
|
||||
}
|
||||
tools.push(testRunScriptTool)
|
||||
tools.push(getLintErrorsTool)
|
||||
tools.push(searchFlowsTool)
|
||||
tools.push(getFlowDetailsTool)
|
||||
return tools
|
||||
}
|
||||
|
||||
@@ -909,3 +920,140 @@ export const getLintErrorsTool: Tool<ScriptChatHelpers> = {
|
||||
return formatScriptLintResult(lintResult)
|
||||
}
|
||||
}
|
||||
|
||||
// ============= Flow Search Tools =============
|
||||
|
||||
class WorkspaceFlowsSearch {
|
||||
private uf: uFuzzy
|
||||
private workspace: string | undefined = undefined
|
||||
private flows: Flow[] | undefined = undefined
|
||||
|
||||
constructor() {
|
||||
this.uf = new uFuzzy()
|
||||
}
|
||||
|
||||
private async init(workspace: string) {
|
||||
if (this.flows === undefined || this.workspace !== workspace) {
|
||||
this.flows = await FlowService.listFlows({ workspace })
|
||||
this.workspace = workspace
|
||||
}
|
||||
}
|
||||
|
||||
async search(query: string, workspace: string) {
|
||||
await this.init(workspace)
|
||||
const flows = this.flows
|
||||
if (!flows) return []
|
||||
|
||||
const results = this.uf.search(
|
||||
flows.map((f) => (emptyString(f.summary) ? f.path : f.summary + ' (' + f.path + ')')),
|
||||
query.trim()
|
||||
)
|
||||
return (
|
||||
results[2]?.map((id) => ({
|
||||
path: flows[id].path,
|
||||
summary: flows[id].summary,
|
||||
description: flows[id].description
|
||||
})) ?? []
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
const workspaceFlowsSearch = new WorkspaceFlowsSearch()
|
||||
|
||||
const SEARCH_FLOWS_TOOL: ChatCompletionFunctionTool = {
|
||||
type: 'function',
|
||||
function: {
|
||||
name: 'search_flows',
|
||||
description:
|
||||
'Search for flows in the workspace. Use this when the user mentions a flow, wants to find existing flows, or wants to reuse inline script code from a flow.',
|
||||
parameters: {
|
||||
type: 'object',
|
||||
properties: {
|
||||
query: {
|
||||
type: 'string',
|
||||
description: 'The search query (e.g. "invoice processing", "stripe webhook")'
|
||||
}
|
||||
},
|
||||
required: ['query'],
|
||||
additionalProperties: false
|
||||
},
|
||||
strict: true
|
||||
}
|
||||
}
|
||||
|
||||
export const searchFlowsTool: Tool<ScriptChatHelpers> = {
|
||||
def: SEARCH_FLOWS_TOOL,
|
||||
fn: async ({ args, workspace, toolId, toolCallbacks }) => {
|
||||
toolCallbacks.setToolStatus(toolId, {
|
||||
content: 'Searching for flows related to "' + args.query + '"...'
|
||||
})
|
||||
const flowResults = await workspaceFlowsSearch.search(args.query, workspace)
|
||||
toolCallbacks.setToolStatus(toolId, {
|
||||
content: 'Found ' + flowResults.length + ' flow(s) related to "' + args.query + '"'
|
||||
})
|
||||
return JSON.stringify(flowResults)
|
||||
}
|
||||
}
|
||||
|
||||
const MAX_INLINE_SCRIPT_LENGTH = 2000
|
||||
|
||||
const GET_FLOW_DETAILS_TOOL: ChatCompletionFunctionTool = {
|
||||
type: 'function',
|
||||
function: {
|
||||
name: 'get_flow_details',
|
||||
description:
|
||||
'Get the details of a flow including its modules and inline script code. Use after search_flows to inspect a specific flow and potentially reuse its inline scripts.',
|
||||
parameters: {
|
||||
type: 'object',
|
||||
properties: {
|
||||
path: {
|
||||
type: 'string',
|
||||
description: 'The path of the flow (e.g. "f/ops/process_invoices")'
|
||||
}
|
||||
},
|
||||
required: ['path'],
|
||||
additionalProperties: false
|
||||
},
|
||||
strict: true
|
||||
}
|
||||
}
|
||||
|
||||
export const getFlowDetailsTool: Tool<ScriptChatHelpers> = {
|
||||
def: GET_FLOW_DETAILS_TOOL,
|
||||
fn: async ({ args, workspace, toolId, toolCallbacks }) => {
|
||||
toolCallbacks.setToolStatus(toolId, {
|
||||
content: 'Fetching flow details for "' + args.path + '"...'
|
||||
})
|
||||
const flow = await FlowService.getFlowByPath({ workspace, path: args.path })
|
||||
const modules = extractAllModules(flow.value.modules)
|
||||
const moduleDetails = modules.map((m) => {
|
||||
const base: Record<string, unknown> = {
|
||||
id: m.id,
|
||||
summary: m.summary,
|
||||
type: m.value.type
|
||||
}
|
||||
if (m.value.type === 'rawscript') {
|
||||
base.language = m.value.language
|
||||
const content = m.value.content ?? ''
|
||||
base.content =
|
||||
content.length > MAX_INLINE_SCRIPT_LENGTH
|
||||
? content.slice(0, MAX_INLINE_SCRIPT_LENGTH) + '...(truncated)'
|
||||
: content
|
||||
} else if (m.value.type === 'script') {
|
||||
base.path = m.value.path
|
||||
}
|
||||
return base
|
||||
})
|
||||
const result = {
|
||||
path: flow.path,
|
||||
summary: flow.summary,
|
||||
description: flow.description,
|
||||
schema: flow.schema,
|
||||
modules: moduleDetails
|
||||
}
|
||||
toolCallbacks.setToolStatus(toolId, {
|
||||
content: 'Retrieved flow details for "' + args.path + '"'
|
||||
})
|
||||
return JSON.stringify(result)
|
||||
}
|
||||
}
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user