From dfcdb2d8ff6dc1ba299d8883b9a5bd18edece2c2 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Wed, 13 Sep 2023 21:24:54 +0200 Subject: [PATCH] feat: worker groups admin panel (#2277) * merge * merge * merge * wg * progress * all * all * all * all * all * all * fix * fix --- ...cf6473960377455b5fd7ac3b578a1d36c0cc6.json | 19 + ...7bad6882c86e9422b715844aa92b67ed05174.json | 22 ++ ...915ae2734b90b8c4cd35fe37783b1d4dd0b0.json} | 12 +- ...15280f32b70a617fde87f70ea53ca9ade39f.json} | 5 +- ...9148a75e1bd130e7614463487e2ba6957dfdf.json | 17 - ...edbc78f0bf3f815676d8b343152bfce2dfc4c.json | 22 ++ ...15bf29ba5150cd59732288ceb72ce9cecf987.json | 26 ++ ...bb9e132d310d19ca744531a42ca2c7ac56f56.json | 15 + backend/Cargo.lock | 2 + .../20230908155455_add_worker_group.down.sql | 4 + .../20230908155455_add_worker_group.up.sql | 8 + backend/src/main.rs | 99 +++-- backend/src/monitor.rs | 67 +++- backend/tests/worker.rs | 2 + backend/windmill-api/openapi.yaml | 62 +++ backend/windmill-api/src/groups.rs | 2 +- backend/windmill-api/src/jobs.rs | 26 +- backend/windmill-api/src/lib.rs | 6 +- backend/windmill-api/src/workers.rs | 141 ++++--- backend/windmill-api/src/workspaces.rs | 18 +- backend/windmill-common/Cargo.toml | 2 + .../windmill-common/src/global_settings.rs | 6 +- backend/windmill-common/src/lib.rs | 25 +- backend/windmill-common/src/worker.rs | 200 ++++++++++ backend/windmill-queue/src/jobs.rs | 71 +--- backend/windmill-worker/src/common.rs | 2 +- backend/windmill-worker/src/config.rs | 0 backend/windmill-worker/src/lib.rs | 2 +- backend/windmill-worker/src/worker.rs | 95 +++-- frontend/src/lib/components/PageHeader.svelte | 4 +- .../src/lib/components/ScriptBuilder.svelte | 4 +- frontend/src/lib/components/Tooltip.svelte | 19 +- .../src/lib/components/WorkspaceGroup.svelte | 246 ++++++++++++ .../components/common/button/Button.svelte | 8 +- .../flows/header/FlowPreviewButtons.svelte | 2 + frontend/src/lib/utils.ts | 27 +- .../(root)/(logged)/workers/+page.svelte | 365 +++++++++++++----- 37 files changed, 1273 insertions(+), 380 deletions(-) create mode 100644 backend/.sqlx/query-03c7f098ad795d216d58ded0bf4cf6473960377455b5fd7ac3b578a1d36c0cc6.json create mode 100644 backend/.sqlx/query-210b7fa50246d9b8fd1a24ade787bad6882c86e9422b715844aa92b67ed05174.json rename backend/.sqlx/{query-4f6b3b472b4b78c0325cf3755f9ef1806d2e82328ceccbeade8cc2333c6dfe47.json => query-240ce8c9b5c7530999642190c6f7915ae2734b90b8c4cd35fe37783b1d4dd0b0.json} (78%) rename backend/.sqlx/{query-07551a32c49da8c0693dd39c6a63b5b2a596ccc0e52e8918160604a5e133dd32.json => query-47beea5cd6324b53bfb349665fb215280f32b70a617fde87f70ea53ca9ade39f.json} (62%) delete mode 100644 backend/.sqlx/query-61e6aac871b482b6e36f866b4ec9148a75e1bd130e7614463487e2ba6957dfdf.json create mode 100644 backend/.sqlx/query-6c0136f7965f1e01620a7d5efd7edbc78f0bf3f815676d8b343152bfce2dfc4c.json create mode 100644 backend/.sqlx/query-7998b23eb72f5967a0fd376fa1015bf29ba5150cd59732288ceb72ce9cecf987.json create mode 100644 backend/.sqlx/query-903f2f62f3829274f5dfa0cb1ebbb9e132d310d19ca744531a42ca2c7ac56f56.json create mode 100644 backend/migrations/20230908155455_add_worker_group.down.sql create mode 100644 backend/migrations/20230908155455_add_worker_group.up.sql create mode 100644 backend/windmill-common/src/worker.rs create mode 100644 backend/windmill-worker/src/config.rs create mode 100644 frontend/src/lib/components/WorkspaceGroup.svelte diff --git a/backend/.sqlx/query-03c7f098ad795d216d58ded0bf4cf6473960377455b5fd7ac3b578a1d36c0cc6.json b/backend/.sqlx/query-03c7f098ad795d216d58ded0bf4cf6473960377455b5fd7ac3b578a1d36c0cc6.json new file mode 100644 index 0000000000..2f9eb8c94f --- /dev/null +++ b/backend/.sqlx/query-03c7f098ad795d216d58ded0bf4cf6473960377455b5fd7ac3b578a1d36c0cc6.json @@ -0,0 +1,19 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO worker_ping (worker_instance, worker, ip, custom_tags, worker_group, dedicated_worker) VALUES ($1, $2, $3, $4, $5, $6) ON CONFLICT (worker) DO UPDATE set ip = $3, custom_tags = $4, worker_group = $5", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Varchar", + "Varchar", + "TextArray", + "Varchar", + "Varchar" + ] + }, + "nullable": [] + }, + "hash": "03c7f098ad795d216d58ded0bf4cf6473960377455b5fd7ac3b578a1d36c0cc6" +} diff --git a/backend/.sqlx/query-210b7fa50246d9b8fd1a24ade787bad6882c86e9422b715844aa92b67ed05174.json b/backend/.sqlx/query-210b7fa50246d9b8fd1a24ade787bad6882c86e9422b715844aa92b67ed05174.json new file mode 100644 index 0000000000..ae14447a5e --- /dev/null +++ b/backend/.sqlx/query-210b7fa50246d9b8fd1a24ade787bad6882c86e9422b715844aa92b67ed05174.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM worker_group_config WHERE name = $1 RETURNING name", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "name", + "type_info": "Varchar" + } + ], + "parameters": { + "Left": [ + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "210b7fa50246d9b8fd1a24ade787bad6882c86e9422b715844aa92b67ed05174" +} diff --git a/backend/.sqlx/query-4f6b3b472b4b78c0325cf3755f9ef1806d2e82328ceccbeade8cc2333c6dfe47.json b/backend/.sqlx/query-240ce8c9b5c7530999642190c6f7915ae2734b90b8c4cd35fe37783b1d4dd0b0.json similarity index 78% rename from backend/.sqlx/query-4f6b3b472b4b78c0325cf3755f9ef1806d2e82328ceccbeade8cc2333c6dfe47.json rename to backend/.sqlx/query-240ce8c9b5c7530999642190c6f7915ae2734b90b8c4cd35fe37783b1d4dd0b0.json index e7276c191a..d6bbf3583d 100644 --- a/backend/.sqlx/query-4f6b3b472b4b78c0325cf3755f9ef1806d2e82328ceccbeade8cc2333c6dfe47.json +++ b/backend/.sqlx/query-240ce8c9b5c7530999642190c6f7915ae2734b90b8c4cd35fe37783b1d4dd0b0.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT worker, worker_instance, EXTRACT(EPOCH FROM (now() - ping_at))::integer as last_ping, started_at, ip, jobs_executed, custom_tags FROM worker_ping ORDER BY ping_at desc LIMIT $1 OFFSET $2", + "query": "SELECT worker, worker_instance, EXTRACT(EPOCH FROM (now() - ping_at))::integer as last_ping, started_at, ip, jobs_executed, custom_tags, worker_group FROM worker_ping ORDER BY ping_at desc LIMIT $1 OFFSET $2", "describe": { "columns": [ { @@ -37,6 +37,11 @@ "ordinal": 6, "name": "custom_tags", "type_info": "TextArray" + }, + { + "ordinal": 7, + "name": "worker_group", + "type_info": "Varchar" } ], "parameters": { @@ -52,8 +57,9 @@ false, false, false, - true + true, + false ] }, - "hash": "4f6b3b472b4b78c0325cf3755f9ef1806d2e82328ceccbeade8cc2333c6dfe47" + "hash": "240ce8c9b5c7530999642190c6f7915ae2734b90b8c4cd35fe37783b1d4dd0b0" } diff --git a/backend/.sqlx/query-07551a32c49da8c0693dd39c6a63b5b2a596ccc0e52e8918160604a5e133dd32.json b/backend/.sqlx/query-47beea5cd6324b53bfb349665fb215280f32b70a617fde87f70ea53ca9ade39f.json similarity index 62% rename from backend/.sqlx/query-07551a32c49da8c0693dd39c6a63b5b2a596ccc0e52e8918160604a5e133dd32.json rename to backend/.sqlx/query-47beea5cd6324b53bfb349665fb215280f32b70a617fde87f70ea53ca9ade39f.json index 0bef008c89..59f45509b9 100644 --- a/backend/.sqlx/query-07551a32c49da8c0693dd39c6a63b5b2a596ccc0e52e8918160604a5e133dd32.json +++ b/backend/.sqlx/query-47beea5cd6324b53bfb349665fb215280f32b70a617fde87f70ea53ca9ade39f.json @@ -1,15 +1,16 @@ { "db_name": "PostgreSQL", - "query": "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1 WHERE worker = $2", + "query": "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2 WHERE worker = $3", "describe": { "columns": [], "parameters": { "Left": [ "Int4", + "TextArray", "Text" ] }, "nullable": [] }, - "hash": "07551a32c49da8c0693dd39c6a63b5b2a596ccc0e52e8918160604a5e133dd32" + "hash": "47beea5cd6324b53bfb349665fb215280f32b70a617fde87f70ea53ca9ade39f" } diff --git a/backend/.sqlx/query-61e6aac871b482b6e36f866b4ec9148a75e1bd130e7614463487e2ba6957dfdf.json b/backend/.sqlx/query-61e6aac871b482b6e36f866b4ec9148a75e1bd130e7614463487e2ba6957dfdf.json deleted file mode 100644 index 6152229295..0000000000 --- a/backend/.sqlx/query-61e6aac871b482b6e36f866b4ec9148a75e1bd130e7614463487e2ba6957dfdf.json +++ /dev/null @@ -1,17 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "INSERT INTO worker_ping (worker_instance, worker, ip, custom_tags) VALUES ($1, $2, $3, $4) ON CONFLICT (worker) DO NOTHING", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Varchar", - "Varchar", - "Varchar", - "TextArray" - ] - }, - "nullable": [] - }, - "hash": "61e6aac871b482b6e36f866b4ec9148a75e1bd130e7614463487e2ba6957dfdf" -} diff --git a/backend/.sqlx/query-6c0136f7965f1e01620a7d5efd7edbc78f0bf3f815676d8b343152bfce2dfc4c.json b/backend/.sqlx/query-6c0136f7965f1e01620a7d5efd7edbc78f0bf3f815676d8b343152bfce2dfc4c.json new file mode 100644 index 0000000000..9796d12680 --- /dev/null +++ b/backend/.sqlx/query-6c0136f7965f1e01620a7d5efd7edbc78f0bf3f815676d8b343152bfce2dfc4c.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT config FROM worker_group_config WHERE name = $1", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "config", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [ + "Text" + ] + }, + "nullable": [ + true + ] + }, + "hash": "6c0136f7965f1e01620a7d5efd7edbc78f0bf3f815676d8b343152bfce2dfc4c" +} diff --git a/backend/.sqlx/query-7998b23eb72f5967a0fd376fa1015bf29ba5150cd59732288ceb72ce9cecf987.json b/backend/.sqlx/query-7998b23eb72f5967a0fd376fa1015bf29ba5150cd59732288ceb72ce9cecf987.json new file mode 100644 index 0000000000..a3846f69f7 --- /dev/null +++ b/backend/.sqlx/query-7998b23eb72f5967a0fd376fa1015bf29ba5150cd59732288ceb72ce9cecf987.json @@ -0,0 +1,26 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT * FROM worker_group_config", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "name", + "type_info": "Varchar" + }, + { + "ordinal": 1, + "name": "config", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + false, + true + ] + }, + "hash": "7998b23eb72f5967a0fd376fa1015bf29ba5150cd59732288ceb72ce9cecf987" +} diff --git a/backend/.sqlx/query-903f2f62f3829274f5dfa0cb1ebbb9e132d310d19ca744531a42ca2c7ac56f56.json b/backend/.sqlx/query-903f2f62f3829274f5dfa0cb1ebbb9e132d310d19ca744531a42ca2c7ac56f56.json new file mode 100644 index 0000000000..e24ca9042f --- /dev/null +++ b/backend/.sqlx/query-903f2f62f3829274f5dfa0cb1ebbb9e132d310d19ca744531a42ca2c7ac56f56.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO worker_group_config (name, config) VALUES ($1, $2) ON CONFLICT (name) DO UPDATE SET config = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar", + "Jsonb" + ] + }, + "nullable": [] + }, + "hash": "903f2f62f3829274f5dfa0cb1ebbb9e132d310d19ca744531a42ca2c7ac56f56" +} diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 7bb40805a3..b053d84c24 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -7212,9 +7212,11 @@ dependencies = [ "hex", "hmac", "hyper", + "itertools 0.10.5", "lazy_static", "prometheus", "rand 0.8.5", + "regex", "reqwest", "serde", "serde_json", diff --git a/backend/migrations/20230908155455_add_worker_group.down.sql b/backend/migrations/20230908155455_add_worker_group.down.sql new file mode 100644 index 0000000000..298b72a264 --- /dev/null +++ b/backend/migrations/20230908155455_add_worker_group.down.sql @@ -0,0 +1,4 @@ +-- Add down migration script here +DROP TABLE worker_group_config; + +ALTER TABLE worker_ping DROP COLUMN worker_group; \ No newline at end of file diff --git a/backend/migrations/20230908155455_add_worker_group.up.sql b/backend/migrations/20230908155455_add_worker_group.up.sql new file mode 100644 index 0000000000..6e8325fc2e --- /dev/null +++ b/backend/migrations/20230908155455_add_worker_group.up.sql @@ -0,0 +1,8 @@ +-- Add up migration script here +CREATE TABLE worker_group_config ( + name VARCHAR(255) PRIMARY KEY, + config JSONB DEFAULT '{}'::jsonb +); + +ALTER TABLE worker_ping ADD COLUMN IF NOT EXISTS worker_group VARCHAR(255) NOT NULL DEFAULT 'default'; +ALTER TABLE worker_ping ADD COLUMN IF NOT EXISTS dedicated_worker VARCHAR(255); \ No newline at end of file diff --git a/backend/src/main.rs b/backend/src/main.rs index 823bb0e329..f6cbd90bb6 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -8,15 +8,14 @@ use gethostname::gethostname; use git_version::git_version; -use monitor::handle_zombie_jobs_periodically; use sqlx::{Pool, Postgres}; use std::{ net::{IpAddr, Ipv4Addr, SocketAddr}, sync::Arc, + time::Duration, }; use tokio::{ fs::{metadata, DirBuilder}, - join, sync::RwLock, }; use windmill_api::{LICENSE_KEY, OAUTH_CLIENTS, SMTP_CLIENT}; @@ -28,6 +27,8 @@ use windmill_worker::{ PIP_CACHE_DIR, ROOT_TMP_CACHE_DIR, TAR_PIP_TMP_CACHE_DIR, }; +use crate::monitor::monitor_db; + const GIT_VERSION: &str = git_version!(args = ["--tag", "--always"], fallback = "unknown-version"); const DEFAULT_NUM_WORKERS: usize = 1; const DEFAULT_PORT: u16 = 8000; @@ -135,7 +136,10 @@ Windmill Community Edition {GIT_VERSION} tracing::info!("Smtp client connected."); } } - if server_mode || num_workers > 0 { + + let worker_mode = num_workers > 0; + + if server_mode || worker_mode { let port_var = std::env::var("PORT").ok().and_then(|x| x.parse().ok()); let port = if server_mode { @@ -144,6 +148,19 @@ Windmill Community Edition {GIT_VERSION} port_var.unwrap_or(0) }; + // since it's only on server mode, the port is statically defined + let base_internal_url: String = format!("http://localhost:{}", port.to_string()); + + monitor_db( + &db, + tx.clone(), + &base_internal_url, + rsmq.clone(), + worker_mode, + server_mode, + ) + .await; + if std::env::var("BASE_INTERNAL_URL").is_ok() { tracing::warn!("BASE_INTERNAL_URL is now unecessary and ignored, you can remove it."); } @@ -161,10 +178,11 @@ Windmill Community Edition {GIT_VERSION} let workers_f = async { let port = port_rx.await?; let base_internal_url: String = format!("http://localhost:{}", port.to_string()); - if num_workers > 0 { + if worker_mode { run_workers( db.clone(), rx.resubscribe(), + tx.clone(), num_workers, base_internal_url.clone(), rsmq.clone(), @@ -176,29 +194,48 @@ Windmill Community Edition {GIT_VERSION} Ok(()) as anyhow::Result<()> }; - let rsmq2 = rsmq.clone(); let monitor_f = async { - if server_mode { - // since it's only on server mode, the port is statically defined - let base_internal_url: String = format!("http://localhost:{}", port.to_string()); - monitor_db(&db, rx.resubscribe(), &base_internal_url, rsmq2).await; - } + let db = db.clone(); + let tx = tx.clone(); + let rsmq = rsmq.clone(); + + let mut rx = rx.resubscribe(); + let base_internal_url = base_internal_url.to_string(); + tokio::spawn(async move { + //monitor_db is applied at start, no need to apply it twice + tokio::time::sleep(Duration::from_secs(30)).await; + loop { + monitor_db( + &db, + tx.clone(), + &base_internal_url, + rsmq.clone(), + worker_mode, + server_mode, + ) + .await; + tokio::select! { + _ = tokio::time::sleep(Duration::from_secs(30)) => (), + _ = rx.recv() => { + println!("received killpill for monitor job"); + break; + } + } + } + }); + Ok(()) as anyhow::Result<()> }; let metrics_f = async { - match metrics_addr { - Some(_addr) => { - #[cfg(not(feature = "enterprise"))] - panic!("Metrics are only available in the Enterprise Edition"); + if let Some(_addr) = metrics_addr { + #[cfg(not(feature = "enterprise"))] + panic!("Metrics are only available in the Enterprise Edition"); - #[cfg(feature = "enterprise")] - windmill_common::serve_metrics(_addr, rx.resubscribe(), num_workers > 0) - .await - .map_err(anyhow::Error::from) - } - None => Ok(()), + #[cfg(feature = "enterprise")] + windmill_common::serve_metrics(_addr, rx.resubscribe(), num_workers > 0).await; } + Ok(()) as anyhow::Result<()> }; futures::try_join!(shutdown_signal, server_f, metrics_f, workers_f, monitor_f)?; @@ -225,28 +262,10 @@ fn display_config(envs: &[&str]) { ) } -pub async fn monitor_db( - db: &Pool, - rx: tokio::sync::broadcast::Receiver<()>, - base_internal_url: &str, - rsmq: Option, -) -> tokio::task::JoinHandle<()> { - let db1 = db.clone(); - let db2 = db.clone(); - - let rx2 = rx.resubscribe(); - let base_internal_url = base_internal_url.to_string(); - tokio::spawn(async move { - join!( - handle_zombie_jobs_periodically(&db1, rx, &base_internal_url, rsmq), - windmill_api::delete_expired_items_perdiodically(&db2, rx2) - ); - }) -} - pub async fn run_workers( db: Pool, rx: tokio::sync::broadcast::Receiver<()>, + tx: tokio::sync::broadcast::Sender<()>, num_workers: i32, base_internal_url: String, rsmq: Option, @@ -320,6 +339,7 @@ pub async fn run_workers( +pub async fn monitor_db( db: &Pool, - mut rx: tokio::sync::broadcast::Receiver<()>, + tx: tokio::sync::broadcast::Sender<()>, base_internal_url: &str, rsmq: Option, + worker_mode: bool, + server_mode: bool, ) { - loop { - handle_zombie_jobs(db, base_internal_url, rsmq.clone()).await; - - tokio::select! { - _ = tokio::time::sleep(Duration::from_secs(30)) => (), - _ = rx.recv() => { - println!("received killpill for monitor job"); - break; + let zombie_jobs_f = async { + if server_mode { + handle_zombie_jobs(db, base_internal_url, rsmq.clone()).await; + } + }; + let expired_items_f = async { + if server_mode { + windmill_api::delete_expired_items(&db).await; + } + }; + let reload_worker_config_f = async { + if worker_mode { + reload_worker_config(&db, tx).await; + } + }; + let reload_custom_tags_f = async { + if server_mode { + if let Err(e) = reload_custom_tags_setting(db).await { + tracing::error!("Error reloading custom tags: {:?}", e) } } + }; + join!( + expired_items_f, + zombie_jobs_f, + reload_worker_config_f, + reload_custom_tags_f + ); +} + +pub async fn reload_worker_config(db: &Pool, tx: tokio::sync::broadcast::Sender<()>) { + let config = load_worker_config(&db).await; + if let Err(e) = config { + tracing::error!("Error reloading worker config: {:?}", e) + } else { + let wc = WORKER_CONFIG.read().await; + let config = config.unwrap(); + if *wc != config { + if (*wc).dedicated_worker != config.dedicated_worker { + tracing::info!("Dedicated worker config changed, sending killpill. Expecting to be restarted by supervisor."); + let _ = tx.send(()); + } + drop(wc); + + let mut wc = WORKER_CONFIG.write().await; + tracing::info!("Reloading worker config..."); + *wc = config + } } } diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index ec8daa4f0d..30bbd861bb 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -956,6 +956,7 @@ fn spawn_test_worker( let worker_name: String = next_worker_name(); let ip: &str = Default::default(); + let tx2 = tx.clone(); let future = async move { let base_internal_url = format!("http://localhost:{}", port); windmill_worker::run_worker::( @@ -966,6 +967,7 @@ fn spawn_test_worker( 1, ip, rx, + tx2, &base_internal_url, None, Arc::new(RwLock::new(None)), diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index e9ce067f64..da6db7bfee 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -5252,6 +5252,65 @@ paths: items: $ref: "#/components/schemas/WorkerPing" + /workers/list_worker_groups: + get: + summary: list workers + operationId: listWorkerGroups + tags: + - worker + responses: + "200": + description: a list of workers + content: + application/json: + schema: + type: array + items: + type: object + properties: + name: + type: string + config: {} + required: + - name + - config + + /workers/worker_group/{name}: + post: + summary: Update Worker Group + operationId: updateWorkerGroup + tags: + - worker + parameters: + - $ref: "#/components/parameters/Name" + requestBody: + description: worker group + required: true + content: + application/json: + schema: {} + responses: + "200": + description: Update a worker group + content: + text/plain: + schema: + type: string + delete: + summary: Delete Worker Group + operationId: deleteWorkerGroup + tags: + - worker + parameters: + - $ref: "#/components/parameters/Name" + responses: + "200": + description: Delete a worker group + content: + text/plain: + schema: + type: string + /w/{workspace}/acls/get/{kind}/{path}: get: summary: get granular acls @@ -7172,6 +7231,8 @@ components: type: array items: type: string + worker_group: + type: string required: - worker - worker_instance @@ -7179,6 +7240,7 @@ components: - started_at - ip - jobs_executed + - worker_group UserWorkspaceList: type: object diff --git a/backend/windmill-api/src/groups.rs b/backend/windmill-api/src/groups.rs index 580ffce75c..29f364b55e 100644 --- a/backend/windmill-api/src/groups.rs +++ b/backend/windmill-api/src/groups.rs @@ -15,6 +15,7 @@ use axum::{ Json, Router, }; use windmill_audit::{audit_log, ActionKind}; +use windmill_common::worker::CLOUD_HOSTED; use windmill_common::{db::UserDB, users::username_to_permissioned_as}; use windmill_common::{ error::{Error, JsonResult, Result}, @@ -23,7 +24,6 @@ use windmill_common::{ use serde::{Deserialize, Serialize}; use sqlx::{query_scalar, FromRow, Postgres, Transaction}; -use windmill_queue::CLOUD_HOSTED; pub fn workspaced_service() -> Router { Router::new() diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index b6fbb72260..7e1d505fd4 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -13,7 +13,6 @@ use crate::{ users::{check_scopes, require_owner_of_path, OptAuthed}, utils::require_super_admin, variables::get_workspace_key, - workers::{CUSTOM_TAGS, CUSTOM_TAGS_PER_WORKSPACE}, BASE_URL, }; use anyhow::Context; @@ -34,6 +33,7 @@ use sqlx::{query_scalar, types::Uuid, FromRow, Postgres, Transaction}; use tower_http::cors::{Any, CorsLayer}; use urlencoding::encode; use windmill_audit::{audit_log, ActionKind}; +use windmill_common::worker::CUSTOM_TAGS_PER_WORKSPACE; use windmill_common::{ db::UserDB, error::{self, to_anyhow, Error}, @@ -1607,12 +1607,12 @@ fn add_raw_string( return args; } -fn check_tag_available_for_workspace(w_id: &str, tag: &Option) -> error::Result<()> { +async fn check_tag_available_for_workspace(w_id: &str, tag: &Option) -> error::Result<()> { if let Some(tag) = tag { if tag == "" { return Ok(()); } - let custom_tags_per_w = &*CUSTOM_TAGS_PER_WORKSPACE; + let custom_tags_per_w = CUSTOM_TAGS_PER_WORKSPACE.read().await; if custom_tags_per_w.0.contains(&tag.to_string()) { Ok(()) } else if custom_tags_per_w.1.contains_key(tag) @@ -1626,7 +1626,7 @@ fn check_tag_available_for_workspace(w_id: &str, tag: &Option) -> error: } else { return Err(error::Error::BadRequest(format!( "Tag {tag} cannot be used on workspace {w_id}: (CUSTOM_TAGS: {:?})", - *CUSTOM_TAGS + custom_tags_per_w ))); } } else { @@ -1655,7 +1655,7 @@ pub async fn run_flow_by_path( .fetch_optional(&db) .await? .flatten(); - check_tag_available_for_workspace(&w_id, &tag)?; + check_tag_available_for_workspace(&w_id, &tag).await?; let scheduled_for = run_query.get_scheduled_for(&db).await?; let args = run_query.add_include_headers(headers, args.unwrap_or_default()); let args = add_raw_string(raw_string, args); @@ -1705,7 +1705,7 @@ pub async fn run_job_by_path( let args = run_query.add_include_headers(headers, args.unwrap_or_default()); let args = add_raw_string(raw_string, args); - check_tag_available_for_workspace(&w_id, &tag)?; + check_tag_available_for_workspace(&w_id, &tag).await?; let tx = PushIsolationLevel::Isolated(user_db.clone(), authed.clone().into(), rsmq); let (uuid, tx) = push( @@ -1908,7 +1908,7 @@ pub async fn run_wait_result_job_by_path_get( check_scopes(&authed, || format!("run:script/{script_path}"))?; let (job_payload, tag) = script_path_to_payload(script_path, &db, &w_id).await?; - check_tag_available_for_workspace(&w_id, &tag)?; + check_tag_available_for_workspace(&w_id, &tag).await?; let tx = PushIsolationLevel::Isolated(user_db.clone(), authed.clone().into(), rsmq); let (uuid, tx) = push( @@ -2079,7 +2079,7 @@ async fn run_wait_result_script_by_path_internal( let args = run_query.add_include_headers(headers, args.unwrap_or_default()); let args = add_raw_string(raw_string, args); - check_tag_available_for_workspace(&w_id, &tag)?; + check_tag_available_for_workspace(&w_id, &tag).await?; let tx = PushIsolationLevel::Isolated(user_db.clone(), authed.clone().into(), rsmq); let (uuid, tx) = push( @@ -2143,7 +2143,7 @@ pub async fn run_wait_result_script_by_hash( let args = run_query.add_include_headers(headers, args.unwrap_or_default()); let args = add_raw_string(raw_string, args); - check_tag_available_for_workspace(&w_id, &tag)?; + check_tag_available_for_workspace(&w_id, &tag).await?; let tx = PushIsolationLevel::Isolated(user_db.clone(), authed.clone().into(), rsmq); let (uuid, tx) = push( @@ -2258,7 +2258,7 @@ async fn run_wait_result_flow_by_path_internal( .fetch_optional(&db) .await? .flatten(); - check_tag_available_for_workspace(&w_id, &tag)?; + check_tag_available_for_workspace(&w_id, &tag).await?; let tx = PushIsolationLevel::Isolated(user_db.clone(), authed.clone().into(), rsmq); let (uuid, tx) = push( @@ -2314,7 +2314,7 @@ async fn run_preview_job( } let scheduled_for = run_query.get_scheduled_for(&db).await?; let args = run_query.add_include_headers(headers, preview.args.unwrap_or_default()); - check_tag_available_for_workspace(&w_id, &preview.tag)?; + check_tag_available_for_workspace(&w_id, &preview.tag).await?; let tx = PushIsolationLevel::Isolated(user_db.clone(), authed.clone().into(), rsmq); let (uuid, tx) = push( @@ -2422,7 +2422,7 @@ async fn run_preview_flow_job( } let scheduled_for = run_query.get_scheduled_for(&db).await?; let args = run_query.add_include_headers(headers, raw_flow.args.unwrap_or_default()); - check_tag_available_for_workspace(&w_id, &raw_flow.tag)?; + check_tag_available_for_workspace(&w_id, &raw_flow.tag).await?; let tx = PushIsolationLevel::Isolated(user_db.clone(), authed.clone().into(), rsmq); let (uuid, tx) = push( @@ -2479,7 +2479,7 @@ pub async fn run_job_by_hash( let scheduled_for = run_query.get_scheduled_for(&db).await?; let args = run_query.add_include_headers(headers, args.unwrap_or_default()); let args = add_raw_string(raw_string, args); - check_tag_available_for_workspace(&w_id, &tag)?; + check_tag_available_for_workspace(&w_id, &tag).await?; let tx = PushIsolationLevel::Isolated(user_db.clone(), authed.clone().into(), rsmq); let (uuid, tx) = push( diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index 1be5fe6c95..f76800fb4a 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -11,7 +11,6 @@ use crate::oauth2::AllClients; use crate::saml::{SamlSsoLogin, ServiceProviderExt}; use crate::scim::has_scim_token; use crate::tracing_init::MyOnFailure; -use crate::workers::ALL_TAGS; use crate::{ oauth2::{build_oauth_clients, SlackVerifier}, tracing_init::{MyMakeSpan, MyOnResponse}, @@ -36,6 +35,7 @@ use tower_http::{ }; use windmill_common::db::UserDB; use windmill_common::utils::rd_string; +use windmill_common::worker::ALL_TAGS; use windmill_common::error::AppError; @@ -72,7 +72,7 @@ mod workspaces; pub const GIT_VERSION: &str = git_version!(args = ["--tag", "--always"], fallback = "unknown-version"); -pub use users::delete_expired_items_perdiodically; +pub use users::delete_expired_items; pub const DEFAULT_BODY_LIMIT: usize = 2097152; // 2MB lazy_static::lazy_static! { @@ -153,7 +153,7 @@ pub async fn run_server( port_tx: tokio::sync::oneshot::Sender, ) -> anyhow::Result<()> { if let Some(mut rsmq) = rsmq.clone() { - for tag in ALL_TAGS.clone() { + for tag in ALL_TAGS.read().await.iter() { let r = rsmq_async::RsmqConnection::create_queue(&mut rsmq, &tag, None, None, None).await; if let Err(e) = r { diff --git a/backend/windmill-api/src/workers.rs b/backend/windmill-api/src/workers.rs index 44d6e3c3b5..fe642ca781 100644 --- a/backend/windmill-api/src/workers.rs +++ b/backend/windmill-api/src/workers.rs @@ -7,74 +7,44 @@ */ use axum::{ - extract::{Extension, Query}, + extract::{Extension, Path, Query}, routing::get, Json, Router, }; -use itertools::Itertools; -use regex::Regex; use serde::{Deserialize, Serialize}; use sqlx::FromRow; use windmill_common::{ db::UserDB, - error::JsonResult, + error::{self, JsonResult}, utils::{paginate, Pagination}, + worker::ALL_TAGS, + DB, }; -use std::collections::HashMap; #[cfg(feature = "benchmark")] use std::sync::atomic::Ordering; #[cfg(feature = "benchmark")] use windmill_queue::IDLE_WORKERS; -use crate::db::ApiAuthed; +use crate::{db::ApiAuthed, utils::require_super_admin}; -#[cfg(not(feature = "benchmark"))] pub fn global_service() -> Router { - Router::new() + use axum::routing::post; + + let router = Router::new() .route("/list", get(list_worker_pings)) .route("/custom_tags", get(get_custom_tags)) -} + .route("/list_worker_groups", get(get_worker_groups)) + .route( + "/worker_group/:name", + post(update_worker_group).delete(delete_worker_group), + ); + #[cfg(feature = "benchmark")] + return router.route("/toggle", get(toggle)); -#[cfg(feature = "benchmark")] -pub fn global_service() -> Router { - Router::new() - .route("/toggle", get(toggle)) - .route("/list", get(list_worker_pings)) - .route("/custom_tags", get(get_custom_tags)) -} - -lazy_static::lazy_static! { - pub static ref CUSTOM_TAGS: Vec = std::env::var("CUSTOM_TAGS") - .ok() - .map(|x| x.split(',').map(|x| x.to_string()).collect::>()).unwrap_or_default(); - - pub static ref CUSTOM_TAGS_PER_WORKSPACE: (Vec, HashMap>) = process_custom_tags(std::env::var("CUSTOM_TAGS") - .ok()); - - pub static ref ALL_TAGS: Vec = [CUSTOM_TAGS_PER_WORKSPACE.0.clone(), CUSTOM_TAGS_PER_WORKSPACE.1.keys().map(|x| x.to_string()).collect_vec()].concat(); - -} - -fn process_custom_tags(o: Option) -> (Vec, HashMap>) { - let regex = Regex::new(r"^(\w+)\(((?:\w+)\+?)+\)$").unwrap(); - if let Some(s) = o { - let mut global = vec![]; - let mut specific: HashMap> = HashMap::new(); - for e in s.split(",") { - if let Some(cap) = regex.captures(e) { - let tag = cap.get(1).unwrap().as_str().to_string(); - let workspaces = cap.get(2).unwrap().as_str().split("+"); - specific.insert(tag, workspaces.map(|x| x.to_string()).collect_vec()); - } else { - global.push(e.to_string()); - } - } - (global, specific) - } else { - (vec![], HashMap::new()) - } + #[cfg(not(feature = "benchmark"))] + return router; } #[derive(FromRow, Serialize, Deserialize)] @@ -86,6 +56,7 @@ struct WorkerPing { ip: String, jobs_executed: i32, custom_tags: Option>, + worker_group: String, } #[derive(Serialize, Deserialize)] @@ -104,7 +75,7 @@ async fn list_worker_pings( let rows = sqlx::query_as!( WorkerPing, - "SELECT worker, worker_instance, EXTRACT(EPOCH FROM (now() - ping_at))::integer as last_ping, started_at, ip, jobs_executed, custom_tags FROM worker_ping ORDER BY ping_at desc LIMIT $1 OFFSET $2", + "SELECT worker, worker_instance, EXTRACT(EPOCH FROM (now() - ping_at))::integer as last_ping, started_at, ip, jobs_executed, custom_tags, worker_group FROM worker_ping ORDER BY ping_at desc LIMIT $1 OFFSET $2", per_page as i64, offset as i64 ) @@ -121,5 +92,77 @@ async fn toggle(Query(query): Query) -> JsonResult { } async fn get_custom_tags() -> Json> { - Json(ALL_TAGS.clone()) + Json(ALL_TAGS.read().await.clone().into()) +} + +#[derive(Serialize, Deserialize, FromRow)] +struct WorkerGroup { + name: String, + config: serde_json::Value, +} + +async fn get_worker_groups( + authed: ApiAuthed, + Extension(db): Extension, + Extension(user_db): Extension, +) -> error::JsonResult> { + let mut tx = user_db.begin(&authed).await?; + + require_super_admin(&db, &authed.email).await?; + + let rows = sqlx::query_as!(WorkerGroup, "SELECT * FROM worker_group_config") + .fetch_all(&mut *tx) + .await?; + tx.commit().await?; + Ok(Json(rows)) +} + +async fn update_worker_group( + Path(name): Path, + Extension(db): Extension, + Extension(user_db): Extension, + authed: ApiAuthed, + Json(config): Json, +) -> error::Result { + let tx = user_db.begin(&authed).await?; + + require_super_admin(&db, &authed.email).await?; + tx.commit().await?; + + sqlx::query!( + "INSERT INTO worker_group_config (name, config) VALUES ($1, $2) ON CONFLICT (name) DO UPDATE SET config = $2", + &name, + config + ) + .execute(&db) + .await?; + + Ok(format!("Updated worker group {name}")) +} + +async fn delete_worker_group( + Path(name): Path, + Extension(db): Extension, + Extension(user_db): Extension, + authed: ApiAuthed, +) -> error::Result { + let tx = user_db.begin(&authed).await?; + + require_super_admin(&db, &authed.email).await?; + tx.commit().await?; + + let deleted = sqlx::query!( + "DELETE FROM worker_group_config WHERE name = $1 RETURNING name", + name, + ) + .fetch_all(&db) + .await?; + + if deleted.len() == 0 { + return Err(error::Error::NotFound(format!( + "Worker group {name} not found", + name = name + ))); + } + Ok(format!("Deleted worker group {name}")) } diff --git a/backend/windmill-api/src/workspaces.rs b/backend/windmill-api/src/workspaces.rs index bbae417424..bc0a830a55 100644 --- a/backend/windmill-api/src/workspaces.rs +++ b/backend/windmill-api/src/workspaces.rs @@ -74,13 +74,19 @@ pub fn workspaced_service() -> Router { .route("/edit_error_handler", post(edit_error_handler)); #[cfg(feature = "enterprise")] - tracing::info!("stripe enabled"); - - #[cfg(feature = "enterprise")] - let router = router - .route("/checkout", get(stripe_checkout)) - .route("/billing_portal", get(stripe_portal)); + { + if std::env::var("STRIPE_KEY").is_err() { + return router; + } else { + tracing::info!("stripe enabled"); + + return router + .route("/checkout", get(stripe_checkout)) + .route("/billing_portal", get(stripe_portal)); + } + } + #[cfg(not(feature = "enterprise"))] router } pub fn global_service() -> Router { diff --git a/backend/windmill-common/Cargo.toml b/backend/windmill-common/Cargo.toml index c0e37e8fa5..caccedb1d5 100644 --- a/backend/windmill-common/Cargo.toml +++ b/backend/windmill-common/Cargo.toml @@ -44,3 +44,5 @@ reqwest = { workspace = true, optional = true } tracing-subscriber = { workspace = true, optional = true } lazy_static.workspace = true tracing-flame = { version = "^0", optional = true } +itertools.workspace = true +regex.workspace = true \ No newline at end of file diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index e01a8fd152..490bd47338 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -1,6 +1,7 @@ pub const WORKER_S3_BUCKET_SYNC: &str = "worker_s3_bucket_sync"; +pub const CUSTOM_TAGS_SETTING: &str = "custom_tags"; -pub const ENV_SETTINGS: [&str; 55] = [ +pub const ENV_SETTINGS: [&str; 54] = [ "DISABLE_NSJAIL", "DISABLE_SERVER", "NUM_WORKERS", @@ -41,8 +42,6 @@ pub const ENV_SETTINGS: [&str; 55] = [ "INSTANCE_EVENTS_WEBHOOK", "CLOUD_HOSTED", "GLOBAL_CACHE_INTERVAL", - "WORKER_TAGS", - "CUSTOM_TAGS", "JOB_RETENTION_SECS", "WAIT_RESULT_FAST_POLL_DURATION_SECS", "WAIT_RESULT_SLOW_POLL_INTERVAL_MS", @@ -56,4 +55,5 @@ pub const ENV_SETTINGS: [&str; 55] = [ "CREATE_WORKSPACE_REQUIRE_SUPERADMIN", "GLOBAL_ERROR_HANDLER_PATH_IN_ADMINS_WORKSPACE", "MAX_WAIT_FOR_SIGTERM", + "WORKER_GROUP", ]; diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index 8058c3a75e..2451930ab1 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -27,6 +27,7 @@ pub mod scripts; pub mod users; pub mod utils; pub mod variables; +pub mod worker; #[cfg(feature = "tracing_init")] pub mod tracing_init; @@ -77,12 +78,15 @@ pub async fn shutdown_signal( Ok(()) } +#[cfg(feature = "prometheus")] +use tokio::task::JoinHandle; + #[cfg(feature = "prometheus")] pub async fn serve_metrics( addr: SocketAddr, mut rx: tokio::sync::broadcast::Receiver<()>, ready_worker_endpoint: bool, -) -> Result<(), hyper::Error> { +) -> JoinHandle<()> { use std::sync::atomic::Ordering; use axum::{routing::get, Router}; @@ -104,13 +108,18 @@ pub async fn serve_metrics( router }; - axum::Server::bind(&addr) - .serve(router.into_make_service()) - .with_graceful_shutdown(async { - rx.recv().await.ok(); - println!("Graceful shutdown of metrics"); - }) - .await + tokio::spawn(async move { + if let Err(e) = axum::Server::bind(&addr) + .serve(router.into_make_service()) + .with_graceful_shutdown(async { + rx.recv().await.ok(); + println!("Graceful shutdown of metrics"); + }) + .await + { + tracing::error!("Error serving metrics: {}", e); + } + }) } async fn metrics() -> Result { diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs new file mode 100644 index 0000000000..2dc7fa9de3 --- /dev/null +++ b/backend/windmill-common/src/worker.rs @@ -0,0 +1,200 @@ +use std::{collections::HashMap, sync::Arc}; + +use itertools::Itertools; +use regex::Regex; +use serde::{Deserialize, Serialize}; +use tokio::sync::RwLock; + +use crate::{error, global_settings::CUSTOM_TAGS_SETTING, DB}; + +lazy_static::lazy_static! { + pub static ref WORKER_GROUP: String = std::env::var("WORKER_GROUP").unwrap_or_else(|_| "default".to_string()); + + pub static ref DEFAULT_TAGS : Vec = vec![ + "deno".to_string(), + "python3".to_string(), + "go".to_string(), + "bash".to_string(), + "powershell".to_string(), + "nativets".to_string(), + "mysql".to_string(), + "graphql".to_string(), + "bun".to_string(), + "postgresql".to_string(), + "bigquery".to_string(), + "snowflake".to_string(), + "graphql".to_string(), + "dependency".to_string(), + "flow".to_string(), + "hub".to_string(), + "other".to_string()]; + + + pub static ref WORKER_CONFIG: Arc> = Arc::new(RwLock::new(WorkerConfig { + worker_tags: Default::default(), + dedicated_worker: Default::default(), + })); + + + pub static ref CLOUD_HOSTED: bool = std::env::var("CLOUD_HOSTED").is_ok(); + + pub static ref CUSTOM_TAGS: Vec = std::env::var("CUSTOM_TAGS") + .ok() + .map(|x| x.split(',').map(|x| x.to_string()).collect::>()).unwrap_or_default(); + + + pub static ref CUSTOM_TAGS_PER_WORKSPACE: Arc, HashMap>)>> = Arc::new(RwLock::new((vec![], HashMap::new()))); + + pub static ref ALL_TAGS: Arc>> = Arc::new(RwLock::new(vec![])); + + static ref CUSTOM_TAG_REGEX: Regex = Regex::new(r"^(\w+)\(((?:\w+)\+?)+\)$").unwrap(); + +} + +pub async fn reload_custom_tags_setting(db: &DB) -> error::Result<()> { + let q = sqlx::query!( + "SELECT value FROM global_settings WHERE name = $1", + CUSTOM_TAGS_SETTING + ) + .fetch_optional(db) + .await?; + + let tags = if let Some(q) = q { + if let Ok(v) = serde_json::from_value::>(q.value.clone()) { + v + } else { + tracing::error!( + "Could not parse custom tags setting as vec of strings, found: {:#?}", + &q.value + ); + vec![] + } + } else { + CUSTOM_TAGS.clone() + }; + + let custom_tags = process_custom_tags(tags); + + { + let l = CUSTOM_TAGS_PER_WORKSPACE.read().await; + if l.clone() == custom_tags { + tracing::info!("Custom tags setting unchanged, skipping update"); + return Ok(()); + } else { + tracing::info!("Custom tags setting changed, updating"); + } + } + { + let mut l = CUSTOM_TAGS_PER_WORKSPACE.write().await; + *l = custom_tags.clone() + } + { + let mut l = ALL_TAGS.write().await; + *l = [ + custom_tags.0.clone(), + custom_tags.1.keys().map(|x| x.to_string()).collect_vec(), + ] + .concat(); + } + Ok(()) + // pub static ref CUSTOM_TAGS_PER_WORKSPACE: (Vec, HashMap>) = process_custom_tags(std::env::var("CUSTOM_TAGS") + // .ok()); + + // pub static ref ALL_TAGS: Vec = [CUSTOM_TAGS_PER_WORKSPACE.0.clone(), CUSTOM_TAGS_PER_WORKSPACE.1.keys().map(|x| x.to_string()).collect_vec()].concat(); +} + +fn process_custom_tags(tags: Vec) -> (Vec, HashMap>) { + let mut global = vec![]; + let mut specific: HashMap> = HashMap::new(); + for e in tags { + if let Some(cap) = CUSTOM_TAG_REGEX.captures(&e) { + let tag = cap.get(1).unwrap().as_str().to_string(); + let workspaces = cap.get(2).unwrap().as_str().split("+"); + specific.insert(tag, workspaces.map(|x| x.to_string()).collect_vec()); + } else { + global.push(e.to_string()); + } + } + (global, specific) +} + +pub async fn update_ping(worker_instance: &str, worker_name: &str, ip: &str, db: &DB) { + let wc = WORKER_CONFIG.read().await; + let tags = wc.worker_tags.as_slice(); + sqlx::query!( + "INSERT INTO worker_ping (worker_instance, worker, ip, custom_tags, worker_group, dedicated_worker) VALUES ($1, $2, $3, $4, $5, $6) ON CONFLICT (worker) DO UPDATE set ip = $3, custom_tags = $4, worker_group = $5", + worker_instance, + worker_name, + ip, + tags, + *WORKER_GROUP, + wc.dedicated_worker.as_ref().map(|x| format!("{}:{}", x.workspace_id, x.path)) + ) + .execute(db) + .await + .expect("insert worker_ping initial value"); +} + +pub async fn load_worker_config(db: &DB) -> error::Result { + let config: WorkerConfigOpt = sqlx::query_scalar!( + "SELECT config FROM worker_group_config WHERE name = $1", + *WORKER_GROUP + ) + .fetch_optional(db) + .await? + .flatten() + .map(|x| serde_json::from_value(x).ok()) + .flatten() + .unwrap_or_default(); + let dedicated_worker = config.dedicated_worker.map(|x| { + let splitted = x.split(':').to_owned().collect_vec(); + if splitted.len() != 2 { + panic!("DEDICATED_WORKER setting should be in the form of :") + } else { + let workspace = splitted[0]; + let script_path = splitted[1]; + WorkspacedPath { workspace_id: workspace.to_string(), path: script_path.to_string() } + } + }); + Ok(WorkerConfig { + worker_tags: config + .worker_tags + .or_else(|| { + if let Some(ref dedicated_worker) = dedicated_worker.as_ref() { + Some(vec![format!( + "{}:{}", + dedicated_worker.workspace_id, dedicated_worker.path + )]) + } else { + std::env::var("WORKER_TAGS") + .ok() + .map(|x| x.split(',').map(|x| x.to_string()).collect()) + } + }) + .unwrap_or_else(|| DEFAULT_TAGS.clone()), + dedicated_worker, + }) +} + +#[derive(Clone, PartialEq)] +pub struct WorkspacedPath { + pub workspace_id: String, + pub path: String, +} +#[derive(Serialize, Deserialize)] +pub struct WorkerConfigOpt { + pub worker_tags: Option>, + pub dedicated_worker: Option, +} + +impl Default for WorkerConfigOpt { + fn default() -> Self { + Self { worker_tags: Default::default(), dedicated_worker: Default::default() } + } +} + +#[derive(PartialEq)] +pub struct WorkerConfig { + pub worker_tags: Vec, + pub dedicated_worker: Option, +} diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 9db9d129f6..8886391a4a 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -14,7 +14,6 @@ use std::time::Instant; use anyhow::Context; use async_recursion::async_recursion; use chrono::{DateTime, Duration, Utc}; -use itertools::Itertools; use reqwest::Client; use rsmq_async::RsmqConnection; use serde_json::json; @@ -37,9 +36,13 @@ use windmill_common::{ schedule::{schedule_to_user, Schedule}, scripts::{ScriptHash, ScriptLang}, users::{username_to_permissioned_as, SUPERADMIN_SECRET_EMAIL}, + worker::WORKER_CONFIG, DB, METRICS_ENABLED, }; +#[cfg(feature = "enterprise")] +use windmill_common::worker::CLOUD_HOSTED; + use crate::{ schedule::{get_schedule_opt, push_scheduled_job}, QueueTransaction, @@ -66,57 +69,6 @@ lazy_static::lazy_static! { "Total number of jobs pulled from the queue." ) .unwrap(); - pub static ref CLOUD_HOSTED: bool = std::env::var("CLOUD_HOSTED").is_ok(); - - pub static ref DEFAULT_TAGS : Vec = vec![ - "deno".to_string(), - "python3".to_string(), - "go".to_string(), - "bash".to_string(), - "powershell".to_string(), - "nativets".to_string(), - "mysql".to_string(), - "graphql".to_string(), - "bun".to_string(), - "postgresql".to_string(), - "bigquery".to_string(), - "snowflake".to_string(), - "graphql".to_string(), - "dependency".to_string(), - "flow".to_string(), - "hub".to_string(), - "other".to_string()]; - - pub static ref DEDICATED_WORKER: Option<(String, String)> = std::env::var("DEDICATED_WORKER") - .ok() - .map(|x| { - let splitted = x.split(':').to_owned().collect_vec(); - if splitted.len() != 2 { - panic!("DEDICATED_WORKER should be in the form of :") - } else { - let workspace = splitted[0]; - let script_path = splitted[1]; - (workspace.to_string(), script_path.to_string()) - } - }); - - pub static ref ACCEPTED_TAGS: Vec = { - let worker_tags = std::env::var("WORKER_TAGS") - .ok() - .map(|x| x.split(',').map(|x| x.to_string()).collect()) - .unwrap_or_else(|| DEFAULT_TAGS.clone()); - if let Some(ref dedicated_worker) = DEDICATED_WORKER.as_ref() { - vec![format!("{}:{}", dedicated_worker.0, dedicated_worker.1)] - } else { - worker_tags - } - }; - - pub static ref IS_WORKER_TAGS_DEFINED: bool = std::env::var("WORKER_TAGS").ok().is_some(); - - - - // When compiled in 'benchmark' mode, this flags is exposed via the /workers/toggle endpoint // and make it possible to disable to current active workers (such that they don't pull any) @@ -425,7 +377,7 @@ pub async fn add_completed_job( if !is_flow && _duration > 1000 { let additional_usage = _duration / 1000; let w_id = &queued_job.workspace_id; - let premium_workspace = *CLOUD_HOSTED + let premium_workspace = *windmill_common::worker::CLOUD_HOSTED && sqlx::query_scalar!("SELECT premium FROM workspace WHERE id = $1", w_id) .fetch_one(db) .await @@ -1153,7 +1105,7 @@ async fn pull_single_job_and_mark_as_running_no_concurrency_limit< // TODO: REDIS: Race conditions / replace last_ping // TODO: shuffle this list to have fairness - let mut all_tags = ACCEPTED_TAGS.clone(); + let mut all_tags = WORKER_CONFIG.read().await.worker_tags.clone(); let mut msg: Option<_> = None; let mut tag = None; @@ -1210,6 +1162,8 @@ async fn pull_single_job_and_mark_as_running_no_concurrency_limit< * suspend_until is non-null * and suspend = 0 when the resume messages are received * or suspend_until <= now() if it has timed out */ + let config = WORKER_CONFIG.read().await; + let tags = config.worker_tags.as_slice(); let r = if suspend_first { sqlx::query_as::<_, QueuedJob>("UPDATE queue @@ -1226,17 +1180,21 @@ async fn pull_single_job_and_mark_as_running_no_concurrency_limit< LIMIT 1 ) RETURNING *") - .bind(ACCEPTED_TAGS.as_slice()) + .bind(tags) .fetch_optional(db) .await? } else { None }; + drop(config); if r.is_none() { // #[cfg(feature = "benchmark")] // let instant = Instant::now(); + let config = WORKER_CONFIG.read().await; + let tags = config.worker_tags.as_slice(); + let r = sqlx::query_as::<_, QueuedJob>( "UPDATE queue SET running = true @@ -1253,10 +1211,9 @@ async fn pull_single_job_and_mark_as_running_no_concurrency_limit< ) RETURNING *", ) - .bind(ACCEPTED_TAGS.as_slice()) + .bind(tags) .fetch_optional(db) .await?; - // #[cfg(feature = "benchmark")] // println!("pull query: {:?}", instant.elapsed()); diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index e25ba71a4c..fd9ea8db24 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -5,12 +5,12 @@ use serde_json::{json, Value}; use sqlx::{Pool, Postgres}; use tokio::{fs::File, io::AsyncReadExt}; use windmill_api_client::{types::CreateResource, Client}; +use windmill_common::worker::CLOUD_HOSTED; use windmill_common::{ error::{self, Error}, jobs::QueuedJob, variables::ContextualVariable, }; -use windmill_queue::CLOUD_HOSTED; use anyhow::Result; use std::{ diff --git a/backend/windmill-worker/src/config.rs b/backend/windmill-worker/src/config.rs new file mode 100644 index 0000000000..e69de29bb2 diff --git a/backend/windmill-worker/src/lib.rs b/backend/windmill-worker/src/lib.rs index 94c582eeaa..5325201ee4 100644 --- a/backend/windmill-worker/src/lib.rs +++ b/backend/windmill-worker/src/lib.rs @@ -6,6 +6,7 @@ mod snowflake_executor; mod bash_executor; mod bun_executor; mod common; +mod config; mod dedicated_worker; mod deno_executor; mod global_cache; @@ -17,5 +18,4 @@ mod pg_executor; mod python_executor; mod worker; mod worker_flow; - pub use worker::*; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 8b30f749ba..1832a31b09 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -28,12 +28,10 @@ use windmill_common::{ scripts::{get_full_hub_script_by_path, ScriptHash, ScriptLang}, users::SUPERADMIN_SECRET_EMAIL, utils::{rd_string, StripPath}, + worker::{update_ping, CLOUD_HOSTED, WORKER_CONFIG}, DB, IS_READY, METRICS_ENABLED, }; -use windmill_queue::{ - canceled_job_to_result, get_queued_job, pull, ACCEPTED_TAGS, CLOUD_HOSTED, DEDICATED_WORKER, - HTTP_CLIENT, IS_WORKER_TAGS_DEFINED, -}; +use windmill_queue::{canceled_job_to_result, get_queued_job, pull, HTTP_CLIENT}; use serde_json::{json, Value}; @@ -350,6 +348,7 @@ pub async fn run_worker, + killpill_tx: tokio::sync::broadcast::Sender<()>, base_internal_url: &str, rsmq: Option, _sync_barrier: Arc>>, @@ -388,7 +387,7 @@ pub async fn run_worker(MAX_BUFFERED_DEDICATED_JOBS); - let killpill_rx = killpill_rx.resubscribe(); + let mut killpill_rx = killpill_rx.resubscribe(); let db = db.clone(); let worker_dir = worker_dir.clone(); let base_internal_url = base_internal_url.to_string(); @@ -645,19 +648,52 @@ pub async fn run_worker, Option, Option>)>( + let (content, lock, _language, envs) = { + let r; + loop { + let q = sqlx::query_as::<_, (String, Option, Option, Option>)>( "SELECT content, lock, language, envs FROM script WHERE path = $1 AND workspace_id = $2 AND created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND deleted = false AND lock IS not NULL AND lock_error_logs IS NULL)", ) - .bind(&_script_path) - .bind(&_workspace) + .bind(&_wp.path) + .bind(&_wp.workspace_id) .fetch_optional(&db) - .await.expect("Failed to fetch script for dedicated worker") - .expect(&format!("Failed to fetch script `{_script_path}` in workspace {_workspace} for dedicated worker")); + .await; + if let Ok(q) = q { + if let Some(wp) = q { + r = wp; + break; + } else { + tracing::error!( + "Failed to fetch script `{}` in workspace {} for dedicated worker. Retrying in 10s.", + _wp.path, + _wp.workspace_id + ); + tokio::select! { + biased; + _ = killpill_rx.recv() => { + tracing::info!("Killing dedicated worker while it was attempting to fetch script"); + return; + } + _ = tokio::time::sleep(Duration::from_secs(10)) => { + continue; + } + + } + } + } else { + tracing::error!("Failed to fetch script for dedicated worker"); + killpill_tx.send(()).expect("send"); + return; + } + } + r + }; let worker_envs = build_envs(envs).expect("failed to build envs"); if let Err(e) = start_worker( @@ -668,8 +704,8 @@ pub async fn run_worker NUM_SECS_PING { + let wc = WORKER_CONFIG.read().await; + let tags = wc.worker_tags.as_slice(); + sqlx::query!( - "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1 WHERE worker = $2", + "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2 WHERE worker = $3", jobs_executed, + tags, &worker_name ) .execute(db) @@ -1269,25 +1309,6 @@ pub async fn handle_job_error, -) { - let tags = ACCEPTED_TAGS.clone(); - sqlx::query!( - "INSERT INTO worker_ping (worker_instance, worker, ip, custom_tags) VALUES ($1, $2, $3, $4) ON CONFLICT (worker) DO NOTHING", - worker_instance, - worker_name, - ip, - if *IS_WORKER_TAGS_DEFINED { Some(tags.as_slice()) } else { None } - ) - .execute(db) - .await - .expect("insert worker_ping initial value"); -} - fn extract_error_value(log_lines: &str, i: i32) -> serde_json::Value { return json!({"message": format!("ExitCode: {i}, last log lines:\n{}", ANSI_ESCAPE_RE.replace_all(log_lines.trim(), "").to_string()), "name": "ExecutionErr"}); } diff --git a/frontend/src/lib/components/PageHeader.svelte b/frontend/src/lib/components/PageHeader.svelte index 00996dedfe..384bcbcd51 100644 --- a/frontend/src/lib/components/PageHeader.svelte +++ b/frontend/src/lib/components/PageHeader.svelte @@ -12,7 +12,7 @@

{title}

{#if tooltip != '' || documentationLink} - + {tooltip} {/if} @@ -21,7 +21,7 @@

{title}

{#if tooltip != '' || documentationLink} - + {tooltip} {/if} diff --git a/frontend/src/lib/components/ScriptBuilder.svelte b/frontend/src/lib/components/ScriptBuilder.svelte index b25ff85f10..3aaaebf276 100644 --- a/frontend/src/lib/components/ScriptBuilder.svelte +++ b/frontend/src/lib/components/ScriptBuilder.svelte @@ -483,7 +483,7 @@ > {title} {#if desc} - + {desc} {/if} diff --git a/frontend/src/lib/components/Tooltip.svelte b/frontend/src/lib/components/Tooltip.svelte index f62b5eff2b..a97fb3b67e 100644 --- a/frontend/src/lib/components/Tooltip.svelte +++ b/frontend/src/lib/components/Tooltip.svelte @@ -1,25 +1,22 @@ - +
+ +
{#if documentationLink} diff --git a/frontend/src/lib/components/WorkspaceGroup.svelte b/frontend/src/lib/components/WorkspaceGroup.svelte new file mode 100644 index 0000000000..f18a03cbb3 --- /dev/null +++ b/frontend/src/lib/components/WorkspaceGroup.svelte @@ -0,0 +1,246 @@ + + + { + open = false + }} + on:confirmed={async () => { + deleteWorkerGroup() + open = false + }} +> +
+ Are you sure you want to remove this worker nconfig? +
+
+ +

{name}

+ {#if $superadmin} + + + + + { + dirty = true + if (nconfig == undefined) { + nconfig = {} + } + console.log(e.detail) + if (e.detail == 'dedicated') { + nconfig.dedicated_worker = '' + nconfig.worker_tags = undefined + } else { + nconfig.dedicated_worker = undefined + nconfig.worker_tags = [] + } + }} + class="mb-4" + > + + + + {#if selected == 'normal'} + {#if nconfig?.worker_tags != undefined} +
+ {#each nconfig.worker_tags as tag} +
- {tag}
+
+ {/each} +
+ +
+ +
+ + + +
+ {/if} + {:else if selected == 'dedicated'} + {#if nconfig?.dedicated_worker != undefined} + { + dirty = true + }} + bind:value={nconfig.dedicated_worker} + /> +

Workers will get killed upon detecting this setting change. It is assumed they are in + an environment where the supervisor will restart them. Upon restart, they will pick the + new dedicated worker config.

+ {/if} + {/if} +
+
+ + {#if !$enterpriseLicense}{selected == 'dedicated' + ? 'Dedicated workers are an enterprise only feature' + : 'The Worker Group Manager UI is an enterprise only feature. However, workers can still have their WORKER_TAGS passed as env'}{/if} +
+ + {#if config} + + {/if} + {/if} +
diff --git a/frontend/src/lib/components/common/button/Button.svelte b/frontend/src/lib/components/common/button/Button.svelte index be5cd0d316..0030d55588 100644 --- a/frontend/src/lib/components/common/button/Button.svelte +++ b/frontend/src/lib/components/common/button/Button.svelte @@ -89,9 +89,9 @@ }, light: { border: - 'border bg-surface hover:bg-surface-hover focus:bg-surface-hover text-primary hover:text-secondary focus:text-secondary focus:ring-surface-selected', + 'border bg-surface hover:bg-surface-hover focus:bg-surface-hover text-primary hover:text-secondary focus:text-secondary focus:ring-surface-selected', contained: - 'bg-surface hover:bg-surface-hover focus:bg-surface-hover text-primary focus:ring-surface-selected', + 'bg-surface border-transparent hover:bg-surface-hover focus:bg-surface-hover text-primary focus:ring-surface-selected', divider: 'divide-x divide-gray-200 dark:divide-gray-700' } } @@ -174,9 +174,7 @@ on:blur class={twMerge( buttonClass, - disabled - ? '!bg-surface-disabled !text-tertiary border border-disabled !cursor-not-allowed' - : '' + disabled ? '!bg-surface-disabled !text-tertiary border !cursor-not-allowed' : 'border' )} {id} tabindex={disabled ? -1 : 0} diff --git a/frontend/src/lib/components/flows/header/FlowPreviewButtons.svelte b/frontend/src/lib/components/flows/header/FlowPreviewButtons.svelte index 1350309cc0..e2a8dae238 100644 --- a/frontend/src/lib/components/flows/header/FlowPreviewButtons.svelte +++ b/frontend/src/lib/components/flows/header/FlowPreviewButtons.svelte @@ -26,6 +26,8 @@ 'settings-same-worker', 'settings-graph', 'settings-worker-group', + 'settings-cache', + 'settings-concurrency', 'inputs', 'schedules', 'failure', diff --git a/frontend/src/lib/utils.ts b/frontend/src/lib/utils.ts index 756bc0e772..048cd4e022 100644 --- a/frontend/src/lib/utils.ts +++ b/frontend/src/lib/utils.ts @@ -228,28 +228,29 @@ export function setQueryWithoutLoad( }, bounceTime ?? 200) } -export function groupBy( - items: T[], - toGroup: (t: T) => string, - toSort: (t: T) => string, - dflts: string[] = [] -): [string, T[]][] { - let r: Record = {} +export function groupBy( + items: V[], + toGroup: (t: V) => K, + toSort: (t: V) => string, + dflts: K[] = [] +): [K, V[]][] { + let r: Map = new Map() for (const dflt of dflts) { - r[dflt] = [] + r.set(dflt as K, []) } items.forEach((sc) => { let section = toGroup(sc) - if (section in r) { - r[section].push(sc) - r[section].sort((a, b) => toSort(a).localeCompare(toSort(b))) + if (r.has(section)) { + let arr = r.get(section)! + arr.push(sc) + arr.sort((a, b) => toSort(a).localeCompare(toSort(b))) } else { - r[section] = [sc] + r.set(section, [sc]) } }) - return Object.entries(r).sort((s1, s2) => { + return [...r.entries()].sort((s1, s2) => { let n1 = s1[0] let n2 = s2[0] diff --git a/frontend/src/routes/(root)/(logged)/workers/+page.svelte b/frontend/src/routes/(root)/(logged)/workers/+page.svelte index 8b179b40d5..7493006f34 100644 --- a/frontend/src/routes/(root)/(logged)/workers/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/workers/+page.svelte @@ -1,6 +1,6 @@ + > + {#if $superadmin} +
+
+ { + try { + console.log('Setting global cache to', e.detail) + await SettingService.setGlobal({ + key: WORKER_S3_BUCKET_SYNC_SETTING, + requestBody: { value: e.detail } + }) + globalCache = e.detail + } catch (err) { + sendUserToast(`Could not set global cache: ${err}`, true) + } + }} + options={{ right: 'global cache to s3' }} + disabled={!$enterpriseLicense} + /> +

global cache to s3 is an enterprise feature that enable workers to do fast cold start + and share a single cache backed by s3 to ensure that even with a high number of + workers, dependencies for python/deno/bun/go are only downloaded for the first time + only once by the whole fleet. +

require S3_CACHE_BUCKET to be set and has NO effect otherwise (even if this setting + is on)
+
+
+ + + +
+ {#if customTags == undefined} + + {:else} +
+ {#each customTags as customTag} +
+
- {customTag}
+ +
+ {/each} +
+ + + + For tags specific to some workspaces, use
tag(workspace1+workspace2)
+ For dynamic tags based on the workspace, use
$workspace
, e.g: +
tag-$workspace
+ {/if} +
+
+
+
+ {/if} +
- {#if $superadmin} -
-

global cache to s3 is an enterprise feature that enable workers to do fast cold start and - share a single cache backed by s3 to ensure that even with a high number of workers, - dependencies for python/deno/bun/go are only downloaded for the first time only once by - the whole fleet. -

require S3_CACHE_BUCKET to be set and has NO effect otherwise (even if this setting is - on)
- { - try { - console.log('Setting global cache to', e.detail) - await SettingService.setGlobal({ - key: worker_s3_bucket_sync, - requestBody: { value: e.detail } - }) - globalCache = e.detail - } catch (err) { - sendUserToast(`Could not set global cache: ${err}`, true) - } - }} - options={{ right: 'global cache to s3' }} - disabled={!$enterpriseLicense} - /> -
- {/if} {#if workers != undefined} {#if groupedWorkers.length == 0}

No workers seems to be available

{/if} - - - - Worker - -
- Custom Tags - - If defined, the workers only pull jobs with the same corresponding tag - -
-
- Last ping - Worker start - Nb of jobs executed - Liveness - - - - {#each groupedWorkers as [section, workers]} - - - Instance: {section} - IP: {workers[0].ip} - - +

Worker Groups Worker groups are groups of workers that share a config and are meant to be identical. + Worker groups are meant to be used with tags. Tags can be assigned to scripts and flows + and can be seen as dedicated queues. Only the corresponding +

+
+ {#if $superadmin} +
+ + + Worker Group configs are propagated to every workers in the worker group +
+ {/if}
+ {#each groupedWorkers as worker_group} + { + loadWorkerGroups() + }} + /> - {#if workers} - {#each workers as { worker, custom_tags, last_ping, started_at, jobs_executed }} - - {worker} - {custom_tags?.join(', ') ?? ''} - {last_ping != undefined ? last_ping + timeSinceLastPing : -1}s ago - {displayDate(started_at)} - {jobs_executed} - - + + + Worker + +
+ Worker Tags + + If defined, the workers only pull jobs with the same corresponding tag + +
+
+ Last ping + Worker start + Nb of jobs executed + Liveness + + + + {#each worker_group[1] as [section, workers]} + + + Instance: {section[0]} + IP: {workers[0].ip} + + + + {#if workers} + {#each workers as { worker, custom_tags, last_ping, started_at, jobs_executed }} + + {worker} + {#if custom_tags && custom_tags?.length > 2}{truncate( + custom_tags?.join(', ') ?? '', + 10 + )} + {custom_tags?.join(', ')}{:else}{custom_tags?.join(', ') ?? + ''}{/if} - {last_ping != undefined ? (last_ping < 60 ? 'Alive' : 'Dead') : 'Unknown'} -
-
- - {/each} - {/if} - {/each} - - + {last_ping != undefined ? last_ping + timeSinceLastPing : -1}s ago + {displayDate(started_at)} + {jobs_executed} + + + {last_ping != undefined ? (last_ping < 60 ? 'Alive' : 'Dead') : 'Unknown'} + + + + {/each} + {/if} + {/each} + + +
+ {/each} + + {#each Object.entries(workerGroups ?? {}).filter((x) => !groupedWorkers.some((y) => y[0] == x[0])) as worker_group} + { + loadWorkerGroups() + }} + name={worker_group[0]} + config={worker_group[1]} + /> +
No workers currently in this worker group
+ {/each} {:else}
{#each new Array(4) as _}