From 3d7b3d74bcfd69ac6d82a23402bca9e2e7c171ea Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sat, 11 Nov 2023 15:19:27 +0100 Subject: [PATCH] feat: add dedicated worker support for deno --- backend/windmill-api/src/jobs.rs | 8 +- backend/windmill-worker/src/bun_executor.rs | 2 - backend/windmill-worker/src/deno_executor.rs | 222 ++++++++++++++---- backend/windmill-worker/src/worker.rs | 16 ++ .../src/lib/components/ScriptBuilder.svelte | 3 +- .../(root)/(logged)/run/[...run]/+page.svelte | 57 +++-- 6 files changed, 235 insertions(+), 73 deletions(-) diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 309ed26361..5c09ae0edf 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -364,7 +364,7 @@ async fn get_job( async fn get_job_internal(db: &DB, workspace_id: &str, job_id: Uuid) -> error::Result { let cjob_maybe = sqlx::query_as::<_, CompletedJob>("SELECT id, workspace_id, parent_job, created_by, created_at, duration_ms, success, script_hash, script_path, - CASE WHEN pg_column_size(args) < 2000000 THEN args ELSE '{\"reason\": \"WINDMILL_TOO_BIG\"}'::jsonb END as args, CASE WHEN pg_column_size(result) < 2000000 THEN result ELSE '\"WINDMILL_TOO_BIG\"'::jsonb END as result, logs, deleted, raw_code, canceled, canceled_by, canceled_reason, job_kind, env_id, + CASE WHEN args is null or pg_column_size(args) < 2000000 THEN args ELSE '{\"reason\": \"WINDMILL_TOO_BIG\"}'::jsonb END as args, CASE WHEN result is null or pg_column_size(result) < 2000000 THEN result ELSE '\"WINDMILL_TOO_BIG\"'::jsonb END as result, logs, deleted, raw_code, canceled, canceled_by, canceled_reason, job_kind, env_id, schedule_path, permissioned_as, flow_status, raw_flow, is_flow_step, language, started_at, is_skipped, raw_lock, email, visible_to_owner, mem_peak, tag, priority FROM completed_job WHERE id = $1 AND workspace_id = $2") @@ -378,7 +378,7 @@ async fn get_job_internal(db: &DB, workspace_id: &str, job_id: Uuid) -> error::R } else { let job_o = sqlx::query_as::<_, QueuedJob>( "SELECT id, workspace_id, parent_job, created_by, created_at, started_at, scheduled_for, running, - script_hash, script_path, CASE WHEN pg_column_size(args) < 2000000 THEN args ELSE '{\"reason\": \"WINDMILL_TOO_BIG\"}'::jsonb END as args, logs, raw_code, canceled, canceled_by, canceled_reason, last_ping, + script_hash, script_path, CASE WHEN args is null or pg_column_size(args) < 2000000 THEN args ELSE '{\"reason\": \"WINDMILL_TOO_BIG\"}'::jsonb END as args, logs, raw_code, canceled, canceled_by, canceled_reason, last_ping, job_kind, env_id, schedule_path, permissioned_as, flow_status, raw_flow, is_flow_step, language, suspend, suspend_until, same_worker, raw_lock, pre_run_error, email, visible_to_owner, mem_peak, root_job, leaf_jobs, tag, concurrent_limit, concurrency_time_window_s, timeout, flow_step_id, cache_ttl, priority @@ -2825,7 +2825,7 @@ async fn get_completed_job<'a>( Path((w_id, id)): Path<(String, Uuid)>, ) -> error::Result { let job_o = sqlx::query("SELECT id, workspace_id, parent_job, created_by, created_at, duration_ms, success, script_hash, script_path, - CASE WHEN pg_column_size(args) < 2000000 THEN args ELSE '\"WINDMILL_TOO_BIG\"'::jsonb END as args, CASE WHEN pg_column_size(result) < 2000000 THEN result ELSE '\"WINDMILL_TOO_BIG\"'::jsonb END as result, logs, deleted, raw_code, canceled, canceled_by, canceled_reason, job_kind, env_id, + CASE WHEN args is null or pg_column_size(args) < 2000000 THEN args ELSE '\"WINDMILL_TOO_BIG\"'::jsonb END as args, CASE WHEN result is null or pg_column_size(result) < 2000000 THEN result ELSE '\"WINDMILL_TOO_BIG\"'::jsonb END as result, logs, deleted, raw_code, canceled, canceled_by, canceled_reason, job_kind, env_id, schedule_path, permissioned_as, flow_status, raw_flow, is_flow_step, language, started_at, is_skipped, raw_lock, email, visible_to_owner, mem_peak, tag, priority FROM completed_job WHERE id = $1 AND workspace_id = $2") .bind(id) @@ -2944,7 +2944,7 @@ async fn delete_completed_job<'a>( require_admin(authed.is_admin, &authed.username)?; let job_o = sqlx::query( - "UPDATE completed_job SET logs = '', result = null, deleted = true WHERE id = $1 AND workspace_id = $2 \ + "UPDATE completed_job SET args = null, logs = '', result = null, deleted = true WHERE id = $1 AND workspace_id = $2 \ RETURNING *", ) .bind(id) diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 015fcd0476..f1c3864f6c 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -671,14 +671,12 @@ BigInt.prototype.toJSON = function () {{ {dates} let stdout = Bun.stdout.writer(); -// let stdout = Bun.file("output.txt").writer(); stdout.write('start\n'); for await (const chunk of Bun.stdin.stream()) {{ const lines = Buffer.from(chunk).toString(); let exit = false; for (const line of lines.trim().split("\n")) {{ - // stdout.write('s: ' + line + 'EE\n'); if (line === "end") {{ exit = true; break; diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index 1f645c64ab..0389cdb841 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -58,6 +58,8 @@ async fn get_common_deno_proc_envs( (String::from("PATH"), PATH_ENV.clone()), (String::from("HOME"), HOME_ENV.clone()), (String::from("TZ"), TZ_ENV.clone()), + (String::from("RUST_LOG"), "info".to_string()), + (String::from("DENO_DIR"), DENO_CACHE_DIR.to_string()), (String::from("DENO_AUTH_TOKENS"), deno_auth_tokens), ( String::from("BASE_INTERNAL_URL"), @@ -216,43 +218,12 @@ run().catch(async (e) => {{ Ok(()) as error::Result<()> }; - let write_import_map_f = async { - let w_id = job.workspace_id.clone(); - let script_path_split = job.script_path().split("/"); - let script_path_parts_len = script_path_split.clone().count(); - let mut relative_mounts = "".to_string(); - for c in 0..script_path_parts_len { - relative_mounts += ",\n "; - relative_mounts += &format!( - "\"./{}\": \"{base_internal_url}/api/w/{w_id}/scripts/raw/p/{}{}\"", - (0..c).map(|_| "../").join(""), - &script_path_split - .clone() - .take(script_path_parts_len - c - 1) - .join("/"), - if c == script_path_parts_len - 1 { - "" - } else { - "/" - }, - ); - } - let extra_import_map = DENO_EXTRA_IMPORT_MAP.as_str(); - let import_map = format!( - r#"{{ - "imports": {{ - "{base_internal_url}/api/w/{w_id}/scripts/raw/p/": "{base_internal_url}/api/w/{w_id}/scripts/raw/p/", - "{base_internal_url}": "{base_internal_url}/api/w/{w_id}/scripts/raw/p/", - "/": "{base_internal_url}/api/w/{w_id}/scripts/raw/p/", - "./wrapper.ts": "./wrapper.ts", - "./main.ts": "./main.ts"{relative_mounts} - {extra_import_map} - }} - }}"#, - ); - write_file(job_dir, "import_map.json", &import_map).await?; - Ok(()) as error::Result<()> - }; + let write_import_map_f = build_import_map( + &job.workspace_id, + job.script_path(), + base_internal_url, + job_dir, + ); let reserved_variables_args_out_f = async { let args_and_out_f = async { @@ -261,8 +232,7 @@ run().catch(async (e) => {{ }; let reserved_variables_f = async { let client = client.get_authed().await; - let mut vars = get_reserved_variables(job, &client.token, db).await?; - vars.insert("RUST_LOG".to_string(), "info".to_string()); + let vars = get_reserved_variables(job, &client.token, db).await?; Ok((vars, client.token)) as Result<(HashMap, String)> }; let (_, reserved_variables) = tokio::try_join!(args_and_out_f, reserved_variables_f)?; @@ -283,8 +253,8 @@ run().catch(async (e) => {{ } //do not cache local dependencies - let reload = format!("--reload={base_internal_url}"); let child = { + let reload = format!("--reload={base_internal_url}"); let script_path = format!("{job_dir}/wrapper.ts"); let import_map_path = format!("{job_dir}/import_map.json"); let mut args = Vec::with_capacity(12); @@ -321,7 +291,6 @@ run().catch(async (e) => {{ .envs(envs) .envs(reserved_variables) .envs(common_deno_proc_envs) - .env("DENO_DIR", DENO_CACHE_DIR) .args(args) .stdout(Stdio::piped()) .stderr(Stdio::piped()); @@ -351,3 +320,174 @@ run().catch(async (e) => {{ } read_result(job_dir).await } + +async fn build_import_map( + w_id: &str, + script_path: &str, + base_internal_url: &str, + job_dir: &str, +) -> error::Result<()> { + let script_path_split = script_path.split("/"); + let script_path_parts_len = script_path_split.clone().count(); + let mut relative_mounts = "".to_string(); + for c in 0..script_path_parts_len { + relative_mounts += ",\n "; + relative_mounts += &format!( + "\"./{}\": \"{base_internal_url}/api/w/{w_id}/scripts/raw/p/{}{}\"", + (0..c).map(|_| "../").join(""), + &script_path_split + .clone() + .take(script_path_parts_len - c - 1) + .join("/"), + if c == script_path_parts_len - 1 { + "" + } else { + "/" + }, + ); + } + let extra_import_map = DENO_EXTRA_IMPORT_MAP.as_str(); + let import_map = format!( + r#"{{ + "imports": {{ + "{base_internal_url}/api/w/{w_id}/scripts/raw/p/": "{base_internal_url}/api/w/{w_id}/scripts/raw/p/", + "{base_internal_url}": "{base_internal_url}/api/w/{w_id}/scripts/raw/p/", + "/": "{base_internal_url}/api/w/{w_id}/scripts/raw/p/", + "./wrapper.ts": "./wrapper.ts", + "./main.ts": "./main.ts"{relative_mounts} + {extra_import_map} + }} + }}"#, + ); + write_file(job_dir, "import_map.json", &import_map).await?; + Ok(()) as error::Result<()> +} + +#[cfg(feature = "enterprise")] +use crate::{dedicated_worker::handle_dedicated_process, JobCompletedSender}; +#[cfg(feature = "enterprise")] +use std::sync::Arc; +#[cfg(feature = "enterprise")] +use tokio::sync::mpsc::Receiver; + +#[cfg(feature = "enterprise")] +pub async fn start_worker( + 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: JobCompletedSender, + jobs_rx: Receiver>, + killpill_rx: tokio::sync::broadcast::Receiver<()>, +) -> Result<()> { + use windmill_common::variables; + + use crate::common::build_envs_map; + + let _ = write_file(job_dir, "main.ts", inner_content).await?; + let common_deno_proc_envs = get_common_deno_proc_envs(&token, base_internal_url).await; + + 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; + let context_envs = build_envs_map(context.to_vec()); + + { + // let mut start = Instant::now(); + let args = windmill_parser_ts::parse_deno_signature(inner_content, true)?.args; + let dates = args + .iter() + .enumerate() + .filter_map(|(i, x)| { + if matches!(x.typ, Typ::Datetime) { + Some(i) + } else { + None + } + }) + .map(|x| return format!("args[{x}] = args[{x}] ? new Date(args[{x}]) : undefined")) + .join("\n"); + + let spread = args.into_iter().map(|x| x.name).join(","); + // 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 {{ main }} from "./main.ts"; + +BigInt.prototype.toJSON = function () {{ + return this.toString(); +}}; + +{dates} + +console.log('start\n'); + +const decoder = new TextDecoder(); +for await (const chunk of Deno.stdin.readable) {{ + const lines = decoder.decode(chunk); + let exit = false; + for (const line of lines.trim().split("\n")) {{ + if (line === "end") {{ + exit = true; + break; + }} + try {{ + let {{ {spread} }} = JSON.parse(line) + let res: any = await main(...[ {spread} ]); + console.log("wm_res:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value) + '\n'); + }} catch (e) {{ + console.log("wm_res:" + JSON.stringify({{ error: {{ message: e.message, name: e.name, stack: e.stack, line: line }}}}) + '\n'); + }} + }} + if (exit) {{ + break; + }} +}} +"#, + ); + write_file(job_dir, "wrapper.ts", &wrapper_content).await?; + } + + build_import_map(w_id, script_path, base_internal_url, job_dir).await?; + + handle_dedicated_process( + &*DENO_PATH, + job_dir, + context_envs, + envs, + context, + common_deno_proc_envs, + vec![ + "run", + "--no-check", + "--import-map", + &format!("{job_dir}/import_map.json"), + &format!("--reload={base_internal_url}"), + "--unstable", + "-A", + &format!("{job_dir}/wrapper.ts"), + ], + killpill_rx, + job_completed_tx, + token, + jobs_rx, + worker_name, + ) + .await +} diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index cb952443e7..5de0fda598 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -1212,6 +1212,22 @@ pub async fn run_worker { + crate::deno_executor::start_worker( + &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 b3468d1141..3067f1ac83 100644 --- a/frontend/src/lib/components/ScriptBuilder.svelte +++ b/frontend/src/lib/components/ScriptBuilder.svelte @@ -639,7 +639,8 @@ disabled={!$enterpriseLicense || isCloudHosted() || (script.language != Script.language.BUN && - script.language != Script.language.PYTHON3)} + script.language != Script.language.PYTHON3 && + script.language != Script.language.DENO)} size="sm" checked={Boolean(script.dedicated_worker)} on:change={() => { diff --git a/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte b/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte index edc8029996..17d72530ef 100644 --- a/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte @@ -33,7 +33,7 @@ import HighlightCode from '$lib/components/HighlightCode.svelte' import TestJobLoader from '$lib/components/TestJobLoader.svelte' import LogViewer from '$lib/components/LogViewer.svelte' - import { ActionRow, Button, Popup, Skeleton, Tab, Alert } from '$lib/components/common' + import { ActionRow, Button, Popup, Skeleton, Tab, Alert, MenuItem } from '$lib/components/common' import FlowMetadata from '$lib/components/FlowMetadata.svelte' import JobArgs from '$lib/components/JobArgs.svelte' import FlowProgressBar from '$lib/components/flows/FlowProgressBar.svelte' @@ -43,6 +43,7 @@ import { goto } from '$app/navigation' import { sendUserToast } from '$lib/toast' import { forLater } from '$lib/forLater' + import ButtonDropdown from '$lib/components/common/button/ButtonDropdown.svelte' let job: Job | undefined @@ -186,30 +187,35 @@ {@const isScript = job?.job_kind === 'script'} {@const runsHref = `/runs/${job?.script_path}${!isScript ? '?jobKind=flow' : ''}`} - {#if job && 'deleted' in job && !job?.deleted && ($superadmin || ($userStore?.is_admin ?? false))} - - {#if job?.job_kind === 'script' || job?.job_kind === 'flow'} - +
+ {#if job && 'deleted' in job && !job?.deleted && ($superadmin || ($userStore?.is_admin ?? false))} + + + + {/if} {/if} - {/if} +
{@const stem = `/${job?.job_kind}s`} @@ -415,10 +421,11 @@ The content of this run was deleted (by an admin, no less) +
{/if} -
+