improve bun dedicated worker to use nodejs

This commit is contained in:
Ruben Fiszel
2024-05-01 02:10:44 +02:00
parent b854165b0d
commit 017710f17d
4 changed files with 214 additions and 136 deletions
+208 -134
View File
@@ -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<sqlx::Postgres>,
timeout: Option<i32>,
mem_peak: &mut i32,
canceled_by: &mut Option<CanceledBy>,
common_bun_proc_envs: &HashMap<String, String>,
) -> 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<String>,
@@ -520,66 +632,16 @@ try {{
Ok(reserved_variables) as error::Result<HashMap<String, String>>
};
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<String, String> =
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
}
}
@@ -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<QueuedJob> = jobs.pop_front().expect("pop");
tracing::info!("job completed on dedicated worker {script_path}: {}", job.id);
@@ -528,6 +528,7 @@ for await (const chunk of Deno.stdin.readable) {{
worker_name,
db,
script_path,
"deno",
)
.await
}
@@ -1198,6 +1198,7 @@ for line in sys.stdin:
worker_name,
db,
script_path,
"python",
)
.await
}