From aeacaa2b01b6ca17ca00c9cc5f8b2ebc996b851f Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Mon, 6 Nov 2023 11:21:36 +0100 Subject: [PATCH] fix: improve dedicated workers --- backend/tests/worker.rs | 16 +++-- backend/windmill-api/openapi.yaml | 3 +- backend/windmill-api/src/jobs.rs | 2 + backend/windmill-common/src/jobs.rs | 1 + backend/windmill-common/src/worker.rs | 9 ++- backend/windmill-queue/src/jobs.rs | 5 +- backend/windmill-worker/src/bun_executor.rs | 45 ++++++++------ .../windmill-worker/src/dedicated_worker.rs | 28 +++++++-- backend/windmill-worker/src/worker.rs | 62 ++++++++++++------- backend/windmill-worker/src/worker_flow.rs | 5 +- .../lib/components/sidebar/UserMenu.svelte | 2 +- 11 files changed, 121 insertions(+), 57 deletions(-) diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 8637a5f04a..c18ca800b0 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1658,7 +1658,8 @@ func main(derp string) (string, error) { language: ScriptLang::Go, concurrent_limit: None, concurrency_time_window_s: None, - cache_ttl: None + cache_ttl: None, + dedicated_worker: None })) .arg("derp", json!("world")) .run_until_complete(&db, port) @@ -1688,7 +1689,8 @@ echo "hello $msg" language: ScriptLang::Bash, concurrent_limit: None, concurrency_time_window_s: None, - cache_ttl: None + cache_ttl: None, + dedicated_worker: None })) .arg("msg", json!("world")) .run_until_complete(&db, port) @@ -1715,7 +1717,8 @@ def main(): lock: None, concurrent_limit: None, concurrency_time_window_s: None, - cache_ttl: None + cache_ttl: None, + dedicated_worker: None }); let result = run_job_in_new_worker_until_complete(&db, job, port) @@ -1748,7 +1751,8 @@ def main(): lock: None, concurrent_limit: None, concurrency_time_window_s: None, - cache_ttl: None + cache_ttl: None, + dedicated_worker: None }); let result = run_job_in_new_worker_until_complete(&db, job, port) @@ -1780,7 +1784,8 @@ def main(): lock: None, concurrent_limit: None, concurrency_time_window_s: None, - cache_ttl: None + cache_ttl: None, + dedicated_worker: None }); let result = run_job_in_new_worker_until_complete(&db, job, port) @@ -3177,6 +3182,7 @@ async fn run_preview_relative_imports(db: &Pool, script_content: Strin concurrent_limit: None, concurrency_time_window_s: None, cache_ttl: None, + dedicated_worker: None })).push(&db2).await; diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 49b29d1cc7..08239168a3 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -7452,7 +7452,8 @@ components: kind: type: string enum: [code, identity, http] - + dedicated_worker: + type: boolean required: - args diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 0153d51419..46991e58db 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -1468,6 +1468,7 @@ struct Preview { args: Option>, language: Option, tag: Option, + dedicated_worker: Option, } #[derive(Deserialize)] @@ -2262,6 +2263,7 @@ async fn run_preview_job( concurrent_limit: None, // TODO(gbouv): once I find out how to store limits in the content of a script, should be easy to plug limits here concurrency_time_window_s: None, // TODO(gbouv): same as above cache_ttl: None, + dedicated_worker: preview.dedicated_worker, }), }, preview.args.unwrap_or_default(), diff --git a/backend/windmill-common/src/jobs.rs b/backend/windmill-common/src/jobs.rs index c9fde19184..09a84d0b2c 100644 --- a/backend/windmill-common/src/jobs.rs +++ b/backend/windmill-common/src/jobs.rs @@ -315,6 +315,7 @@ pub struct RawCode { pub concurrent_limit: Option, pub concurrency_time_window_s: Option, pub cache_ttl: Option, + pub dedicated_worker: Option, } type Tag = String; diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index f8faba8a54..55cc46e1fe 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -209,10 +209,14 @@ pub async fn load_worker_config(db: &DB) -> error::Result { .worker_tags .or_else(|| { if let Some(ref dedicated_worker) = dedicated_worker.as_ref() { - Some(vec![format!( + let mut dedi_tags = vec![format!( "{}:{}", dedicated_worker.workspace_id, dedicated_worker.path - )]) + )]; + if std::env::var("ADD_FLOW_TAG").is_ok() { + dedi_tags.push("flow".to_string()); + } + Some(dedi_tags) } else { std::env::var("WORKER_TAGS") .ok() @@ -220,6 +224,7 @@ pub async fn load_worker_config(db: &DB) -> error::Result { } }) .unwrap_or_else(|| DEFAULT_TAGS.clone()); + let mut priority_tags_sorted: Vec = Vec::new(); let priority_tags_map = config.priority_tags.unwrap_or_else(HashMap::new); if priority_tags_map.len() > 0 { diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index d81e3a0915..6b55ef7971 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -1453,7 +1453,7 @@ async fn pull_single_job_and_mark_as_running_no_concurrency_limit< LIMIT 1 ) RETURNING *") - .bind(tags.clone()) + .bind(tags) .fetch_optional(db) .await? } else { @@ -2208,6 +2208,7 @@ pub async fn push<'c, T: Serialize + Send + Sync, R: rsmq_async::RsmqConnection concurrent_limit, concurrency_time_window_s, cache_ttl, + dedicated_worker, }) => ( None, path, @@ -2219,7 +2220,7 @@ pub async fn push<'c, T: Serialize + Send + Sync, R: rsmq_async::RsmqConnection concurrent_limit, concurrency_time_window_s, cache_ttl, - None, + dedicated_worker, None, ), JobPayload::Dependencies { hash, dependencies, language, path } => ( diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index f9d6ffb12f..ee05c3657d 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -577,28 +577,37 @@ pub async fn start_worker( let _ = write_file(job_dir, "package.json", &splitted[0]).await?; let lockb = splitted[1]; if lockb != EMPTY_FILE { - let _ = write_file_binary( + let has_trusted_deps = &splitted[0].contains("trustedDependencies"); + + if !has_trusted_deps { + let _ = write_file_binary( + job_dir, + "bun.lockb", + &base64::engine::general_purpose::STANDARD + .decode(&splitted[1]) + .map_err(|_| { + error::Error::InternalErr("Could not decode bun.lockb".to_string()) + })?, + ) + .await?; + } + + install_lockfile( + &mut logs, + &mut mem_peak, + &Uuid::nil(), + &w_id, + db, job_dir, - "bun.lockb", - &base64::engine::general_purpose::STANDARD - .decode(&splitted[1]) - .map_err(|_| { - error::Error::InternalErr("Could not decode bun.lockb".to_string()) - })?, + worker_name, + common_bun_proc_envs.clone(), ) .await?; + if !has_trusted_deps { + remove_dir_all(format!("{}/node_modules", job_dir)).await?; + } + tracing::info!("dedicated worker requirements installed: {reqs}"); } - install_lockfile( - &mut logs, - &mut mem_peak, - &Uuid::nil(), - &w_id, - db, - job_dir, - worker_name, - common_bun_proc_envs.clone(), - ) - .await?; } else if !*DISABLE_NSJAIL { let trusted_deps = get_trusted_deps(inner_content); logs.push_str("\n\n--- BUN INSTALL ---\n"); diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index e43248d279..53915d3f61 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -8,7 +8,7 @@ use tokio::{ io::{AsyncBufReadExt, AsyncWriteExt, BufReader}, process::Command, }; -use windmill_common::{error, jobs::QueuedJob, variables}; +use windmill_common::{error, jobs::QueuedJob, variables, worker::to_raw_value}; use std::{collections::VecDeque, process::Stdio, sync::Arc}; @@ -82,8 +82,15 @@ pub async fn handle_dedicated_process( .take() .expect("child did not have a handle to stdout"); + let stderr = child + .stderr + .take() + .expect("child did not have a handle to stderr"); + let mut reader = BufReader::new(stdout).lines(); + let mut err_reader = BufReader::new(stderr).lines(); + let mut stdin = child .stdin .take() @@ -115,6 +122,14 @@ pub async fn handle_dedicated_process( tracing::info!("Could not write end message to stdin: {e:?}") } }, + line = err_reader.next_line() => { + if let Some(line) = line.expect("line is ok") { + tracing::error!("dedicated worker process stderr: {:?}", line); + } else { + tracing::info!("dedicated worker process exited"); + break; + } + }, line = reader.next_line() => { // j += 1; @@ -123,11 +138,16 @@ pub async fn handle_dedicated_process( tracing::info!("dedicated worker process started"); continue; } - tracing::debug!("processed job"); + tracing::debug!("processed job: {line}"); - let result = serde_json::from_str(&line).expect("json is ok"); let job: Arc = jobs.pop_front().expect("pop"); - job_completed_tx.send(JobCompleted { job , result, logs: "".to_string(), mem_peak: 0, success: true, cached_res_path: None, token: token.to_string() }).await.unwrap(); + match serde_json::from_str::>(&line) { + Ok(result) => job_completed_tx.send(JobCompleted { job , result, logs: "".to_string(), mem_peak: 0, success: true, cached_res_path: None, token: token.to_string() }).await.unwrap(), + Err(e) => { + tracing::error!("Could not deserialize job result `{line}`: {e:?}"); + job_completed_tx.send(JobCompleted { job , result: to_raw_value(&serde_json::json!({"error": format!("Could not deserialize job result `{line}`: {e:?}")})), logs: "".to_string(), mem_peak: 0, success: false, cached_res_path: None, token: token.to_string() }).await.unwrap(); + }, + }; } else { tracing::info!("dedicated worker process exited"); break; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 7fb6939e39..3450dedb25 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -35,7 +35,8 @@ use windmill_common::{ users::SUPERADMIN_SECRET_EMAIL, utils::{rd_string, StripPath}, worker::{ - to_raw_value, to_raw_value_owned, update_ping, CLOUD_HOSTED, WORKER_CONFIG, WORKER_GROUP, + to_raw_value, to_raw_value_owned, update_ping, WorkspacedPath, CLOUD_HOSTED, WORKER_CONFIG, + WORKER_GROUP, }, DB, IS_READY, METRICS_DEBUG_ENABLED, METRICS_ENABLED, }; @@ -1079,7 +1080,7 @@ pub async fn run_worker>>, Option>) + (None, None, None) + as ( + Option, + Option>>, + Option>, + ) }; #[cfg(feature = "benchmark")] @@ -1447,28 +1455,35 @@ pub async fn run_worker( concurrent_limit: None, concurrency_time_window_s: None, cache_ttl: None, + dedicated_worker: None, }), PushArgs::empty(), worker_name, diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index b601317533..c5850ab746 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -1426,7 +1426,9 @@ async fn push_next_flow_job match json_value { Ok(serde_json::Value::Number(n)) => { if !n.is_u64() { - return Err(Error::ExecutionErr(format!("Expected an integer, found: {n}"))); + return Err(Error::ExecutionErr(format!( + "Expected an integer, found: {n}" + ))); } n.as_u64().map(|x| from_now(Duration::from_secs(x))) @@ -2554,6 +2556,7 @@ fn raw_script_to_payload( concurrent_limit: *concurrent_limit, concurrency_time_window_s: *concurrency_time_window_s, cache_ttl: module.cache_ttl.map(|x| x as i32), + dedicated_worker: None, }), tag: tag.clone(), } diff --git a/frontend/src/lib/components/sidebar/UserMenu.svelte b/frontend/src/lib/components/sidebar/UserMenu.svelte index 905b8fda4a..185ee1c767 100644 --- a/frontend/src/lib/components/sidebar/UserMenu.svelte +++ b/frontend/src/lib/components/sidebar/UserMenu.svelte @@ -1,7 +1,7 @@