diff --git a/backend/src/main.rs b/backend/src/main.rs index 93ebd1d13b..83dd87b720 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -54,9 +54,9 @@ use windmill_common::METRICS_ADDR; use windmill_common::global_settings::OBJECT_STORE_CACHE_CONFIG_SETTING; use windmill_worker::{ - BUN_CACHE_DIR, BUN_TAR_CACHE_DIR, DENO_CACHE_DIR, DENO_CACHE_DIR_DEPS, DENO_CACHE_DIR_NPM, - GO_BIN_CACHE_DIR, GO_CACHE_DIR, HUB_CACHE_DIR, LOCK_CACHE_DIR, PIP_CACHE_DIR, - POWERSHELL_CACHE_DIR, TAR_PIP_CACHE_DIR, TMP_LOGS_DIR, + BUN_BUNDLE_CACHE_DIR, BUN_CACHE_DIR, BUN_DEPSTAR_CACHE_DIR, DENO_CACHE_DIR, + DENO_CACHE_DIR_DEPS, DENO_CACHE_DIR_NPM, GO_BIN_CACHE_DIR, GO_CACHE_DIR, HUB_CACHE_DIR, + LOCK_CACHE_DIR, PIP_CACHE_DIR, POWERSHELL_CACHE_DIR, TAR_PIP_CACHE_DIR, TMP_LOGS_DIR, }; use crate::monitor::{ @@ -778,7 +778,8 @@ pub async fn run_workers Annotations { Annotations { npm_mode, nodejs_mode } } +pub async fn load_cache(bin_path: &str, remote_path: &str) -> (bool, String) { + if tokio::fs::metadata(&bin_path).await.is_ok() { + (true, format!("loaded from local cache: {}\n", bin_path)) + } else { + #[cfg(all(feature = "enterprise", feature = "parquet"))] + if let Some(os) = crate::s3_helpers::OBJECT_STORE_CACHE_SETTINGS + .read() + .await + .clone() + { + use crate::s3_helpers::attempt_fetch_bytes; + + if let Ok(mut x) = attempt_fetch_bytes(os, remote_path).await { + if let Err(e) = write_binary_file(bin_path, &mut x) { + tracing::error!("could not write binary file: {e:?}"); + return ( + false, + "error writing binary file from object store".to_string(), + ); + } + tracing::info!("loaded from object store {}", bin_path); + return (true, format!("loaded bin from object store {}", bin_path)); + } + } + (false, "".to_string()) + } +} + +pub async fn exists_in_cache(bin_path: &str, remote_path: &str) -> bool { + if tokio::fs::metadata(&bin_path).await.is_ok() { + return true; + } else { + #[cfg(all(feature = "enterprise", feature = "parquet"))] + if let Some(os) = crate::s3_helpers::OBJECT_STORE_CACHE_SETTINGS + .read() + .await + .clone() + { + return os + .get(&object_store::path::Path::from(remote_path)) + .await + .is_ok(); + } + return false; + } +} + +pub async fn save_cache( + local_cache_path: &str, + remote_cache_path: &str, + origin: &str, +) -> crate::error::Result { + let mut _cached_to_s3 = false; + #[cfg(all(feature = "enterprise", feature = "parquet"))] + if let Some(os) = crate::s3_helpers::OBJECT_STORE_CACHE_SETTINGS + .read() + .await + .clone() + { + use object_store::path::Path; + + if let Err(e) = os + .put( + &Path::from(remote_cache_path), + std::fs::read(origin)?.into(), + ) + .await + { + tracing::error!( + "Failed to put go bin to object store: {remote_cache_path}. Error: {:?}", + e + ); + } else { + _cached_to_s3 = true; + } + } + + // if !*CLOUD_HOSTED { + if true { + std::fs::copy(origin, local_cache_path)?; + Ok(format!( + "\nwrote cached binary: {} (backed by EE distributed object store: {_cached_to_s3})\n", + local_cache_path + )) + } else if _cached_to_s3 { + Ok(format!( + "wrote cached binary to object store {}\n", + local_cache_path + )) + } else { + Ok("".to_string()) + } +} + +#[cfg(all(feature = "enterprise", feature = "parquet"))] +fn write_binary_file(main_path: &str, byts: &mut bytes::Bytes) -> error::Result<()> { + use std::fs::{File, Permissions}; + use std::io::Write; + use std::os::unix::fs::PermissionsExt; + + let mut file = File::create(main_path)?; + file.write_all(byts)?; + file.set_permissions(Permissions::from_mode(0o755))?; + file.flush()?; + Ok(()) +} + fn get_cgroupv2_path() -> Option { let cgroup_path: String = parse_file("/proc/self/cgroup")?; diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 45cac13b64..5868bf46d7 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -14,11 +14,12 @@ use crate::common::build_envs_map; use crate::{ common::{ create_args_and_out_file, get_main_override, get_reserved_variables, handle_child, - parse_npm_config, read_result, start_child_process, write_file, write_file_binary, + parse_npm_config, read_file_content, read_result, start_child_process, write_file, + write_file_binary, }, - AuthedClientBackgroundTask, BUNFIG_INSTALL_SCOPES, BUN_CACHE_DIR, BUN_PATH, BUN_TAR_CACHE_DIR, - DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NODE_PATH, NPM_CONFIG_REGISTRY, NPM_PATH, NSJAIL_PATH, - PATH_ENV, TZ_ENV, + AuthedClientBackgroundTask, BUNFIG_INSTALL_SCOPES, BUN_BUNDLE_CACHE_DIR, BUN_CACHE_DIR, + BUN_DEPSTAR_CACHE_DIR, BUN_PATH, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NODE_PATH, + NPM_CONFIG_REGISTRY, NPM_PATH, NSJAIL_PATH, PATH_ENV, TZ_ENV, }; use tokio::{fs::File, process::Command}; @@ -34,6 +35,8 @@ use windmill_common::variables; use windmill_common::{ error::{self, Result}, jobs::QueuedJob, + worker::{exists_in_cache, get_annotation, save_cache}, + DB, }; #[cfg(all(feature = "enterprise", feature = "parquet"))] @@ -281,13 +284,20 @@ pub async fn install_lockfile( Ok(()) } -pub async fn build_loader( +#[derive(PartialEq)] +enum LoaderMode { + Node, + Bun, + BunBundle, + NodeBundle, +} +async fn build_loader( job_dir: &str, base_internal_url: &str, token: &str, w_id: &str, current_path: &str, - nodejs_mode: bool, + mode: LoaderMode, ) -> Result<()> { let loader = RELATIVE_BUN_LOADER .replace("W_ID", w_id) @@ -298,7 +308,7 @@ pub async fn build_loader( &crate::common::use_flow_root_path(current_path), ) .replace("RAW_GET_ENDPOINT", "raw_unpinned"); - if nodejs_mode { + if mode == LoaderMode::Node { write_file( &job_dir, "node_builder.ts", @@ -321,6 +331,7 @@ const bo = await Bun.build({{ target: "node", plugins: [p], external: fileNames, + minify: true, }}); if (!bo.success) {{ @@ -332,7 +343,7 @@ if (!bo.success) {{ ), ) .await?; - } else { + } else if mode == LoaderMode::Bun { write_file( &job_dir, "loader.bun.js", @@ -348,7 +359,42 @@ plugin(p) ), ) .await?; - }; + } else if mode == LoaderMode::BunBundle || mode == LoaderMode::NodeBundle { + write_file( + &job_dir, + "node_builder.ts", + &format!( + r#" +{} + +const bo = await Bun.build({{ + entrypoints: ["{job_dir}/main.ts"], + outdir: "./", + target: "{}", + plugins: [p], + external: [], + minify: {{ + identifiers: false, + syntax: true, + whitespace: false + }}, + }}); + +if (!bo.success) {{ + bo.logs.forEach((l) => console.log(l)); + process.exit(1); +}} +"#, + loader, + if mode == LoaderMode::BunBundle { + "bun" + } else { + "node" + } + ), + ) + .await?; + } Ok(()) } @@ -387,15 +433,52 @@ pub async fn generate_wrapper_mjs( false, ) .await?; - tokio::fs::rename( + 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(()) } +pub async fn generate_bun_bundle( + 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?; + Ok(()) +} + #[cfg(all(feature = "enterprise", feature = "parquet"))] pub async fn pull_codebase(w_id: &str, id: &str, job_dir: &str) -> Result<()> { let path = windmill_common::s3_helpers::bundle(&w_id, &id); @@ -470,6 +553,77 @@ pub fn copy_recursively( Ok(()) } +pub async fn prebundle_script( + inner_content: &str, + lockfile: Option, + script_path: &str, + job_id: &Uuid, + w_id: &str, + db: &DB, + job_dir: &str, + base_internal_url: &str, + worker_name: &str, + token: &str, +) -> Result<()> { + let (local_path, remote_path) = compute_bundle_local_and_remote_path(inner_content, &lockfile); + if exists_in_cache(&local_path, &remote_path).await { + return Ok(()); + } + let annotation = get_annotation(inner_content); + let origin = format!("{job_dir}/main.js"); + write_file(job_dir, "main.ts", &remove_pinned_imports(inner_content)?).await?; + build_loader( + job_dir, + base_internal_url, + &token, + w_id, + script_path, + if annotation.nodejs_mode { + LoaderMode::NodeBundle + } else { + LoaderMode::BunBundle + }, + ) + .await?; + + let common_bun_proc_envs: HashMap = + get_common_bun_proc_envs(&base_internal_url).await; + + generate_bun_bundle( + job_dir, + w_id, + job_id, + worker_name, + db, + None, + &mut 0, + &mut None, + &common_bun_proc_envs, + ) + .await?; + save_cache(&local_path, &remote_path, &origin).await?; + Ok(()) +} + +pub const BUN_BUNDLE_OBJECT_STORE_PREFIX: &str = "bun_bundle/"; + +fn compute_bundle_local_and_remote_path( + inner_content: &str, + requirements_o: &Option, +) -> (String, String) { + let hash = windmill_common::utils::calculate_hash(&format!( + "{}{}", + inner_content, + requirements_o + .as_ref() + .map(|x| x.to_string()) + .unwrap_or_default() + )); + let local_path = format!("{BUN_BUNDLE_CACHE_DIR}/{hash}"); + let remote_path = format!("{BUN_BUNDLE_OBJECT_STORE_PREFIX}{hash}"); + (local_path, remote_path) +} + #[tracing::instrument(level = "trace", skip_all)] pub async fn handle_bun_job( requirements_o: Option, @@ -486,7 +640,18 @@ pub async fn handle_bun_job( envs: HashMap, shared_mount: &str, ) -> error::Result> { - if !codebase.is_some() { + let (mut bundle_cache, cache_logs, local_path, remote_path) = if requirements_o.is_some() + && codebase.is_none() + { + let (local_path, remote_path) = + compute_bundle_local_and_remote_path(inner_content, &requirements_o); + let (cache, logs) = windmill_common::worker::load_cache(&local_path, &remote_path).await; + (cache, logs, local_path, remote_path) + } else { + (false, "".to_string(), "".to_string(), "".to_string()) + }; + + if !codebase.is_some() && !bundle_cache { let _ = write_file(job_dir, "main.ts", inner_content).await?; } else { let _ = write_file(job_dir, "package.json", r#"{ "type": "module" }"#).await?; @@ -510,15 +675,23 @@ pub async fn handle_bun_job( } let mut gbuntar_name = None; - if let Some(codebase) = codebase.as_ref() { + if bundle_cache { + let target = format!("{job_dir}/main.js"); + std::os::unix::fs::symlink(&local_path, &target).map_err(|e| { + error::Error::ExecutionErr(format!( + "could not copy cached binary from {local_path} to {job_dir}/main: {e:?}" + )) + })?; + } else if let Some(codebase) = codebase.as_ref() { pull_codebase(&job.workspace_id, codebase, job_dir).await?; - } else if let Some(reqs) = requirements_o { + } else if let Some(reqs) = requirements_o.as_ref() { let splitted = reqs.split(BUN_LOCKB_SPLIT).collect::>(); if splitted.len() != 2 && !annotation.npm_mode { return Err(error::Error::ExecutionErr( format!("Invalid requirements, expected to find //bun.lockb split pattern in reqs. Found: |{reqs}|") )); } + let _ = write_file(job_dir, "package.json", &splitted[0]).await?; let lockb = if annotation.npm_mode { "" } else { splitted[1] }; if lockb != EMPTY_FILE { @@ -543,7 +716,7 @@ pub async fn handle_bun_job( let buntar_name = base64::engine::general_purpose::URL_SAFE.encode(sha_path.finalize()); - buntar_path = format!("{BUN_TAR_CACHE_DIR}/{buntar_name}"); + buntar_path = format!("{BUN_DEPSTAR_CACHE_DIR}/{buntar_name}"); #[cfg(unix)] if tokio::fs::metadata(&buntar_path).await.is_ok() { @@ -617,13 +790,19 @@ pub async fn handle_bun_job( // } } - let _ = write_file(job_dir, "main.ts", &remove_pinned_imports(inner_content)?).await?; - - let mut init_logs = if codebase.is_some() { - "\n\n--- NODE SNAPSHOT EXECUTION ---\n".to_string() + let mut init_logs = if bundle_cache { + if annotation.nodejs_mode { + "\n\n--- NODE BUNDLE SNAPSHOT EXECUTION ---\n".to_string() + } else { + "\n\n--- BUN BUNDLE SNAPSHOT EXECUTION ---\n".to_string() + } + } else if codebase.is_some() { + "\n\n--- NODE CODEBASE SNAPSHOT EXECUTION ---\n".to_string() } else if annotation.nodejs_mode { + write_file(job_dir, "main.ts", &remove_pinned_imports(inner_content)?).await?; "\n\n--- NODE CODE EXECUTION ---\n".to_string() } else { + write_file(job_dir, "main.ts", &remove_pinned_imports(inner_content)?).await?; "\n\n--- BUN CODE EXECUTION ---\n".to_string() }; @@ -634,7 +813,9 @@ pub async fn handle_bun_job( ); } - append_logs(&job.id, &job.workspace_id, init_logs, db).await; + if bundle_cache { + init_logs = format!("\n{}{}", cache_logs, init_logs); + } let write_wrapper_f = async { // let mut start = Instant::now(); @@ -659,7 +840,7 @@ pub async fn handle_bun_job( // we cannot use Bun.read and Bun.write because it results in an EBADF error on cloud let main_name = main_override.unwrap_or("main".to_string()); - let main_import = if codebase.is_some() { + let main_import = if codebase.is_some() || bundle_cache { "./main.js" } else { "./main.ts" @@ -716,15 +897,37 @@ try {{ Ok(reserved_variables) as error::Result> }; + let build_cache = !bundle_cache && !codebase.is_some() && requirements_o.is_some(); + let write_loader_f = async { - if !codebase.is_some() { + if build_cache { build_loader( job_dir, base_internal_url, &client.get_token().await, &job.workspace_id, &job.script_path(), - annotation.nodejs_mode, + if annotation.nodejs_mode { + LoaderMode::NodeBundle + } else { + LoaderMode::BunBundle + }, + ) + .await?; + + Ok(()) + } else if !codebase.is_some() && !bundle_cache { + build_loader( + job_dir, + base_internal_url, + &client.get_token().await, + &job.workspace_id, + &job.script_path(), + if annotation.nodejs_mode { + LoaderMode::Node + } else { + LoaderMode::Bun + }, ) .await } else { @@ -737,21 +940,61 @@ try {{ write_wrapper_f, write_loader_f )?; - - if annotation.nodejs_mode && !codebase.is_some() { - generate_wrapper_mjs( - job_dir, - &job.workspace_id, - &job.id, - worker_name, - db, - job.timeout, - mem_peak, - canceled_by, - &common_bun_proc_envs, - ) - .await?; + if !codebase.is_some() && !bundle_cache { + if build_cache { + generate_bun_bundle( + job_dir, + &job.workspace_id, + &job.id, + worker_name, + db, + job.timeout, + mem_peak, + canceled_by, + &common_bun_proc_envs, + ) + .await?; + match save_cache(&local_path, &remote_path, &format!("{job_dir}/main.js")).await { + Err(e) => { + let em = format!("could not save {local_path} to go cache: {e:?}"); + tracing::error!(em) + } + Ok(logs) => { + init_logs.push_str(&"\n"); + init_logs.push_str(&logs); + init_logs.push_str(&"\n"); + tracing::info!("saved bun bundle cache: {logs}") + } + } + let ex_wrapper = read_file_content(&format!("{job_dir}/wrapper.mjs")).await?; + write_file( + job_dir, + "wrapper.mjs", + &ex_wrapper.replace( + "import * as Main from \"./main.ts\"", + "import * as Main from \"./main.js\"", + ), + ) + .await?; + write_file(job_dir, "package.json", r#"{ "type": "module" }"#).await?; + fs::remove_file(format!("{job_dir}/main.ts"))?; + bundle_cache = true; + } else if annotation.nodejs_mode { + generate_wrapper_mjs( + job_dir, + &job.workspace_id, + &job.id, + worker_name, + db, + job.timeout, + mem_peak, + canceled_by, + &common_bun_proc_envs, + ) + .await?; + } } + append_logs(&job.id, &job.workspace_id, init_logs, db).await; //do not cache local dependencies let child = if !*DISABLE_NSJAIL { @@ -793,13 +1036,14 @@ try {{ &NODE_PATH, "/tmp/nodejs/wrapper.mjs", ] - } else if codebase.is_some() { + } else if codebase.is_some() || bundle_cache { vec![ "--config", "run.config.proto", "--", &BUN_PATH, "run", + "--preserve-symlinks", "/tmp/bun/wrapper.mjs", ] } else { @@ -838,7 +1082,7 @@ try {{ .envs(envs) .envs(reserved_variables) .envs(common_bun_proc_envs) - .args(vec![&script_path]) + .args(vec!["--preserve-symlinks", &script_path]) .stdout(Stdio::piped()) .stderr(Stdio::piped()); bun_cmd @@ -846,7 +1090,7 @@ try {{ let script_path = format!("{job_dir}/wrapper.mjs"); let mut bun_cmd = Command::new(&*BUN_PATH); - let args = if codebase.is_some() { + let args = if codebase.is_some() || bundle_cache { vec!["run", &script_path] } else { vec![ @@ -1105,7 +1349,11 @@ for await (const line of Readline.createInterface({{ input: process.stdin }})) { token, w_id, script_path, - annotation.nodejs_mode, + if annotation.nodejs_mode { + LoaderMode::Node + } else { + LoaderMode::Bun + }, ) .await?; } diff --git a/backend/windmill-worker/src/go_executor.rs b/backend/windmill-worker/src/go_executor.rs index 398e55cb5f..3a71170cb7 100644 --- a/backend/windmill-worker/src/go_executor.rs +++ b/backend/windmill-worker/src/go_executor.rs @@ -12,7 +12,7 @@ use windmill_common::{ error::{self, Error}, jobs::QueuedJob, utils::calculate_hash, - worker::CLOUD_HOSTED, + worker::{save_cache, CLOUD_HOSTED}, }; use windmill_parser_go::{parse_go_imports, REQUIRE_PARSE}; use windmill_queue::{append_logs, CanceledBy}; @@ -33,111 +33,7 @@ lazy_static::lazy_static! { static ref GO_PATH: String = std::env::var("GO_PATH").unwrap_or_else(|_| "/usr/bin/go".to_string()); } -pub async fn save_cache( - bin_path: &str, - job_dir: &str, - _hash: &str, - job: &QueuedJob, - db: &sqlx::Pool, -) -> windmill_common::error::Result<()> { - let job_main_path = format!("{job_dir}/main"); - let mut _cached_to_s3 = false; - #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(os) = windmill_common::s3_helpers::OBJECT_STORE_CACHE_SETTINGS - .read() - .await - .clone() - { - use object_store::path::Path; - - let hash_path = hash_to_os_path(_hash); - if let Err(e) = os - .put( - &Path::from(hash_path.clone()), - std::fs::read(&job_main_path)?.into(), - ) - .await - { - tracing::error!( - "Failed to put go bin to object store: {hash_path}. Error: {:?}", - e - ); - } else { - _cached_to_s3 = true; - } - } - - if !*CLOUD_HOSTED { - tokio::fs::copy(&job_main_path, bin_path).await?; - append_logs( - &job.id, - &job.workspace_id, - format!( - "\nwrite cached binary: {} (backed by object store: {_cached_to_s3})\n", - bin_path - ), - db, - ) - .await; - } else if _cached_to_s3 { - append_logs( - &job.id, - &job.workspace_id, - format!("write cached binary to object store {}\n", bin_path), - db, - ) - .await; - } - - Ok(()) -} - -#[cfg(all(feature = "enterprise", feature = "parquet"))] -async fn write_binary_file(main_path: &str, byts: &mut bytes::Bytes) -> error::Result<()> { - use std::fs::Permissions; - use std::os::unix::fs::PermissionsExt; - use tokio::io::AsyncWriteExt; - - let mut file = File::create(main_path).await?; - file.write_all_buf(byts).await?; - file.set_permissions(Permissions::from_mode(0o755)).await?; - file.flush().await?; - Ok(()) -} - -#[cfg(all(feature = "enterprise", feature = "parquet"))] -fn hash_to_os_path(hash: &str) -> String { - format!("gobin/{hash}") -} - -async fn load_cache(bin_path: &str, _hash: &str) -> (bool, String) { - if tokio::fs::metadata(&bin_path).await.is_ok() { - (true, format!("loaded bin from local cache: {}\n", bin_path)) - } else { - #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(os) = windmill_common::s3_helpers::OBJECT_STORE_CACHE_SETTINGS - .read() - .await - .clone() - { - use windmill_common::s3_helpers::attempt_fetch_bytes; - - if let Ok(mut x) = attempt_fetch_bytes(os, &hash_to_os_path(_hash)).await { - if let Err(e) = write_binary_file(bin_path, &mut x).await { - tracing::error!("could not write binary file: {e:?}"); - return ( - false, - "error writing binary file from object store".to_string(), - ); - } - tracing::info!("loaded bin from object store {}", bin_path); - return (true, format!("loaded bin from object store {}", bin_path)); - } - } - (false, "".to_string()) - } -} - +pub const GO_OBJECT_STORE_PREFIX: &str = "gobin/"; #[tracing::instrument(level = "trace", skip_all)] pub async fn handle_go_job( mem_peak: &mut i32, @@ -163,9 +59,9 @@ pub async fn handle_go_job( .map(|x| x.to_string()) .unwrap_or_default() )); - let bin_path = format!("{}/{hash}", GO_BIN_CACHE_DIR,); - - let (cache, cache_logs) = load_cache(&bin_path, &hash).await; + let bin_path = format!("{}/{hash}", GO_BIN_CACHE_DIR); + let remote_path = format!("{GO_OBJECT_STORE_PREFIX}{hash}"); + let (cache, cache_logs) = windmill_common::worker::load_cache(&bin_path, &remote_path).await; let (skip_go_mod, skip_tidy) = if cache { create_dir(job_dir).await?; @@ -309,13 +205,23 @@ func Run(req Req) (interface{{}}, error){{ ) .await?; - if let Err(e) = save_cache(&bin_path, &job_dir, &hash, &job, db).await { - tracing::error!("could not save {bin_path} to go cache: {e:?}"); + match save_cache( + &bin_path, + &format!("{GO_OBJECT_STORE_PREFIX}{hash}"), + &format!("{job_dir}/main"), + ) + .await + { + Err(e) => { + let em = format!("could not save {bin_path} to go cache: {e:?}"); + tracing::error!(em); + em + } + Ok(logs) => logs, } - "".to_string() } else { let target = format!("{job_dir}/main"); - tokio::fs::symlink(&bin_path, &target).await.map_err(|e| { + std::os::unix::fs::symlink(&bin_path, &target).map_err(|e| { Error::ExecutionErr(format!( "could not copy cached binary from {bin_path} to {job_dir}/main: {e:?}" )) diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 0ce67c01e7..ef6087ba27 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -237,7 +237,8 @@ pub const DENO_CACHE_DIR_NPM: &str = concatcp!(ROOT_CACHE_DIR, "deno/npm"); pub const GO_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "go"); pub const BUN_CACHE_DIR: &str = concatcp!(ROOT_CACHE_NOMOUNT_DIR, "bun"); -pub const BUN_TAR_CACHE_DIR: &str = concatcp!(ROOT_CACHE_NOMOUNT_DIR, "buntar"); +pub const BUN_BUNDLE_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "bun"); +pub const BUN_DEPSTAR_CACHE_DIR: &str = concatcp!(ROOT_CACHE_NOMOUNT_DIR, "buntar"); pub const HUB_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "hub"); pub const GO_BIN_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "gobin"); diff --git a/backend/windmill-worker/src/worker_lockfiles.rs b/backend/windmill-worker/src/worker_lockfiles.rs index a919b49cb5..051d8430a8 100644 --- a/backend/windmill-worker/src/worker_lockfiles.rs +++ b/backend/windmill-worker/src/worker_lockfiles.rs @@ -1319,6 +1319,21 @@ async fn capture_dependency_job( npm_mode, ) .await?; + if req.is_some() { + crate::bun_executor::prebundle_script( + job_raw_code, + req.clone(), + script_path, + job_id, + w_id, + db, + &job_dir, + base_internal_url, + worker_name, + &token, + ) + .await?; + } Ok(req.unwrap_or_else(String::new)) } ScriptLang::Php => {