From 2699d5065ea3dfe985636733fb6022ca23187291 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Tue, 22 Sep 2026 00:54:34 +0200 Subject: [PATCH] fix: check workflow-as-code task tags against CUSTOM_TAGS (#11273) * fix: check workflow-as-code task tags against CUSTOM_TAGS Co-Authored-By: Claude Opus 5 (1M context) * fix: resolve and check every WAC child before the parent parks Co-Authored-By: Claude Opus 5 (1M context) * fix: build parent-code WAC children at push time, not in the pre-pass Co-Authored-By: Claude Opus 5 (1M context) --------- Co-authored-by: Claude Opus 5 (1M context) --- backend/tests/bun_jobs.rs | 56 ++++ backend/windmill-worker/src/bun_executor.rs | 269 +++++++++++--------- 2 files changed, 209 insertions(+), 116 deletions(-) diff --git a/backend/tests/bun_jobs.rs b/backend/tests/bun_jobs.rs index 63d5341ba3..14f8c422a4 100644 --- a/backend/tests/bun_jobs.rs +++ b/backend/tests/bun_jobs.rs @@ -1303,6 +1303,62 @@ async fn test_bun_wac_inline_task_cache_is_per_task(db: Pool) -> anyho Ok(()) } +/// A task that picks its own worker tag is held to CUSTOM_TAGS like a flow step: with none +/// set, only a superadmin may send it to `gpu`, so the workflow fails instead of queueing it. +#[sqlx::test(fixtures("base"))] +async fn test_bun_wac_task_tag_outside_custom_tags(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#" +import { workflow, task } from "windmill-client"; + +const onGpu = task(async (n: number) => { + return n * 2; +}, { tag: "gpu" }); + +export const main = workflow(async (n: number) => { + return await onGpu(n); +}); +"# + .to_owned(); + + let job = RunJob::from(JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: windmill_common::runnable_settings::ConcurrencySettings::default() + .into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + modules: None, + tag: None, + })) + .arg("n", serde_json::json!(5)) + .as_user("test-user-2", "test2@windmill.dev") + .run_until_complete(&db, false, port) + .await; + + assert!(!job.success, "the workflow must fail"); + let result = job.json_result().unwrap(); + let message = result["error"]["message"].as_str().unwrap_or_default(); + assert!( + message.contains("task 'onGpu' cannot run on tag 'gpu'"), + "unexpected error: {result}" + ); + let children: i64 = sqlx::query_scalar("SELECT count(*) FROM v2_job WHERE parent_job = $1") + .bind(job.id) + .fetch_one(&db) + .await?; + assert_eq!(children, 0, "no child is queued on the refused tag"); + Ok(()) +} + // ============================================================================ // Environment Variable Tests // ============================================================================ diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 117fbda086..7ed7b8d5da 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -2689,13 +2689,13 @@ try {{ /// Resolve a module file from the parent script's modules map. /// For Script jobs, fetches from the `script` table by hash. /// For Preview jobs, fetches from `v2_job.raw_code` (modules stored inline). -fn resolve_parent_module( - modules: &Option>, +fn resolve_parent_module<'a>( + modules: &'a Option>, module_key: &str, -) -> error::Result { +) -> error::Result<&'a windmill_common::scripts::ScriptModule> { if let Some(modules) = modules { if let Some(module) = modules.get(module_key) { - return Ok(module.clone()); + return Ok(module); } } Err(error::Error::ExecutionErr(format!( @@ -2720,7 +2720,10 @@ pub async fn handle_wac_v2_output( }; use serde_json::Value; use windmill_common::get_latest_flow_version_info_for_path; - use windmill_common::jobs::{script_path_to_payload, JobKind, JobPayload, RawCode}; + use windmill_common::jobs::{ + check_tag_available_for_workspace_internal, script_path_to_payload, JobKind, JobPayload, + RawCode, + }; use windmill_common::runnable_settings::{ ConcurrencySettings, ConcurrencySettingsWithCustom, DebouncingSettings, }; @@ -2950,6 +2953,115 @@ pub async fn handle_wac_v2_output( ))), }?; + // Resolve what every task runs, and check any tag it picks, before the parent parks: + // a task that cannot be pushed then fails the dispatch with no sibling queued. The + // parent's code and input are built per child at push time, so a fan-out never holds + // a copy of them for each child at once. + enum ChildRunnable<'a> { + // Re-runs the parent with `_executing_key`. + Parent, + // A `./` module of the parent script. + Module(&'a windmill_common::scripts::ScriptModule), + // A deployed script or flow. + Deployed(JobPayload), + } + struct ResolvedChild<'a> { + runnable: ChildRunnable<'a>, + email: String, + permissioned_as: String, + } + let mut children: Vec = Vec::with_capacity(num_steps); + for step in &steps { + let (runnable, on_behalf_of) = match step.dispatch_type.as_str() { + "script" if step.script.starts_with("./") => { + // Module-relative path: resolve from parent script's modules + let module_key = step.script.strip_prefix("./").unwrap(); + let module = resolve_parent_module(modules, module_key)?; + // Inline module code, not a separate runnable: it has no + // identity of its own and runs as the parent. + (ChildRunnable::Module(module), None) + } + "script" => { + // Resolve script path to job payload (handles hash, lang, etc.) + let (payload, _, _, _, _, on_behalf_of) = script_path_to_payload( + &step.script, + None, // no authed db for background workers + db.clone(), + &job.workspace_id, + Some(true), // skip preprocessor + ) + .await?; + (ChildRunnable::Deployed(payload), on_behalf_of) + } + "flow" => { + let flow_info = get_latest_flow_version_info_for_path( + None, + db, + &job.workspace_id, + &step.script, + true, + ) + .await?; + let payload = JobPayload::Flow { + path: step.script.clone(), + dedicated_worker: flow_info.dedicated_worker, + apply_preprocessor: false, + version: flow_info.version, + labels: flow_info.labels.clone(), + }; + let on_behalf_of = flow_info.on_behalf_of(&job.workspace_id, db).await?; + (ChildRunnable::Deployed(payload), on_behalf_of) + } + // "inline" — re-run parent with _executing_key + _ => (ChildRunnable::Parent, None), + }; + + // A target runnable that opts into on-behalf-of runs under its own + // identity, never the caller's, so a step that reaches it through a + // workflow cannot widen or narrow its permissions. `created_by` still + // credits the caller, matching how the run API pushes these jobs. + let (email, permissioned_as) = match on_behalf_of { + Some(on_behalf_of) => (on_behalf_of.email, on_behalf_of.permissioned_as), + None => ( + job.permissioned_as_email.clone(), + job.permissioned_as.clone(), + ), + }; + + // A task inheriting the parent's tag skips the check: that tag was + // checked when the parent was pushed, and a dedicated worker's tag is + // never in CUSTOM_TAGS. + if let Some(tag) = step + .tag + .as_deref() + .filter(|t| !t.is_empty() && *t != job.tag.as_str()) + { + let is_super_admin = + windmill_common::auth::is_super_admin_email(db, &email).await?; + check_tag_available_for_workspace_internal( + db, + &job.workspace_id, + tag, + is_super_admin, + None, + ) + .warn_after_seconds_with_sql( + 1, + "check_tag_available_for_workspace_internal".to_string(), + ) + .await + .map_err(|e| match e { + error::Error::BadRequest(msg) => error::Error::BadRequest(format!( + "task '{}' cannot run on tag '{tag}': {msg}", + step.name + )), + e => e, + })?; + } + + children.push(ResolvedChild { runnable, email, permissioned_as }); + } + // Step 1: Save checkpoint, suspend parent, and seed child checkpoints // in a single transaction — all BEFORE children become visible. let segment_ms; @@ -3034,106 +3146,46 @@ pub async fn handle_wac_v2_output( // partial failure (e.g. pushing child 3 of 5 fails). let mut pushed_ids: Vec = Vec::with_capacity(num_steps); let push_result: error::Result<()> = async { - for (step, (_, child_uuid)) in steps.iter().zip(job_ids.iter()) { + for ((step, (_, child_uuid)), child) in + steps.iter().zip(job_ids.iter()).zip(children) + { // A task with a runnable of its own (a deployed script or flow) queues // at that runnable's priority; any other task is the parent's code and // queues at the parent's. let own_runnable = matches!(step.dispatch_type.as_str(), "script" | "flow") && !step.script.starts_with("./"); - // Resolve job payload based on dispatch_type - let (job_payload, child_args, is_external, on_behalf_of) = - match step.dispatch_type.as_str() { - "script" if step.script.starts_with("./") => { - // Module-relative path: resolve from parent script's modules - let module_key = step.script.strip_prefix("./").unwrap(); - let module = resolve_parent_module(modules, module_key)?; - let payload = JobPayload::Code(RawCode { - content: module.content, - path: job.runnable_path.clone(), - hash: None, - language: module.language, - lock: module.lock, - cache_ttl: None, - cache_ignore_s3_path: None, - dedicated_worker: None, - concurrency_settings: ConcurrencySettingsWithCustom::default(), - debouncing_settings: DebouncingSettings::default(), - modules: None, - tag: None, - }); - let step_args: HashMap> = step - .args - .iter() - .map(|(k, v)| { - let raw = serde_json::value::to_raw_value(v).unwrap(); - (k.clone(), raw) - }) - .collect(); - // Inline module code, not a separate runnable: it has no - // identity of its own and runs as the parent. - (payload, step_args, true, None) - } - "script" => { - // Resolve script path to job payload (handles hash, lang, etc.) - let (payload, _, _, _, _, on_behalf_of) = script_path_to_payload( - &step.script, - None, // no authed db for background workers - db.clone(), - &job.workspace_id, - Some(true), // skip preprocessor - ) - .await?; - let step_args: HashMap> = step - .args - .iter() - .map(|(k, v)| { - let raw = serde_json::value::to_raw_value(v).unwrap(); - (k.clone(), raw) - }) - .collect(); - (payload, step_args, true, on_behalf_of) - } - "flow" => { - let flow_info = get_latest_flow_version_info_for_path( - None, - db, - &job.workspace_id, - &step.script, - true, - ) - .await?; - let payload = JobPayload::Flow { - path: step.script.clone(), - dedicated_worker: flow_info.dedicated_worker, - apply_preprocessor: false, - version: flow_info.version, - labels: flow_info.labels.clone(), - }; - let on_behalf_of = - flow_info.on_behalf_of(&job.workspace_id, db).await?; - let step_args: HashMap> = step - .args - .iter() - .map(|(k, v)| { - let raw = serde_json::value::to_raw_value(v).unwrap(); - (k.clone(), raw) - }) - .collect(); - (payload, step_args, true, on_behalf_of) - } - _ => { - // "inline" — re-run parent with _executing_key - ( - job_payload_template.clone(), - parent_args.clone(), - false, - None, - ) - } - }; - - let push_args = PushArgs { args: &child_args, extra: None }; + let is_external = !matches!(child.runnable, ChildRunnable::Parent); + let job_payload = match child.runnable { + ChildRunnable::Parent => job_payload_template.clone(), + ChildRunnable::Module(module) => JobPayload::Code(RawCode { + content: module.content.clone(), + path: job.runnable_path.clone(), + hash: None, + language: module.language, + lock: module.lock.clone(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + concurrency_settings: ConcurrencySettingsWithCustom::default(), + debouncing_settings: DebouncingSettings::default(), + modules: None, + tag: None, + }), + ChildRunnable::Deployed(payload) => payload, + }; + let step_args: Option>> = + is_external.then(|| { + step.args + .iter() + .map(|(k, v)| { + let raw = serde_json::value::to_raw_value(v).unwrap(); + (k.clone(), raw) + }) + .collect() + }); + let push_args = + PushArgs { args: step_args.as_ref().unwrap_or(&parent_args), extra: None }; // Apply step-level overrides to payload (cache, concurrency) let mut job_payload = job_payload; @@ -3181,21 +3233,6 @@ pub async fn handle_wac_v2_output( } } - // A target runnable that opts into on-behalf-of runs under its own - // identity, never the caller's, so a step that reaches it through a - // workflow cannot widen or narrow its permissions. `created_by` still - // credits the caller, matching how the run API pushes these jobs. - let (child_email, child_permissioned_as) = match on_behalf_of.as_ref() { - Some(on_behalf_of) => ( - on_behalf_of.email.as_str(), - on_behalf_of.permissioned_as.clone(), - ), - None => ( - job.permissioned_as_email.as_str(), - job.permissioned_as.clone(), - ), - }; - let (_, mut tx) = push( db, PushIsolationLevel::IsolatedRoot(db.clone()), @@ -3203,8 +3240,8 @@ pub async fn handle_wac_v2_output( job_payload, push_args, &job.created_by, - child_email, - child_permissioned_as, + &child.email, + child.permissioned_as, None, None, None,