From 39f30785a04a54c651e532d7ede3b8c17cdec7ea Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Mon, 23 Oct 2023 20:21:16 +0200 Subject: [PATCH] feat: dedicated worker for python (#2492) * update * all --- backend/windmill-api/src/oauth2.rs | 2 +- backend/windmill-common/src/worker.rs | 15 +- backend/windmill-worker/src/bun_executor.rs | 183 +----- backend/windmill-worker/src/common.rs | 1 + .../windmill-worker/src/dedicated_worker.rs | 153 +++++ backend/windmill-worker/src/lib.rs | 1 + .../windmill-worker/src/python_executor.rs | 556 +++++++++++++----- backend/windmill-worker/src/worker.rs | 68 ++- .../src/lib/components/ScriptBuilder.svelte | 3 +- .../(root)/(logged)/run/[...run]/+page.svelte | 2 +- 10 files changed, 640 insertions(+), 344 deletions(-) diff --git a/backend/windmill-api/src/oauth2.rs b/backend/windmill-api/src/oauth2.rs index 5162d70f7d..24a1864208 100644 --- a/backend/windmill-api/src/oauth2.rs +++ b/backend/windmill-api/src/oauth2.rs @@ -287,7 +287,7 @@ pub fn build_oauth_clients( }) .flatten(); let all_clients = AllClients { logins, connects, slack }; - tracing::info!("Final oauth config: {all_clients:#?}"); + tracing::debug!("Final oauth config: {all_clients:#?}"); Ok(all_clients) } diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index 96fe028b50..7cf0d35dfe 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -144,6 +144,7 @@ pub async fn update_ping(worker_instance: &str, worker_name: &str, ip: &str, db: } pub async fn load_worker_config(db: &DB) -> error::Result { + tracing::info!("Loading config from WORKER_GROUP: {}", *WORKER_GROUP); let mut config: WorkerConfigOpt = sqlx::query_scalar!( "SELECT config FROM config WHERE name = $1", format!("worker__{}", *WORKER_GROUP) @@ -155,7 +156,19 @@ pub async fn load_worker_config(db: &DB) -> error::Result { .flatten() .unwrap_or_default(); if config.dedicated_worker.is_none() { - config.dedicated_worker = std::env::var("DEDICATED_WORKER").ok(); + let dw = std::env::var("DEDICATED_WORKER").ok(); + if dw.is_some() { + tracing::info!( + "DEDICATED_WORKER set from env variable: {}", + dw.as_ref().unwrap() + ); + config.dedicated_worker = dw; + } + } else { + tracing::info!( + "DEDICATED_WORKER set from config: {}", + config.dedicated_worker.as_ref().unwrap() + ); } let dedicated_worker = config.dedicated_worker.map(|x| { let splitted = x.split(':').to_owned().collect_vec(); diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 2940398ea8..a2feccec46 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -1,11 +1,5 @@ use std::{collections::HashMap, process::Stdio}; -#[cfg(feature = "enterprise")] -use std::collections::VecDeque; - -#[cfg(feature = "enterprise")] -use anyhow::Context; - use base64::Engine; use itertools::Itertools; use regex::Regex; @@ -24,17 +18,11 @@ use crate::{ NPM_CONFIG_REGISTRY, NSJAIL_PATH, PATH_ENV, TZ_ENV, }; -#[cfg(feature = "enterprise")] -use crate::MAX_BUFFERED_DEDICATED_JOBS; - use tokio::{ fs::{remove_dir_all, File}, process::Command, }; -#[cfg(feature = "enterprise")] -use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; - use tokio::io::AsyncReadExt; #[cfg(feature = "enterprise")] @@ -554,12 +542,10 @@ pub async fn start_worker( script_path: &str, token: &str, job_completed_tx: Sender, - mut jobs_rx: Receiver>, - mut killpill_rx: tokio::sync::broadcast::Receiver<()>, + jobs_rx: Receiver>, + killpill_rx: tokio::sync::broadcast::Receiver<()>, ) -> Result<()> { - use std::task::Poll; - - use futures::{future, Future}; + use crate::dedicated_worker::handle_dedicated_process; let mut logs = "".to_string(); let mut mem_peak: i32 = 0; @@ -579,9 +565,8 @@ pub async fn start_worker( None, None, ) - .await - .to_vec(); - let context_envs = build_envs_map(context); + .await; + let context_envs = build_envs_map(context.to_vec()); if let Some(reqs) = requirements_o { let splitted = reqs.split(BUN_LOCKB_SPLIT).collect::>(); if splitted.len() != 2 { @@ -694,21 +679,6 @@ for await (const chunk of Bun.stdin.stream()) {{ write_file(job_dir, "wrapper.ts", &wrapper_content).await?; } - let reserved_variables = windmill_common::variables::get_reserved_variables( - w_id, - token, - "dedicated_worker", - "dedicated_worker", - Uuid::nil().to_string().as_str(), - "dedicted_worker", - Some(script_path.to_string()), - None, - None, - None, - None, - ) - .await; - let _ = write_file( &job_dir, "loader.bun.ts", @@ -729,136 +699,25 @@ plugin(p) ) .await?; - //do not cache local dependencies - let mut child = { - let script_path = format!("{job_dir}/wrapper.ts"); - let args = vec![ + handle_dedicated_process( + &*BUN_PATH, + job_dir, + context_envs, + envs, + context, + common_bun_proc_envs, + vec![ "run", "-i", "--prefer-offline", "-r", "./loader.bun.ts", - &script_path, - ]; - let mut bun_cmd = Command::new(&*BUN_PATH); - bun_cmd - .current_dir(job_dir) - .env_clear() - .envs(context_envs) - .envs(envs) - .envs( - reserved_variables - .iter() - .map(|x| (x.name.clone(), x.value.clone())) - .collect::>(), - ) - .envs(common_bun_proc_envs) - .args(args) - .stdin(Stdio::piped()) - .stdout(Stdio::piped()) - .stderr(Stdio::piped()); - start_child_process(bun_cmd, &*BUN_PATH).await? - }; - - let stdout = child - .stdout - .take() - .expect("child did not have a handle to stdout"); - - let mut reader = BufReader::new(stdout).lines(); - - let mut stdin = child - .stdin - .take() - .expect("child did not have a handle to stdin"); - - // Ensure the child process is spawned in the runtime so it can - // make progress on its own while we await for any output. - let child = tokio::spawn(async move { - let status = child - .wait() - .await - .expect("child process encountered an error"); - - println!("child status was: {}", status); - }); - - let mut jobs = VecDeque::with_capacity(MAX_BUFFERED_DEDICATED_JOBS); - // let mut i = 0; - // let mut j = 0; - let mut alive = true; - - fn conditional_polling( - fut: impl Future, - predicate: bool, - ) -> impl Future { - let mut fut = Box::pin(fut); - future::poll_fn(move |cx| { - if predicate { - fut.as_mut().poll(cx) - } else { - Poll::Pending - } - }) - } - - loop { - tokio::select! { - biased; - _ = killpill_rx.recv(), if alive => { - println!("received killpill for dedicated worker"); - alive = false; - if let Err(e) = write_stdin(&mut stdin, "end").await { - tracing::info!("Could not write end message to stdin: {e:?}") - } - }, - line = reader.next_line() => { - // j += 1; - - if let Some(line) = line.expect("line is ok") { - if line == "start" { - tracing::info!("dedicated worker process started"); - continue; - } - tracing::debug!("processed job"); - - 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(); - } else { - tracing::info!("dedicated worker process exited"); - break; - } - }, - job = conditional_polling(jobs_rx.recv(), alive && jobs.len() < MAX_BUFFERED_DEDICATED_JOBS) => { - // i += 1; - if let Some(job) = job { - tracing::debug!("received job"); - jobs.push_back(job.clone()); - // write_stdin(&mut stdin, &serde_json::to_string(&job.args.unwrap_or_else(|| serde_json::json!({"x": job.id}))).expect("serialize")).await?; - write_stdin(&mut stdin, &serde_json::to_string(&job.args).expect("serialize")).await?; - stdin.flush().await.context("stdin flush")?; - } else { - tracing::debug!("job channel closed"); - alive = false; - if let Err(e) = write_stdin(&mut stdin, "end").await { - tracing::error!("Could not write end message to stdin: {e:?}") - } - } - } - } - } - - child - .await - .map_err(|e| anyhow::anyhow!("child process encountered an error: {e}"))?; - tracing::info!("dedicated worker child process exited successfully"); - Ok(()) -} - -#[cfg(feature = "enterprise")] -async fn write_stdin(stdin: &mut tokio::process::ChildStdin, s: &str) -> error::Result<()> { - let _ = &stdin.write_all(format!("{s}\n").as_bytes()).await?; - stdin.flush().await.context("stdin flush")?; - Ok(()) + &format!("{job_dir}/wrapper.ts"), + ], + killpill_rx, + job_completed_tx, + token, + jobs_rx, + ) + .await } diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 13d2dce1b6..f53b6622c5 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -17,6 +17,7 @@ use windmill_common::{ }; use anyhow::Result; + use std::{ borrow::Borrow, collections::{hash_map::DefaultHasher, HashMap}, diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index c0ac64fd4c..0571da6e73 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -3,3 +3,156 @@ // pub fn create_dedicated_worker() { // let (job_completed_tx, mut new_job) = mpsc::channel::(100); // } + +use tokio::{ + io::{AsyncBufReadExt, AsyncWriteExt, BufReader}, + process::Command, +}; +use windmill_common::{error, jobs::QueuedJob, variables}; + +use std::{collections::VecDeque, process::Stdio, sync::Arc}; + +use anyhow::Context; + +use crate::{common::start_child_process, JobCompleted, MAX_BUFFERED_DEDICATED_JOBS}; + +use futures::{future, Future}; +use std::{collections::HashMap, task::Poll}; + +use tokio::sync::mpsc::{Receiver, Sender}; + +fn conditional_polling( + fut: impl Future, + predicate: bool, +) -> impl Future { + let mut fut = Box::pin(fut); + future::poll_fn(move |cx| { + if predicate { + fut.as_mut().poll(cx) + } else { + Poll::Pending + } + }) +} + +async fn write_stdin(stdin: &mut tokio::process::ChildStdin, s: &str) -> error::Result<()> { + let _ = &stdin.write_all(format!("{s}\n").as_bytes()).await?; + stdin.flush().await.context("stdin flush")?; + Ok(()) +} + +pub async fn handle_dedicated_process( + command_path: &String, + job_dir: &str, + context_envs: HashMap, + envs: HashMap, + reserved_variables: [variables::ContextualVariable; 15], + common_bun_proc_envs: HashMap, + args: Vec<&str>, + mut killpill_rx: tokio::sync::broadcast::Receiver<()>, + job_completed_tx: Sender, + token: &str, + mut jobs_rx: Receiver>, +) -> std::result::Result<(), error::Error> { + //do not cache local dependencies + let mut child = { + let mut cmd = Command::new(command_path); + cmd.current_dir(job_dir) + .env_clear() + .envs(context_envs) + .envs(envs) + .envs( + reserved_variables + .iter() + .map(|x| (x.name.clone(), x.value.clone())) + .collect::>(), + ) + .envs(common_bun_proc_envs) + .args(args) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); + start_child_process(cmd, command_path).await? + }; + + let stdout = child + .stdout + .take() + .expect("child did not have a handle to stdout"); + + let mut reader = BufReader::new(stdout).lines(); + + let mut stdin = child + .stdin + .take() + .expect("child did not have a handle to stdin"); + + // Ensure the child process is spawned in the runtime so it can + // make progress on its own while we await for any output. + let child = tokio::spawn(async move { + let status = child + .wait() + .await + .expect("child process encountered an error"); + + println!("child status was: {}", status); + }); + + let mut jobs = VecDeque::with_capacity(MAX_BUFFERED_DEDICATED_JOBS); + // let mut i = 0; + // let mut j = 0; + let mut alive = true; + + loop { + tokio::select! { + biased; + _ = killpill_rx.recv(), if alive => { + println!("received killpill for dedicated worker"); + alive = false; + if let Err(e) = write_stdin(&mut stdin, "end").await { + tracing::info!("Could not write end message to stdin: {e:?}") + } + }, + line = reader.next_line() => { + // j += 1; + + if let Some(line) = line.expect("line is ok") { + if line == "start" { + tracing::info!("dedicated worker process started"); + continue; + } + tracing::debug!("processed job"); + + 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(); + } else { + tracing::info!("dedicated worker process exited"); + break; + } + }, + job = conditional_polling(jobs_rx.recv(), alive && jobs.len() < MAX_BUFFERED_DEDICATED_JOBS) => { + // i += 1; + if let Some(job) = job { + tracing::debug!("received job"); + jobs.push_back(job.clone()); + // write_stdin(&mut stdin, &serde_json::to_string(&job.args.unwrap_or_else(|| serde_json::json!({"x": job.id}))).expect("serialize")).await?; + write_stdin(&mut stdin, &serde_json::to_string(&job.args).expect("serialize")).await?; + stdin.flush().await.context("stdin flush")?; + } else { + tracing::debug!("job channel closed"); + alive = false; + if let Err(e) = write_stdin(&mut stdin, "end").await { + tracing::error!("Could not write end message to stdin: {e:?}") + } + } + } + } + } + + child + .await + .map_err(|e| anyhow::anyhow!("child process encountered an error: {e}"))?; + tracing::info!("dedicated worker child process exited successfully"); + Ok(()) +} diff --git a/backend/windmill-worker/src/lib.rs b/backend/windmill-worker/src/lib.rs index cd0c1d07bd..3857dd196a 100644 --- a/backend/windmill-worker/src/lib.rs +++ b/backend/windmill-worker/src/lib.rs @@ -7,6 +7,7 @@ mod bash_executor; mod bun_executor; pub mod common; mod config; +#[cfg(feature = "enterprise")] mod dedicated_worker; mod deno_executor; mod global_cache; diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 4daf2e73db..2671227880 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -15,6 +15,7 @@ use windmill_common::{ jobs::QueuedJob, utils::calculate_hash, worker::WORKER_CONFIG, + DB, }; lazy_static::lazy_static! { @@ -175,168 +176,38 @@ pub async fn handle_python_job( base_internal_url: &str, envs: HashMap, ) -> windmill_common::error::Result> { - create_dependencies_dir(job_dir).await; + let script_path = job.script_path(); + let additional_python_paths = handle_python_deps( + job_dir, + requirements_o, + inner_content, + &job.workspace_id, + script_path, + &job.id, + db, + worker_name, + worker_dir, + logs, + mem_peak, + ) + .await?; - let mut additional_python_paths: Vec = WORKER_CONFIG - .read() - .await - .additional_python_paths - .clone() - .unwrap_or_else(|| vec![]) - .clone(); - - let requirements = match requirements_o { - Some(r) => r, - None => { - let requirements = windmill_parser_py_imports::parse_python_imports( - &inner_content, - &job.workspace_id, - &job.script_path(), - &db, - ) - .await? - .join("\n"); - if requirements.is_empty() { - "".to_string() - } else { - pip_compile( - &job.id, - &requirements, - logs, - mem_peak, - job_dir, - db, - worker_name, - &job.workspace_id, - ) - .await - .map_err(|e| { - Error::ExecutionErr(format!("pip compile failed: {}", e.to_string())) - })? - } - } - }; - - if requirements.len() > 0 { - additional_python_paths = handle_python_reqs( - requirements - .split("\n") - .filter(|x| !x.starts_with("--")) - .collect(), - &job.id, - &job.workspace_id, - logs, - mem_peak, - db, - worker_name, - job_dir, - worker_dir, - ) - .await?; - } logs.push_str("\n\n--- PYTHON CODE EXECUTION ---\n"); set_logs(logs, &job.id, db).await; - let relative_imports = RELATIVE_IMPORT_REGEX.is_match(&inner_content); + let ( + import_loader, + import_base64, + import_datetime, + module_dir_dot, + dirs, + last, + transforms, + spread, + ) = prepare_wrapper(job_dir, inner_content, script_path).await?; - let script_path_splitted = &job.script_path().split("/"); - let dirs_full = script_path_splitted - .clone() - .take(script_path_splitted.clone().count() - 1) - .join("/") - .replace("-", "_") - .replace("@", "."); - let dirs = if dirs_full.len() > 0 { - dirs_full - .strip_prefix("/") - .unwrap_or(&dirs_full) - .to_string() - } else { - "tmp".to_string() - }; - let last = script_path_splitted - .clone() - .last() - .unwrap() - .replace("-", "_") - .replace(" ", "_") - .to_lowercase(); - let module_dir = format!("{}/{}", job_dir, dirs); - tokio::fs::create_dir_all(format!("{module_dir}/")).await?; - let _ = write_file(&module_dir, &format!("{last}.py"), inner_content).await?; - if relative_imports { - let _ = write_file(&job_dir, "loader.py", RELATIVE_PYTHON_LOADER).await?; - } - - let sig = windmill_parser_py::parse_python_signature(inner_content)?; - let transforms = sig - .args - .iter() - .map(|x| match x.typ { - windmill_parser::Typ::Bytes => { - let name = &x.name; - format!( - "if \"{name}\" in kwargs and kwargs[\"{name}\"] is not None:\n \ - kwargs[\"{name}\"] = base64.b64decode(kwargs[\"{name}\"])\n", - ) - } - windmill_parser::Typ::Datetime => { - let name = &x.name; - format!( - "if \"{name}\" in kwargs and kwargs[\"{name}\"] is not None:\n \ - kwargs[\"{name}\"] = datetime.fromisoformat(kwargs[\"{name}\"])\n", - ) - } - _ => "".to_string(), - }) - .collect::>() - .join(""); create_args_and_out_file(&client, job, job_dir, db).await?; - let import_loader = if relative_imports { - "import loader" - } else { - "" - }; - let import_base64 = if sig - .args - .iter() - .any(|x| x.typ == windmill_parser::Typ::Bytes) - { - "import base64" - } else { - "" - }; - let import_datetime = if sig - .args - .iter() - .any(|x| x.typ == windmill_parser::Typ::Datetime) - { - "from datetime import datetime" - } else { - "" - }; - let spread = if sig.star_kwargs { - "args = kwargs".to_string() - } else { - sig.args - .into_iter() - .map(|x| { - let name = &x.name; - if x.default.is_none() { - format!("args[\"{name}\"] = kwargs.get(\"{name}\")") - } else { - format!( - r#"args["{name}"] = kwargs.get("{name}") -if args["{name}"] is None: - del args["{name}"]"# - ) - } - }) - .join("\n") - }; - - let module_dir_dot = dirs.replace("/", ".").replace("-", "_"); let wrapper_content: String = format!( r#" import json @@ -394,6 +265,7 @@ except Exception as e: let client = client.get_authed().await; let mut reserved_variables = get_reserved_variables(job, &client.token, db).await?; let additional_python_paths_folders = additional_python_paths.iter().join(":"); + if !*DISABLE_NSJAIL { let shared_deps = additional_python_paths .into_iter() @@ -446,6 +318,7 @@ mount {{ .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) + .env("BASE_URL", base_internal_url) .args(vec![ "--config", "run.config.proto", @@ -491,6 +364,207 @@ mount {{ read_result(job_dir).await } +async fn prepare_wrapper( + job_dir: &str, + inner_content: &str, + script_path: &str, +) -> error::Result<( + &'static str, + &'static str, + &'static str, + String, + String, + String, + String, + String, +)> { + let relative_imports = RELATIVE_IMPORT_REGEX.is_match(&inner_content); + + let script_path_splitted = script_path.split("/"); + let dirs_full = script_path_splitted + .clone() + .take(script_path_splitted.clone().count() - 1) + .join("/") + .replace("-", "_") + .replace("@", "."); + let dirs = if dirs_full.len() > 0 { + dirs_full + .strip_prefix("/") + .unwrap_or(&dirs_full) + .to_string() + } else { + "tmp".to_string() + }; + let last = script_path_splitted + .clone() + .last() + .unwrap() + .replace("-", "_") + .replace(" ", "_") + .to_lowercase(); + let module_dir = format!("{}/{}", job_dir, dirs); + tokio::fs::create_dir_all(format!("{module_dir}/")).await?; + let _ = write_file(&module_dir, &format!("{last}.py"), inner_content).await?; + if relative_imports { + let _ = write_file(job_dir, "loader.py", RELATIVE_PYTHON_LOADER).await?; + } + + let sig = windmill_parser_py::parse_python_signature(inner_content)?; + let transforms = sig + .args + .iter() + .map(|x| match x.typ { + windmill_parser::Typ::Bytes => { + let name = &x.name; + format!( + "if \"{name}\" in kwargs and kwargs[\"{name}\"] is not None:\n \ + kwargs[\"{name}\"] = base64.b64decode(kwargs[\"{name}\"])\n", + ) + } + windmill_parser::Typ::Datetime => { + let name = &x.name; + format!( + "if \"{name}\" in kwargs and kwargs[\"{name}\"] is not None:\n \ + kwargs[\"{name}\"] = datetime.fromisoformat(kwargs[\"{name}\"])\n", + ) + } + _ => "".to_string(), + }) + .collect::>() + .join(""); + + let import_loader = if relative_imports { + "import loader" + } else { + "" + }; + let import_base64 = if sig + .args + .iter() + .any(|x| x.typ == windmill_parser::Typ::Bytes) + { + "import base64" + } else { + "" + }; + let import_datetime = if sig + .args + .iter() + .any(|x| x.typ == windmill_parser::Typ::Datetime) + { + "from datetime import datetime" + } else { + "" + }; + let spread = if sig.star_kwargs { + "args = kwargs".to_string() + } else { + sig.args + .into_iter() + .map(|x| { + let name = &x.name; + if x.default.is_none() { + format!("args[\"{name}\"] = kwargs.get(\"{name}\")") + } else { + format!( + r#"args["{name}"] = kwargs.get("{name}") +if args["{name}"] is None: + del args["{name}"]"# + ) + } + }) + .join("\n") + }; + + let module_dir_dot = dirs.replace("/", ".").replace("-", "_"); + + Ok(( + import_loader, + import_base64, + import_datetime, + module_dir_dot, + dirs, + last, + transforms, + spread, + )) +} + +async fn handle_python_deps( + job_dir: &str, + requirements_o: Option, + inner_content: &str, + w_id: &str, + script_path: &str, + job_id: &Uuid, + db: &DB, + worker_name: &str, + worker_dir: &str, + logs: &mut String, + mem_peak: &mut i32, +) -> error::Result> { + create_dependencies_dir(job_dir).await; + + let mut additional_python_paths: Vec = WORKER_CONFIG + .read() + .await + .additional_python_paths + .clone() + .unwrap_or_else(|| vec![]) + .clone(); + + let requirements = match requirements_o { + Some(r) => r, + None => { + let requirements = windmill_parser_py_imports::parse_python_imports( + inner_content, + w_id, + script_path, + db, + ) + .await? + .join("\n"); + if requirements.is_empty() { + "".to_string() + } else { + pip_compile( + job_id, + &requirements, + logs, + mem_peak, + job_dir, + db, + worker_name, + w_id, + ) + .await + .map_err(|e| { + Error::ExecutionErr(format!("pip compile failed: {}", e.to_string())) + })? + } + } + }; + + if requirements.len() > 0 { + additional_python_paths = handle_python_reqs( + requirements + .split("\n") + .filter(|x| !x.starts_with("--")) + .collect(), + job_id, + w_id, + logs, + mem_peak, + db, + worker_name, + job_dir, + worker_dir, + ) + .await?; + } + Ok(additional_python_paths) +} + pub async fn handle_python_reqs( requirements: Vec<&str>, job_id: &Uuid, @@ -675,3 +749,175 @@ pub async fn handle_python_reqs( } Ok(req_paths) } + +#[cfg(feature = "enterprise")] +use std::sync::Arc; +#[cfg(feature = "enterprise")] +use tokio::sync::mpsc::Sender; + +#[cfg(feature = "enterprise")] +use crate::{common::build_envs_map, dedicated_worker::handle_dedicated_process, JobCompleted}; +#[cfg(feature = "enterprise")] +use tokio::sync::mpsc::Receiver; +#[cfg(feature = "enterprise")] +use windmill_common::variables; + +#[cfg(feature = "enterprise")] +pub async fn start_worker( + requirements_o: Option, + db: &sqlx::Pool, + inner_content: &str, + base_internal_url: &str, + job_dir: &str, + worker_name: &str, + envs: HashMap, + w_id: &str, + script_path: &str, + token: &str, + job_completed_tx: Sender, + jobs_rx: Receiver>, + killpill_rx: tokio::sync::broadcast::Receiver<()>, +) -> error::Result<()> { + let mut logs = "".to_string(); + let mut mem_peak: i32 = 0; + let context = variables::get_reserved_variables( + w_id, + &token, + "dedicated_worker@windmill.dev", + "dedicated_worker", + "NOT_AVAILABLE", + "dedicated_worker", + Some(script_path.to_string()), + None, + None, + None, + None, + ) + .await + .to_vec(); + + let context_envs = build_envs_map(context); + let additional_python_paths = handle_python_deps( + job_dir, + requirements_o, + inner_content, + w_id, + script_path, + &Uuid::nil(), + db, + worker_name, + job_dir, + &mut logs, + &mut mem_peak, + ) + .await?; + + logs.push_str("\n\n--- PYTHON CODE EXECUTION ---\n"); + set_logs(&mut logs, &Uuid::nil(), db).await; + + let ( + import_loader, + import_base64, + import_datetime, + module_dir_dot, + _dirs, + last, + transforms, + spread, + ) = prepare_wrapper(job_dir, inner_content, script_path).await?; + + { + // logs.push_str(format!("infer args: {:?}\n", start.elapsed().as_micros()).as_str()); + // we cannot use Bun.read and Bun.write because it results in an EBADF error on cloud + let wrapper_content: String = format!( + r#" +import json +{import_loader} +{import_base64} +{import_datetime} +import traceback +import sys +from {module_dir_dot} import {last} as inner_script +import re + + +def to_b_64(v: bytes): + import base64 + b64 = base64.b64encode(v) + return b64.decode('ascii') + +replace_nan = re.compile(r'\bNaN\b') +sys.stdout.write('start\n') + +for line in sys.stdin: + kwargs = json.loads(line, strict=False) + args = {{}} + {transforms} + {spread} + for k, v in list(args.items()): + if v == '': + del args[k] + + try: + res = inner_script.main(**args) + typ = type(res) + if typ.__name__ == 'DataFrame': + if typ.__module__ == 'pandas.core.frame': + res = res.values.tolist() + elif typ.__module__ == 'polars.dataframe.frame': + res = res.rows() + elif typ.__name__ == 'bytes': + res = to_b_64(res) + elif typ.__name__ == 'dict': + for k, v in res.items(): + if type(v).__name__ == 'bytes': + res[k] = to_b_64(v) + res_json = re.sub(replace_nan, ' null ', json.dumps(res, separators=(',', ':'), default=str).replace('\n', '')) + sys.stdout.write(res_json + "\n") + except Exception as e: + exc_type, exc_value, exc_traceback = sys.exc_info() + tb = traceback.format_tb(exc_traceback) + err_json = json.dumps({{ "error": {{ "message": str(e), "name": e.__class__.__name__, "stack": '\n'.join(tb[1:]) }} }}, separators=(',', ':'), default=str).replace('\n', '') + sys.stdout.write(err_json + "\n") + sys.stdout.flush() +"#, + ); + write_file(job_dir, "wrapper.py", &wrapper_content).await?; + } + + let reserved_variables = windmill_common::variables::get_reserved_variables( + w_id, + token, + "dedicated_worker", + "dedicated_worker", + Uuid::nil().to_string().as_str(), + "dedicted_worker", + Some(script_path.to_string()), + None, + None, + None, + None, + ) + .await; + + let mut proc_envs = HashMap::new(); + let additional_python_paths_folders = additional_python_paths.iter().join(":"); + proc_envs.insert("PYTHONPATH".to_string(), additional_python_paths_folders); + proc_envs.insert("PATH".to_string(), PATH_ENV.to_string()); + proc_envs.insert("TZ".to_string(), TZ_ENV.to_string()); + proc_envs.insert("BASE_URL".to_string(), base_internal_url.to_string()); + handle_dedicated_process( + &*PYTHON_PATH, + job_dir, + context_envs, + envs, + reserved_variables, + proc_envs, + ["-u", "-m", "wrapper"].to_vec(), + killpill_rx, + job_completed_tx, + token, + jobs_rx, + ) + .await +} diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index e8dbc437b7..1f668a4d44 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -30,7 +30,9 @@ use windmill_common::{ scripts::{get_full_hub_script_by_path, ScriptHash, ScriptLang}, users::SUPERADMIN_SECRET_EMAIL, utils::{rd_string, StripPath}, - worker::{to_raw_value, to_raw_value_owned, update_ping, CLOUD_HOSTED, WORKER_CONFIG}, + worker::{ + to_raw_value, to_raw_value_owned, update_ping, CLOUD_HOSTED, WORKER_CONFIG, WORKER_GROUP, + }, DB, IS_READY, METRICS_ENABLED, }; use windmill_queue::{ @@ -62,9 +64,6 @@ use crate::global_cache::{ copy_denogo_cache_from_bucket_as_tar, copy_tmp_cache_to_cache, }; -#[cfg(feature = "enterprise")] -use crate::bun_executor::start_worker; - use windmill_queue::{add_completed_job, add_completed_job_error}; use crate::{ @@ -873,7 +872,7 @@ pub async fn run_worker, Option, Option>)>( @@ -961,23 +960,46 @@ pub async fn run_worker { + crate::python_executor::start_worker( + lock, + &db, + &content, + &base_internal_url, + &job_dir, + &worker_name, + worker_envs, + &_wp.workspace_id, + &_wp.path, + &token, + job_completed_tx, + dedicated_worker_rx, + killpill_rx, + ) + .await + } + Some(ScriptLang::Bun) => { + crate::bun_executor::start_worker( + lock, + &db, + &content, + &base_internal_url, + &job_dir, + &worker_name, + worker_envs, + &_wp.workspace_id, + &_wp.path, + &token, + job_completed_tx, + dedicated_worker_rx, + killpill_rx, + ) + .await + } + _ => unreachable!("Non supported language for dedicated worker"), + } { tracing::error!("error in dedicated worker: {:?}", e) } }); diff --git a/frontend/src/lib/components/ScriptBuilder.svelte b/frontend/src/lib/components/ScriptBuilder.svelte index 1077a8a66d..8c250d034d 100644 --- a/frontend/src/lib/components/ScriptBuilder.svelte +++ b/frontend/src/lib/components/ScriptBuilder.svelte @@ -570,7 +570,8 @@ { diff --git a/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte b/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte index 179a44ab88..87169aba7d 100644 --- a/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte @@ -298,7 +298,7 @@ {/if} {#if job.tag && !['deno', 'python3', 'flow', 'other', 'go', 'postgresql', 'mysql', 'bigquery', 'snowflake', 'graphql', 'nativets', 'bash', 'powershell', 'other', 'dependency'].includes(job.tag)}
- Worker group: {job.tag} + Tag: {job.tag}
{/if} {#if !job.visible_to_owner}