diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index bd892a59ef..1390ed3c61 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -338,6 +338,118 @@ fn get_annotation(inner_content: &str) -> Annotations { Annotations { npm_mode, nodejs_mode } } +pub async fn build_loader( + job_dir: &str, + base_internal_url: &str, + token: &str, + w_id: &str, + current_path: &str, + nodejs_mode: bool, +) -> Result<()> { + let loader = RELATIVE_BUN_LOADER + .replace("W_ID", w_id) + .replace("BASE_INTERNAL_URL", base_internal_url) + .replace("TOKEN", token) + .replace("CURRENT_PATH", current_path) + .replace("RAW_GET_ENDPOINT", "raw_unpinned"); + if nodejs_mode { + write_file( + &job_dir, + "node_builder.ts", + &format!( + r#" +{} + +import {{ readdir }} from "node:fs/promises"; + +let fileNames = [] +try {{ + fileNames = await readdir("{job_dir}/node_modules") +}} catch (e) {{ + +}} + +const bo = await Bun.build({{ + entrypoints: ["{job_dir}/wrapper.ts"], + outdir: "./", + target: "node", + plugins: [p], + external: fileNames, + }}); + +if (!bo.success) {{ + bo.logs.forEach((l) => console.log(l)); + process.exit(1); +}} +"#, + loader + ), + ) + .await?; + } else { + write_file( + &job_dir, + "loader.bun.ts", + &format!( + r#" +import {{ plugin }} from "bun"; + +{} + +plugin(p) +"#, + loader + ), + ) + .await?; + }; + Ok(()) +} + +pub async fn generate_wrapper_mjs( + job_dir: &str, + w_id: &str, + job_id: &Uuid, + worker_name: &str, + db: &sqlx::Pool, + timeout: Option, + mem_peak: &mut i32, + canceled_by: &mut Option, + common_bun_proc_envs: &HashMap, +) -> Result<()> { + let mut child = Command::new(&*BUN_PATH); + child + .current_dir(job_dir) + .env_clear() + .envs(common_bun_proc_envs.clone()) + .env("PATH", PATH_ENV.as_str()) + .args(vec!["run", "node_builder.ts"]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); + let child_process = start_child_process(child, &*BUN_PATH).await?; + handle_child( + job_id, + db, + mem_peak, + canceled_by, + child_process, + false, + worker_name, + w_id, + "bun build", + timeout, + false, + ) + .await?; + tokio::fs::rename( + format!("{job_dir}/wrapper.js"), + format!("{job_dir}/wrapper.mjs"), + ) + .await + .map_err(|e| error::Error::InternalErr(format!("Could not move wrapper to mjs: {e}")))?; + Ok(()) +} + #[tracing::instrument(level = "trace", skip_all)] pub async fn handle_bun_job( requirements_o: Option, @@ -520,66 +632,16 @@ try {{ Ok(reserved_variables) as error::Result> }; - let loader = RELATIVE_BUN_LOADER - .replace("W_ID", &job.workspace_id) - .replace("BASE_INTERNAL_URL", base_internal_url) - .replace("TOKEN", &client.get_token().await) - .replace("CURRENT_PATH", job.script_path()) - .replace("RAW_GET_ENDPOINT", "raw_unpinned"); - let write_loader_f = async move { - if annotation.nodejs_mode { - write_file( - &job_dir, - "node_builder.ts", - &format!( - r#" -{} - -import {{ readdir }} from "node:fs/promises"; - -let fileNames = [] -try {{ - fileNames = await readdir("{job_dir}/node_modules") -}} catch (e) {{ - -}} - -const bo = await Bun.build({{ - entrypoints: ["{job_dir}/wrapper.ts"], - outdir: "./", - target: "node", - plugins: [p], - external: fileNames, - }}); - -if (!bo.success) {{ - bo.logs.forEach((l) => console.log(l)); - process.exit(1); -}} -"#, - loader - ), - ) - .await?; - Ok(()) as error::Result<()> - } else { - write_file( - &job_dir, - "loader.bun.ts", - &format!( - r#" -import {{ plugin }} from "bun"; - -{} - -plugin(p) -"#, - loader - ), - ) - .await?; - Ok(()) as error::Result<()> - } + let write_loader_f = async { + build_loader( + job_dir, + base_internal_url, + &client.get_token().await, + &job.workspace_id, + &job.script_path(), + annotation.nodejs_mode, + ) + .await }; let (reserved_variables, _, _) = tokio::try_join!( @@ -589,36 +651,18 @@ plugin(p) )?; if annotation.nodejs_mode { - let mut child = Command::new(&*BUN_PATH); - child - .current_dir(job_dir) - .env_clear() - .envs(common_bun_proc_envs.clone()) - .env("PATH", PATH_ENV.as_str()) - .args(vec!["run", "node_builder.ts"]) - .stdout(Stdio::piped()) - .stderr(Stdio::piped()); - let child_process = start_child_process(child, &*BUN_PATH).await?; - handle_child( + generate_wrapper_mjs( + job_dir, + &job.workspace_id, &job.id, + worker_name, db, + job.timeout, mem_peak, canceled_by, - child_process, - false, - worker_name, - &job.workspace_id, - "bun build", - job.timeout, - false, + &common_bun_proc_envs, ) .await?; - tokio::fs::rename( - format!("{job_dir}/wrapper.js"), - format!("{job_dir}/wrapper.mjs"), - ) - .await - .map_err(|e| error::Error::InternalErr(format!("Could not move wrapper to mjs: {e}")))?; } //do not cache local dependencies @@ -800,6 +844,11 @@ pub async fn start_worker( let common_bun_proc_envs: HashMap = get_common_bun_proc_envs(&base_internal_url).await; + let mut annotation = get_annotation(inner_content); + + //TODO: remove this when bun dedicated workers work without issues + annotation.nodejs_mode = true; + let context = variables::get_reserved_variables( db, w_id, @@ -905,7 +954,7 @@ pub async fn start_worker( let is_debug = std::env::var("RUST_LOG").is_ok_and(|x| x == "windmill=debug"); let print_lines = if is_debug { - r#"stdout.write(line+'\n');"# + r#"console.log(line);"# } else { "" }; @@ -921,73 +970,98 @@ BigInt.prototype.toJSON = function () {{ {dates} -let stdout = Bun.stdout.writer(); -stdout.write('start\n'); +console.log('start'); for await (const line of createInterface({{ input: process.stdin }})) {{ {print_lines} if (line === "end") {{ - break; + process.exit(0); }} try {{ let {{ {spread} }} = JSON.parse(line) let res: any = await main(...[ {spread} ]); - stdout.write("wm_res[success]:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value) + '\n'); + console.log("wm_res[success]:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value)); }} catch (e) {{ - stdout.write("wm_res[error]:" + JSON.stringify({{ message: e.message, name: e.name, stack: e.stack, line: line }}) + '\n'); + console.log("wm_res[error]:" + JSON.stringify({{ message: e.message, name: e.name, stack: e.stack, line: line }})); }} - stdout.flush(); - }} "#, ); write_file(job_dir, "wrapper.ts", &wrapper_content).await?; } - let _ = write_file( - &job_dir, - "loader.bun.ts", - &format!( - r#" -import {{ plugin }} from "bun"; - -{} - -plugin(p) -"#, - RELATIVE_BUN_LOADER - .replace("W_ID", &w_id) - .replace("BASE_INTERNAL_URL", base_internal_url) - .replace("TOKEN", token) - .replace("CURRENT_PATH", script_path) - .replace("RAW_GET_ENDPOINT", "raw_unpinned") - ), + build_loader( + job_dir, + base_internal_url, + token, + w_id, + script_path, + annotation.nodejs_mode, ) .await?; - handle_dedicated_process( - &*BUN_PATH, - job_dir, - context_envs, - envs, - context, - common_bun_proc_envs, - vec![ - "run", - "-i", - "--prefer-offline", - "-r", - "./loader.bun.ts", - &format!("{job_dir}/wrapper.ts"), - ], - killpill_rx, - job_completed_tx, - token, - jobs_rx, - worker_name, - db, - script_path, - ) - .await + if annotation.nodejs_mode { + generate_wrapper_mjs( + job_dir, + w_id, + &Uuid::nil(), + worker_name, + db, + None, + &mut mem_peak, + &mut canceled_by, + &common_bun_proc_envs, + ) + .await?; + } + + if annotation.nodejs_mode { + let script_path = format!("{job_dir}/wrapper.mjs"); + + handle_dedicated_process( + &*NODE_PATH, + job_dir, + context_envs, + envs, + context, + common_bun_proc_envs, + vec![&script_path], + killpill_rx, + job_completed_tx, + token, + jobs_rx, + worker_name, + db, + &script_path, + "nodejs", + ) + .await + } else { + handle_dedicated_process( + &*BUN_PATH, + job_dir, + context_envs, + envs, + context, + common_bun_proc_envs, + vec![ + "run", + "-i", + "--prefer-offline", + "-r", + "./loader.bun.ts", + &format!("{job_dir}/wrapper.ts"), + ], + killpill_rx, + job_completed_tx, + token, + jobs_rx, + worker_name, + db, + script_path, + "bun", + ) + .await + } } diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index e46a89a3ed..426eb2204b 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -59,6 +59,7 @@ pub async fn handle_dedicated_process( worker_name: &str, db: &DB, script_path: &str, + mode: &str, ) -> std::result::Result<(), error::Error> { //do not cache local dependencies let mut child = { @@ -115,7 +116,7 @@ pub async fn handle_dedicated_process( // let mut j = 0; let mut alive = true; - let init_log = format!("dedicated worker: {worker_name}\n\n"); + let init_log = format!("dedicated worker {mode}: {worker_name}\n\n"); let mut logs = init_log.clone(); loop { tokio::select! { @@ -126,6 +127,7 @@ pub async fn handle_dedicated_process( if let Err(e) = write_stdin(&mut stdin, "end").await { tracing::info!("Could not write end message to stdin: {e:?}") } + stdin.flush().await.context("stdin flush")?; }, line = err_reader.next_line() => { if let Some(line) = line.expect("line is ok") { @@ -146,7 +148,7 @@ pub async fn handle_dedicated_process( tracing::info!("dedicated worker process started"); continue; } - tracing::debug!("processed job: {line}"); + tracing::debug!("processed job: |{line}|"); if line.starts_with("wm_res[") { let job: Arc = jobs.pop_front().expect("pop"); tracing::info!("job completed on dedicated worker {script_path}: {}", job.id); diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index b6e8b958f2..5888313b92 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -528,6 +528,7 @@ for await (const chunk of Deno.stdin.readable) {{ worker_name, db, script_path, + "deno", ) .await } diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 3b9c724d29..4f33688491 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -1198,6 +1198,7 @@ for line in sys.stdin: worker_name, db, script_path, + "python", ) .await }