mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-03 08:02:19 +00:00
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) <noreply@anthropic.com> * fix: resolve and check every WAC child before the parent parks Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: build parent-code WAC children at push time, not in the pre-pass Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
2a40f25bda
commit
2699d5065e
@@ -1303,6 +1303,62 @@ async fn test_bun_wac_inline_task_cache_is_per_task(db: Pool<Postgres>) -> 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<Postgres>) -> 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
|
||||
// ============================================================================
|
||||
|
||||
@@ -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<std::collections::HashMap<String, windmill_common::scripts::ScriptModule>>,
|
||||
fn resolve_parent_module<'a>(
|
||||
modules: &'a Option<std::collections::HashMap<String, windmill_common::scripts::ScriptModule>>,
|
||||
module_key: &str,
|
||||
) -> error::Result<windmill_common::scripts::ScriptModule> {
|
||||
) -> 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<ResolvedChild> = 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<Uuid> = 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<String, Box<RawValue>> = 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<String, Box<RawValue>> = 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<String, Box<RawValue>> = 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<HashMap<String, Box<RawValue>>> =
|
||||
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,
|
||||
|
||||
Reference in New Issue
Block a user