Merge branch 'main' into perf-since-arg-historic-inputs

This commit is contained in:
pyranota
2025-02-18 15:33:46 +03:00
committed by GitHub
36 changed files with 337 additions and 180 deletions
+2
View File
@@ -162,6 +162,7 @@ jobs:
-c
https://raw.githubusercontent.com/windmill-labs/windmill/${GITHUB_REF##ref/head/}/benchmarks/suite_config.json
--workers 4
--factor 3
- name: Save benchmark results
uses: actions/upload-artifact@v4
with:
@@ -281,6 +282,7 @@ jobs:
-c
https://raw.githubusercontent.com/windmill-labs/windmill/${GITHUB_REF##ref/head/}/benchmarks/suite_config.json
--workers 8
--factor 3
- name: Save benchmark results
uses: actions/upload-artifact@v4
with:
+16
View File
@@ -1,5 +1,21 @@
# Changelog
## [1.463.5](https://github.com/windmill-labs/windmill/compare/v1.463.4...v1.463.5) (2025-02-18)
### Bug Fixes
* fix teams cleanup preventing start ([1b46e0f](https://github.com/windmill-labs/windmill/commit/1b46e0f08426497d549cf5007c93981df9ab41e5))
## [1.463.4](https://github.com/windmill-labs/windmill/compare/v1.463.3...v1.463.4) (2025-02-17)
### Bug Fixes
* improve queue job indices for faster performances ([9530826](https://github.com/windmill-labs/windmill/commit/953082681e2c4fd71d5ac1acf372265ccc72297b))
* improve teams settings in workspace settings ([#5316](https://github.com/windmill-labs/windmill/issues/5316)) ([935b5b7](https://github.com/windmill-labs/windmill/commit/935b5b799636c0f02597315837268d4a76f6709a))
## [1.463.3](https://github.com/windmill-labs/windmill/compare/v1.463.2...v1.463.3) (2025-02-17)
@@ -1,26 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "\n WITH assigned_teams AS (\n SELECT teams_team_id\n FROM workspace_settings\n ),\n all_teams AS (\n SELECT jsonb_array_elements(value::jsonb) AS team\n FROM global_settings\n WHERE name = 'teams'\n )\n SELECT team->>'team_name' AS team_name, team->>'team_internal_id' AS team_id\n FROM all_teams\n WHERE NOT EXISTS (\n SELECT 1\n FROM assigned_teams\n WHERE assigned_teams.teams_team_id = team->>'team_internal_id'\n )\n AND team->>'team_name' IS NOT NULL\n AND team->>'team_id' IS NOT NULL\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "team_name",
"type_info": "Text"
},
{
"ordinal": 1,
"name": "team_id",
"type_info": "Text"
}
],
"parameters": {
"Left": []
},
"nullable": [
null,
null
]
},
"hash": "2fa27b0a71740e4e168879cf9a58ed593004c22bb5d11170b3196e290dd1d966"
}
@@ -1,12 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "create index concurrently if not exists root_job_index_by_path_2 ON v2_job (workspace_id, runnable_path, created_at desc) WHERE parent_job IS NULL",
"describe": {
"columns": [],
"parameters": {
"Left": []
},
"nullable": []
},
"hash": "3bacf9cd9aa63f4bec5f983f4a0c3030216b5a4ed669f77962509d1c2c6cb780"
}
@@ -0,0 +1,26 @@
{
"db_name": "PostgreSQL",
"query": "\n WITH assigned_teams AS (\n SELECT teams_team_id\n FROM workspace_settings\n ),\n all_teams AS (\n SELECT jsonb_array_elements(CASE\n WHEN jsonb_typeof(value::jsonb) = 'array' THEN value::jsonb\n ELSE '[]'::jsonb\n END) AS team\n FROM global_settings\n WHERE name = 'teams'\n )\n SELECT team->>'team_name' AS team_name, team->>'team_internal_id' AS team_id\n FROM all_teams\n WHERE NOT EXISTS (\n SELECT 1\n FROM assigned_teams\n WHERE assigned_teams.teams_team_id = team->>'team_internal_id'\n )\n AND team->>'team_name' IS NOT NULL\n AND team->>'team_id' IS NOT NULL\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "team_name",
"type_info": "Text"
},
{
"ordinal": 1,
"name": "team_id",
"type_info": "Text"
}
],
"parameters": {
"Left": []
},
"nullable": [
null,
null
]
},
"hash": "50c17c7848760aaf0f869acfc444caecda6335eda5b2e97e5a7370361653ff48"
}
@@ -0,0 +1,12 @@
{
"db_name": "PostgreSQL",
"query": "DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path_2",
"describe": {
"columns": [],
"parameters": {
"Left": []
},
"nullable": []
},
"hash": "6536214f31e9d600e868b01385d8c6395e2440ea27553b7ccb18d7149b106728"
}
@@ -0,0 +1,12 @@
{
"db_name": "PostgreSQL",
"query": "\n UPDATE workspace_settings\n SET teams_command_script = NULL,\n teams_team_id = NULL,\n teams_team_name = NULL\n ",
"describe": {
"columns": [],
"parameters": {
"Left": []
},
"nullable": []
},
"hash": "65c339164e7669360d231d70105849e72bdc197c17c0fc51777c1dc9267e2daf"
}
@@ -0,0 +1,12 @@
{
"db_name": "PostgreSQL",
"query": "\n UPDATE global_settings\n SET value = (\n SELECT COALESCE(jsonb_agg(elem), '[]'::jsonb)\n FROM jsonb_array_elements(value) AS elem\n WHERE NOT (elem ? 'teams_channel')\n )\n WHERE name = 'critical_error_channels'\n ",
"describe": {
"columns": [],
"parameters": {
"Left": []
},
"nullable": []
},
"hash": "81b06122c7a12a314d8905ba5c7c14aa7614f2610e79a8c7302eaa63fb74984d"
}
@@ -0,0 +1,12 @@
{
"db_name": "PostgreSQL",
"query": "create index concurrently if not exists ix_job_root_job_index_by_path_2 ON v2_job (workspace_id, runnable_path, created_at desc) WHERE parent_job IS NULL",
"describe": {
"columns": [],
"parameters": {
"Left": []
},
"nullable": []
},
"hash": "c481e5d63ebf1aa537cc4ce4e84f9a71af5996bc76f328b3ba1cf68a71880462"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM v2_job_queue WHERE workspace_id = $1 AND id = $2 RETURNING 1",
"query": "DELETE FROM v2_job_queue WHERE id = $1 RETURNING 1",
"describe": {
"columns": [
{
@@ -11,7 +11,6 @@
],
"parameters": {
"Left": [
"Text",
"Uuid"
]
},
@@ -19,5 +18,5 @@
null
]
},
"hash": "d25c58d2722ad3dcd91101ce6f66e1d802dd5d82e1cd5f5ed3a15cbc75eb6745"
"hash": "c92cc71e6d10c41368f7aa75b0799c2e1e9ca0ed33077ade8d7560e7cc21fa06"
}
@@ -15,7 +15,7 @@
]
},
"nullable": [
true
null
]
},
"hash": "ddf2eccb78a310ed00c7d8b9c3f05d394a7cbcf0038c72a78add5c7b02ef5927"
@@ -0,0 +1,14 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE global_settings SET value = $1 WHERE name = 'teams'",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Jsonb"
]
},
"nullable": []
},
"hash": "e565f3b2e51059f563d18a8a9442bcae9640cee7b936820cb46c011222a77ff0"
}
+66 -65
View File
@@ -713,7 +713,7 @@ dependencies = [
"percent-encoding",
"pin-project-lite",
"tracing",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -1222,15 +1222,16 @@ dependencies = [
[[package]]
name = "blake3"
version = "1.5.5"
version = "1.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b8ee0c1824c4dea5b5f81736aff91bae041d2c07ee1192bec91054e10e3e601e"
checksum = "1230237285e3e10cde447185e8975408ae24deaa67205ce684805c25bc0c7937"
dependencies = [
"arrayref",
"arrayvec",
"cc",
"cfg-if",
"constant_time_eq",
"memmap2",
]
[[package]]
@@ -1668,9 +1669,9 @@ dependencies = [
[[package]]
name = "clap"
version = "4.5.29"
version = "4.5.30"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8acebd8ad879283633b343856142139f2da2317c96b05b4dd6181c61e2480184"
checksum = "92b7b18d71fad5313a1e320fa9897994228ce274b60faa4d694fe0ea89cd9e6d"
dependencies = [
"clap_builder",
"clap_derive",
@@ -1678,9 +1679,9 @@ dependencies = [
[[package]]
name = "clap_builder"
version = "4.5.29"
version = "4.5.30"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f6ba32cbda51c7e1dfd49acc1457ba1a7dec5b64fe360e828acb13ca8dc9c2f9"
checksum = "a35db2071778a7344791a4fb4f95308b5673d219dee3ae348b86642574ecc90c"
dependencies = [
"anstream",
"anstyle",
@@ -2263,7 +2264,7 @@ dependencies = [
"tokio",
"tokio-util",
"url",
"uuid 1.13.1",
"uuid 1.13.2",
"xz2",
"zstd",
]
@@ -2363,7 +2364,7 @@ dependencies = [
"regex",
"sha2 0.10.8",
"unicode-segmentation",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -2530,7 +2531,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bef552e6f588e446098f6ba40d89ac146c8c7b64aade83c051ee00bb5d2bc18d"
dependencies = [
"serde",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -2825,7 +2826,7 @@ dependencies = [
"serde",
"thiserror 1.0.69",
"tokio",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -5405,7 +5406,7 @@ dependencies = [
"sha2 0.10.8",
"subprocess",
"thiserror 1.0.69",
"uuid 1.13.1",
"uuid 1.13.2",
"zstd",
]
@@ -6366,7 +6367,7 @@ dependencies = [
"postgres-protocol 0.6.8",
"serde",
"serde_json",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -7222,7 +7223,7 @@ dependencies = [
"rkyv_derive",
"seahash",
"tinyvec",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -7613,7 +7614,7 @@ dependencies = [
"serde",
"thiserror 1.0.69",
"url",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -7651,7 +7652,7 @@ dependencies = [
"schemars_derive",
"serde",
"serde_json",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -8386,7 +8387,7 @@ dependencies = [
"tokio-stream",
"tracing",
"url",
"uuid 1.13.1",
"uuid 1.13.2",
"webpki-roots",
]
@@ -8470,7 +8471,7 @@ dependencies = [
"stringprep",
"thiserror 2.0.11",
"tracing",
"uuid 1.13.1",
"uuid 1.13.2",
"whoami",
]
@@ -8511,7 +8512,7 @@ dependencies = [
"stringprep",
"thiserror 2.0.11",
"tracing",
"uuid 1.13.1",
"uuid 1.13.2",
"whoami",
]
@@ -8537,7 +8538,7 @@ dependencies = [
"sqlx-core",
"tracing",
"url",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -9184,7 +9185,7 @@ dependencies = [
"tempfile",
"thiserror 1.0.69",
"time",
"uuid 1.13.1",
"uuid 1.13.2",
"winapi",
]
@@ -9297,9 +9298,9 @@ dependencies = [
[[package]]
name = "tempfile"
version = "3.17.0"
version = "3.17.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a40f762a77d2afa88c2d919489e390a12bdd261ed568e60cfa7e48d4e20f0d33"
checksum = "22e5a0acb1f3f55f65cc4a866c361b2fb2a0ff6366785ae6fbb5f85df07ba230"
dependencies = [
"cfg-if",
"fastrand 2.3.0",
@@ -9414,7 +9415,7 @@ dependencies = [
"tokio-rustls 0.24.1",
"tokio-util",
"tracing",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
@@ -10121,9 +10122,9 @@ dependencies = [
[[package]]
name = "tree-sitter-language"
version = "0.1.4"
version = "0.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "38eee4db33814de3d004de9d8d825627ed3320d0989cce0dea30efaf5be4736c"
checksum = "c4013970217383f67b18aef68f6fb2e8d409bc5755227092d32efb0422ba24b8"
[[package]]
name = "triomphe"
@@ -10195,9 +10196,9 @@ checksum = "6af6ae20167a9ece4bcb41af5b80f8a1f1df981f6391189ce00fd257af04126a"
[[package]]
name = "typenum"
version = "1.17.0"
version = "1.18.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "42ff0bf0c66b8238c6f3b578df37d0b7848e55df8577b3f74f92a69acceeb825"
checksum = "1dccffe3ce07af9386bfd29e80c0ab1a8205a2fc34e4bcd40364df902cfa8f3f"
[[package]]
name = "typify"
@@ -10250,7 +10251,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ab82fc73182c29b02e2926a6df32f2241dbadb5cfc111fd595515b3598f46bb3"
dependencies = [
"rand 0.9.0",
"uuid 1.13.1",
"uuid 1.13.2",
"web-time",
]
@@ -10532,9 +10533,9 @@ dependencies = [
[[package]]
name = "uuid"
version = "1.13.1"
version = "1.13.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ced87ca4be083373936a67f8de945faa23b6b42384bd5b64434850802c6dccd0"
checksum = "8c1f41ffb7cf259f1ecc2876861a17e7142e63ead296f671f81f6ae85903e0d6"
dependencies = [
"getrandom 0.3.1",
"serde",
@@ -10858,7 +10859,7 @@ checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f"
[[package]]
name = "windmill"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"axum",
@@ -10887,7 +10888,7 @@ dependencies = [
"tokio",
"tracing",
"url",
"uuid 1.13.1",
"uuid 1.13.2",
"v8",
"windmill-api",
"windmill-api-client",
@@ -10901,7 +10902,7 @@ dependencies = [
[[package]]
name = "windmill-api"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"argon2",
@@ -10982,7 +10983,7 @@ dependencies = [
"ulid",
"url",
"urlencoding",
"uuid 1.13.1",
"uuid 1.13.2",
"windmill-audit",
"windmill-common",
"windmill-git-sync",
@@ -10995,7 +10996,7 @@ dependencies = [
[[package]]
name = "windmill-api-client"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"base64 0.22.1",
"chrono",
@@ -11008,12 +11009,12 @@ dependencies = [
"serde",
"serde_json",
"syn 1.0.109",
"uuid 1.13.1",
"uuid 1.13.2",
]
[[package]]
name = "windmill-audit"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"chrono",
"serde",
@@ -11026,21 +11027,21 @@ dependencies = [
[[package]]
name = "windmill-autoscaling"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"serde",
"serde_json",
"sqlx",
"tracing",
"uuid 1.13.1",
"uuid 1.13.2",
"windmill-common",
"windmill-queue",
]
[[package]]
name = "windmill-common"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"async-stream",
@@ -11093,27 +11094,27 @@ dependencies = [
"tracing-loki",
"tracing-opentelemetry",
"tracing-subscriber",
"uuid 1.13.1",
"uuid 1.13.2",
"windmill-macros",
]
[[package]]
name = "windmill-git-sync"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"regex",
"serde",
"serde_json",
"sqlx",
"tracing",
"uuid 1.13.1",
"uuid 1.13.2",
"windmill-common",
"windmill-queue",
]
[[package]]
name = "windmill-indexer"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"bytes",
@@ -11130,13 +11131,13 @@ dependencies = [
"tokio",
"tokio-tar",
"tracing",
"uuid 1.13.1",
"uuid 1.13.2",
"windmill-common",
]
[[package]]
name = "windmill-macros"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"itertools 0.14.0",
"lazy_static",
@@ -11148,7 +11149,7 @@ dependencies = [
[[package]]
name = "windmill-parser"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"convert_case 0.6.0",
"serde",
@@ -11157,7 +11158,7 @@ dependencies = [
[[package]]
name = "windmill-parser-bash"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"lazy_static",
@@ -11169,7 +11170,7 @@ dependencies = [
[[package]]
name = "windmill-parser-csharp"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"serde_json",
@@ -11181,7 +11182,7 @@ dependencies = [
[[package]]
name = "windmill-parser-go"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"gosyn",
@@ -11193,7 +11194,7 @@ dependencies = [
[[package]]
name = "windmill-parser-graphql"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"lazy_static",
@@ -11205,7 +11206,7 @@ dependencies = [
[[package]]
name = "windmill-parser-php"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"itertools 0.14.0",
@@ -11216,7 +11217,7 @@ dependencies = [
[[package]]
name = "windmill-parser-py"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"itertools 0.14.0",
@@ -11227,7 +11228,7 @@ dependencies = [
[[package]]
name = "windmill-parser-py-imports"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"async-recursion",
@@ -11247,7 +11248,7 @@ dependencies = [
[[package]]
name = "windmill-parser-rust"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"convert_case 0.6.0",
@@ -11264,7 +11265,7 @@ dependencies = [
[[package]]
name = "windmill-parser-sql"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"lazy_static",
@@ -11276,7 +11277,7 @@ dependencies = [
[[package]]
name = "windmill-parser-ts"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"lazy_static",
@@ -11294,7 +11295,7 @@ dependencies = [
[[package]]
name = "windmill-parser-wasm"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"getrandom 0.2.15",
@@ -11316,7 +11317,7 @@ dependencies = [
[[package]]
name = "windmill-parser-yaml"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"serde_json",
@@ -11326,7 +11327,7 @@ dependencies = [
[[package]]
name = "windmill-queue"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"async-recursion",
@@ -11352,14 +11353,14 @@ dependencies = [
"tokio",
"tracing",
"ulid",
"uuid 1.13.1",
"uuid 1.13.2",
"windmill-audit",
"windmill-common",
]
[[package]]
name = "windmill-sql-datatype-parser-wasm"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"wasm-bindgen",
"wasm-bindgen-test",
@@ -11369,7 +11370,7 @@ dependencies = [
[[package]]
name = "windmill-worker"
version = "1.463.3"
version = "1.463.5"
dependencies = [
"anyhow",
"async-recursion",
@@ -11426,7 +11427,7 @@ dependencies = [
"tokio-util",
"tracing",
"urlencoding",
"uuid 1.13.1",
"uuid 1.13.2",
"windmill-audit",
"windmill-common",
"windmill-git-sync",
+3 -3
View File
@@ -1,6 +1,6 @@
[package]
name = "windmill"
version = "1.463.3"
version = "1.463.5"
authors.workspace = true
edition.workspace = true
@@ -30,7 +30,7 @@ members = [
]
[workspace.package]
version = "1.463.3"
version = "1.463.5"
authors = ["Ruben Fiszel <ruben@windmill.dev>"]
edition = "2021"
@@ -174,7 +174,7 @@ uuid = { version = "^1", features = ["serde", "v4"] }
thiserror = "^2"
anyhow = "^1"
chrono = { version = "0.4.35", features = ["serde"] }
chrono-tz = "^0"
chrono-tz = "^0.10.1"
tracing = "^0"
tracing-subscriber = { version = "^0", features = ["env-filter", "json"] }
tracing-appender = "^0"
+1 -1
View File
@@ -1 +1 @@
703a03ac430b6f603a5a189e5b4ff9a42bb2bd7f
5d25cf2cd15c1953794045fd7debea14a33c7519
View File
+11 -2
View File
@@ -21,7 +21,7 @@ use std::{
net::{IpAddr, Ipv4Addr, SocketAddr},
time::Duration,
};
use tokio::{fs::File, io::AsyncReadExt};
use tokio::{fs::File, io::AsyncReadExt, task::JoinHandle};
use uuid::Uuid;
use windmill_api::HTTP_CLIENT;
@@ -372,6 +372,7 @@ async fn windmill_main() -> anyhow::Result<()> {
let is_agent = mode == Mode::Agent;
let mut migration_handle: Option<JoinHandle<()>> = None;
#[cfg(feature = "parquet")]
let disable_s3_store = std::env::var("DISABLE_S3_STORE")
.ok()
@@ -384,7 +385,7 @@ async fn windmill_main() -> anyhow::Result<()> {
if !skip_migration {
// migration code to avoid break
windmill_api::migrate_db(&db).await?;
migration_handle = windmill_api::migrate_db(&db).await?;
} else {
tracing::info!("SKIP_MIGRATION set, skipping db migration...")
}
@@ -682,6 +683,14 @@ Windmill Community Edition {GIT_VERSION}
loop {
tokio::select! {
biased;
Some(_) = async { if let Some(jh) = migration_handle.take() {
tracing::info!("migration job finished");
Some(jh.await)
} else {
None
}} => {
continue;
},
_ = monitor_killpill_rx.recv() => {
tracing::info!("received killpill for monitor job");
break;
+1 -1
View File
@@ -1,7 +1,7 @@
openapi: "3.0.3"
info:
version: 1.463.3
version: 1.463.5
title: Windmill API
contact:
+95 -22
View File
@@ -15,6 +15,7 @@ use sqlx::{
Executor, PgConnection, Pool, Postgres,
};
use tokio::task::JoinHandle;
use windmill_audit::audit_ee::{AuditAuthor, AuditAuthorable};
use windmill_common::{
db::{Authable, Authed},
@@ -170,7 +171,7 @@ impl Migrate for CustomMigrator {
}
}
pub async fn migrate(db: &DB) -> Result<(), Error> {
pub async fn migrate(db: &DB) -> Result<Option<JoinHandle<()>>, Error> {
let migrator = db.acquire().await?;
let mut custom_migrator = CustomMigrator { inner: migrator };
@@ -225,9 +226,10 @@ pub async fn migrate(db: &DB) -> Result<(), Error> {
}
});
if !has_done_migration(db, "v2_finalize_disable_sync_III").await {
let mut jh = None;
if !has_done_migration(db, "v2_finalize_job_completed").await {
let db2 = db.clone();
let _ = tokio::task::spawn(async move {
let v2jh = tokio::task::spawn(async move {
loop {
if !*MIN_VERSION_IS_AT_LEAST_1_461.read().await {
tracing::info!("Waiting for all workers to be at least version 1.461 before applying v2 finalize migration, sleeping for 5s...");
@@ -245,9 +247,10 @@ pub async fn migrate(db: &DB) -> Result<(), Error> {
break;
}
});
jh = Some(v2jh)
}
Ok(())
Ok(jh)
}
async fn fix_flow_versioning_migration(
@@ -373,29 +376,91 @@ async fn v2_finalize(db: &DB) -> Result<(), Error> {
run_windmill_migration!("v2_finalize_disable_sync_III", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_queue DISABLE ROW LEVEL SECURITY;
ALTER TABLE v2_job_completed DISABLE ROW LEVEL SECURITY;
DROP FUNCTION IF EXISTS v2_job_after_update CASCADE;
DROP FUNCTION IF EXISTS v2_job_completed_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_completed_before_update CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_after_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_before_update CASCADE;
DROP FUNCTION IF EXISTS v2_job_runtime_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_runtime_before_update CASCADE;
DROP FUNCTION IF EXISTS v2_job_status_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_status_before_update CASCADE;
DROP VIEW IF EXISTS completed_job, completed_job_view, job, queue, queue_view CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_2", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_completed DISABLE ROW LEVEL SECURITY;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_3", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_after_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_4", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_completed_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_completed_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_5", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_queue_after_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_queue_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_6", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_runtime IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_runtime_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_runtime_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_7", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_status IN ACCESS EXCLUSIVE MODE;
DROP FUNCTION IF EXISTS v2_job_status_before_insert CASCADE;
DROP FUNCTION IF EXISTS v2_job_status_before_update CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_disable_sync_III_8", db, |tx| {
tx.execute(
r#"
DROP VIEW IF EXISTS completed_job, completed_job_view, job, queue, queue_view CASCADE;
"#,
)
.await?;
});
run_windmill_migration!("v2_finalize_job_queue", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_queue IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_queue
DROP COLUMN IF EXISTS __parent_job CASCADE,
DROP COLUMN IF EXISTS __created_by CASCADE,
@@ -434,6 +499,7 @@ async fn v2_finalize(db: &DB) -> Result<(), Error> {
run_windmill_migration!("v2_finalize_job_completed", db, |tx| {
tx.execute(
r#"
LOCK TABLE v2_job_completed IN ACCESS EXCLUSIVE MODE;
ALTER TABLE v2_job_completed
DROP COLUMN IF EXISTS __parent_job CASCADE,
DROP COLUMN IF EXISTS __created_by CASCADE,
@@ -545,8 +611,8 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> {
.await?;
});
run_windmill_migration!("fix_job_index_1", &db, |tx| {
let migration_job_name = "fix_job_completed_index_4";
run_windmill_migration!("fix_job_index_1_II", &db, |tx| {
let migration_job_name = "fix_job_index_1_II";
let mut i = 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_workspace_id_created_at_new_3 ON v2_job (workspace_id, created_at DESC)")
@@ -579,13 +645,20 @@ async fn fix_job_completed_index(db: &DB) -> Result<(), Error> {
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists root_job_index_by_path_2 ON v2_job (workspace_id, runnable_path, created_at desc) WHERE parent_job IS NULL")
sqlx::query!("create index concurrently if not exists ix_job_root_job_index_by_path_2 ON v2_job (workspace_id, runnable_path, created_at desc) WHERE parent_job IS NULL")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("DROP INDEX CONCURRENTLY IF EXISTS root_job_index_by_path_2")
.execute(db)
.await?;
i += 1;
tracing::info!("step {i} of {migration_job_name} migration");
sqlx::query!("create index concurrently if not exists ix_job_created_at ON v2_job (created_at DESC)")
.execute(db)
.await?;
+5 -3
View File
@@ -34,6 +34,7 @@ use http::HeaderValue;
use reqwest::Client;
#[cfg(feature = "oauth2")]
use std::collections::HashMap;
use tokio::task::JoinHandle;
use windmill_common::global_settings::load_value_from_global_settings;
use windmill_common::global_settings::EMAIL_DOMAIN_SETTING;
use windmill_common::worker::HUB_CACHE_DIR;
@@ -641,7 +642,8 @@ async fn openapi_json() -> &'static str {
include_str!("../openapi-deref.json")
}
pub async fn migrate_db(db: &DB) -> anyhow::Result<()> {
db::migrate(db).await?;
Ok(())
pub async fn migrate_db(db: &DB) -> anyhow::Result<Option<JoinHandle<()>>> {
db::migrate(db)
.await
.map_err(|e| anyhow::anyhow!("Error migrating db: {e:#}"))
}
+13 -13
View File
@@ -281,7 +281,7 @@ pub async fn connect_db(
pub async fn connect(
database_url: &str,
max_connections: u32,
_worker_mode: bool,
worker_mode: bool,
) -> Result<sqlx::Pool<sqlx::Postgres>, error::Error> {
use std::time::Duration;
@@ -289,18 +289,18 @@ pub async fn connect(
.min_connections((max_connections / 5).clamp(3, max_connections))
.max_connections(max_connections)
.max_lifetime(Duration::from_secs(30 * 60)) // 30 mins
// .after_connect(move |conn, _| {
// if worker_mode {
// Box::pin(async move {
// // sqlx::query("SET enable_seqscan = OFF;")
// // .execute(conn)
// // .await?;
// Ok(())
// })
// } else {
// Box::pin(async move { Ok(()) })
// }
// })
.after_connect(move |conn, _| {
if worker_mode {
Box::pin(async move {
sqlx::query("SET enable_seqscan = OFF;")
.execute(conn)
.await?;
Ok(())
})
} else {
Box::pin(async move { Ok(()) })
}
})
.connect_with(
sqlx::postgres::PgConnectOptions::from_str(database_url)?.statement_cache_capacity(400),
)
+10 -10
View File
@@ -566,6 +566,9 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
let result_columns = result_columns.as_ref();
let _job_id = queued_job.id;
let (opt_uuid, _duration, _skip_downstream_error_handlers) = (|| async {
// let start = std::time::Instant::now();
let mut tx = db.begin().await?;
let job_id = queued_job.id;
@@ -663,7 +666,7 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
// tracing::error!("Added completed job {:#?}", queued_job);
let mut _skip_downstream_error_handlers = false;
tx = delete_job(tx, &queued_job.workspace_id, job_id).await?;
tx = delete_job(tx, &job_id).await?;
// tracing::error!("3 {:?}", start.elapsed());
if queued_job.is_flow_step {
@@ -858,6 +861,7 @@ pub async fn add_completed_job<T: Serialize + Send + Sync + ValidableJson>(
"inserted completed job: {} (success: {success})",
queued_job.id
);
// tracing::info!("completed job: {:?}", start.elapsed().as_micros());
Ok((None, _duration, _skip_downstream_error_handlers)) as windmill_common::error::Result<(Option<Uuid>, i64, bool)>
})
.retry(
@@ -2597,21 +2601,17 @@ async fn extract_result_from_job_result(
pub async fn delete_job<'c>(
mut tx: Transaction<'c, Postgres>,
w_id: &str,
job_id: Uuid,
job_id: &Uuid,
) -> windmill_common::error::Result<Transaction<'c, Postgres>> {
#[cfg(feature = "prometheus")]
if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) {
QUEUE_DELETE_COUNT.inc();
}
let job_removed = sqlx::query_scalar!(
"DELETE FROM v2_job_queue WHERE workspace_id = $1 AND id = $2 RETURNING 1",
w_id,
job_id
)
.fetch_optional(&mut *tx)
.await;
let job_removed =
sqlx::query_scalar!("DELETE FROM v2_job_queue WHERE id = $1 RETURNING 1", job_id,)
.fetch_optional(&mut *tx)
.await;
if let Err(job_removed) = job_removed {
tracing::error!(
+1 -1
View File
@@ -201,7 +201,7 @@ export async function main({
kind: "rawscript",
rawscript: {
language: api.RawScript.language.BASH,
content: "# let's bloat that bash script, 3.. 2.. 1.. BOOM\n".repeat(25000) + "echo \"$WM_FLOW_JOB_ID\"\n",
content: "# let's bloat that bash script, 3.. 2.. 1.. BOOM\n".repeat(100) + "echo \"$WM_FLOW_JOB_ID\"\n",
},
});
} else {
+6 -1
View File
@@ -39,6 +39,7 @@ async function main({
workspace,
configPath,
workers,
factor
}: {
host: string;
email?: string;
@@ -47,6 +48,7 @@ async function main({
workspace: string;
configPath: string;
workers: number;
factor?: number;
}) {
async function getConfig(configPath: string): Promise<Config> {
if (configPath.startsWith("http")) {
@@ -77,7 +79,7 @@ async function main({
token,
workspace,
kind: benchmark.kind,
jobs: benchmark.jobs,
jobs: benchmark.jobs * (factor ?? 1),
});
if (benchmark.noSave) {
@@ -153,6 +155,9 @@ await new Command()
"Number of workers that are used to run the benchmarks (only affect graph title)",
{ default: 1 }
)
.option("--factor <factor:number>", "Factor to multiply the number of jobs by.", {
default: 1,
})
.action(main)
.command(
"upgrade",
+3 -3
View File
@@ -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.463.3";
export const VERSION = "v1.463.5";
export async function login(email: string, password: string): Promise<string> {
return await windmill.UserService.login({
@@ -25,7 +25,7 @@ async function waitForDeployment(workspace: string, hash: string) {
if (resp.lock !== null) {
return;
}
} catch (err) {}
} catch (err) { }
await sleep(0.5);
}
throw new Error("Script did not deploy in time");
@@ -246,7 +246,7 @@ export const getFlowPayload = (flowPattern: string): api.FlowPreview => {
input_transforms: {},
language: api.RawScript.language.BASH,
type: "rawscript",
content: "# let's bloat that bash script, 3.. 2.. 1.. BOOM\n".repeat(25000) + "echo \"$WM_FLOW_JOB_ID\"\n",
content: "# let's bloat that bash script, 3.. 2.. 1.. BOOM\n".repeat(100) + `if [[ -z $\{WM_FLOW_JOB_ID+x\} ]]; then\necho "not set"\nelif [[ -z "$WM_FLOW_JOB_ID" ]]; then\necho "empty"\nelse\necho "$WM_FLOW_JOB_ID"\nfi`,
},
}
],
+1 -1
View File
@@ -62,7 +62,7 @@ export {
// }
// });
export const VERSION = "1.463.3";
export const VERSION = "1.463.5";
const command = new Command()
.name("wmill")
+2 -2
View File
@@ -1,12 +1,12 @@
{
"name": "windmill-components",
"version": "1.463.3",
"version": "1.463.5",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "windmill-components",
"version": "1.463.3",
"version": "1.463.5",
"license": "AGPL-3.0",
"dependencies": {
"@anthropic-ai/sdk": "^0.32.1",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "windmill-components",
"version": "1.463.3",
"version": "1.463.5",
"scripts": {
"dev": "vite dev",
"build": "vite build",
+2 -2
View File
@@ -4,8 +4,8 @@ verify_ssl = true
name = "pypi"
[packages]
wmill = ">=1.463.3"
wmill_pg = ">=1.463.3"
wmill = ">=1.463.5"
wmill_pg = ">=1.463.5"
sendgrid = "*"
mysql-connector-python = "*"
pymongo = "*"
+1 -1
View File
@@ -1,7 +1,7 @@
openapi: "3.0.3"
info:
version: 1.463.3
version: 1.463.5
title: OpenFlow Spec
contact:
name: Ruben Fiszel
@@ -12,7 +12,7 @@
RootModule = 'WindmillClient.psm1'
# Version number of this module.
ModuleVersion = '1.463.3'
ModuleVersion = '1.463.5'
# Supported PSEditions
# CompatiblePSEditions = @()
+1 -1
View File
@@ -1,6 +1,6 @@
[tool.poetry]
name = "wmill"
version = "1.463.3"
version = "1.463.5"
description = "A client library for accessing Windmill server wrapping the Windmill client API"
license = "Apache-2.0"
homepage = "https://windmill.dev"
+1 -1
View File
@@ -1,6 +1,6 @@
[tool.poetry]
name = "wmill-pg"
version = "1.463.3"
version = "1.463.5"
description = "An extension client for the wmill client library focused on pg"
license = "Apache-2.0"
homepage = "https://windmill.dev"
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@windmill/windmill",
"version": "1.463.3",
"version": "1.463.5",
"exports": "./src/index.ts",
"publish": {
"exclude": ["!src", "./s3Types.ts", "./client.ts"]
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "windmill-client",
"description": "Windmill SDK client for browsers and Node.js",
"version": "1.463.3",
"version": "1.463.5",
"author": "Ruben Fiszel",
"license": "Apache 2.0",
"devDependencies": {
+1 -1
View File
@@ -1 +1 @@
1.463.3
1.463.5