From 90f9e59321a16ce65ae25b49bda7e632beb313e8 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sun, 27 Sep 2026 16:09:59 +0200 Subject: [PATCH] fix: only let a job's own token claim run lineage (#11367) * fix: only let a job's own token claim its lineage on the run endpoints Co-Authored-By: Claude Opus 5.5 (1M context) * fix: drop an unclaimable run lineage instead of refusing the run Co-Authored-By: Claude Opus 5.5 (1M context) * fix: only let a job's own token run its workflow-as-code tasks Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) --- ...2a50fb1d5072f7177e86bbf0ca942aac0223b.json | 23 +++ ...00311ac84b3f0a0229729b770bcf3cc55f03c.json | 35 +++++ backend/tests/run_lineage_auth.rs | 136 ++++++++++++++++++ backend/windmill-api-jobs/src/execution.rs | 66 +++++++++ backend/windmill-api/src/jobs.rs | 17 +++ 5 files changed, 277 insertions(+) create mode 100644 backend/.sqlx/query-2ab8038052a8fe8ff0e7648909a2a50fb1d5072f7177e86bbf0ca942aac0223b.json create mode 100644 backend/.sqlx/query-2b40bde4c0f78a689ec7e0f5b2600311ac84b3f0a0229729b770bcf3cc55f03c.json create mode 100644 backend/tests/run_lineage_auth.rs diff --git a/backend/.sqlx/query-2ab8038052a8fe8ff0e7648909a2a50fb1d5072f7177e86bbf0ca942aac0223b.json b/backend/.sqlx/query-2ab8038052a8fe8ff0e7648909a2a50fb1d5072f7177e86bbf0ca942aac0223b.json new file mode 100644 index 0000000000..d19056a19a --- /dev/null +++ b/backend/.sqlx/query-2ab8038052a8fe8ff0e7648909a2a50fb1d5072f7177e86bbf0ca942aac0223b.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT id FROM v2_job WHERE id = ANY($1) AND workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "UuidArray", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "2ab8038052a8fe8ff0e7648909a2a50fb1d5072f7177e86bbf0ca942aac0223b" +} diff --git a/backend/.sqlx/query-2b40bde4c0f78a689ec7e0f5b2600311ac84b3f0a0229729b770bcf3cc55f03c.json b/backend/.sqlx/query-2b40bde4c0f78a689ec7e0f5b2600311ac84b3f0a0229729b770bcf3cc55f03c.json new file mode 100644 index 0000000000..13a8e589b4 --- /dev/null +++ b/backend/.sqlx/query-2b40bde4c0f78a689ec7e0f5b2600311ac84b3f0a0229729b770bcf3cc55f03c.json @@ -0,0 +1,35 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT parent_job, root_job, flow_innermost_root_job FROM v2_job\n WHERE id = $1 AND workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "parent_job", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "root_job", + "type_info": "Uuid" + }, + { + "ordinal": 2, + "name": "flow_innermost_root_job", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + true, + true, + true + ] + }, + "hash": "2b40bde4c0f78a689ec7e0f5b2600311ac84b3f0a0229729b770bcf3cc55f03c" +} diff --git a/backend/tests/run_lineage_auth.rs b/backend/tests/run_lineage_auth.rs new file mode 100644 index 0000000000..4f154cdc69 --- /dev/null +++ b/backend/tests/run_lineage_auth.rs @@ -0,0 +1,136 @@ +//! `parent_job` / `root_job` on the run endpoints set the lineage a job is identified by +//! (`WM_FLOW_JOB_ID`, `WM_ROOT_FLOW_JOB_ID`, the OIDC token's flow path). Only a job's own +//! `WM_TOKEN` may claim that job or its ancestors, and a workspace admin any job of the +//! workspace; a lineage anyone else names is dropped from the pushed job. + +use reqwest::StatusCode; +use sqlx::{Pool, Postgres}; +use uuid::Uuid; +use windmill_test_utils::*; + +const ADMIN_JOB: &str = "a0000000-0000-0000-0000-000000000001"; +const ROOT_JOB: &str = "a0000000-0000-0000-0000-000000000002"; +const USER_JOB: &str = "a0000000-0000-0000-0000-000000000003"; + +async fn insert_job( + db: &Pool, + id: &str, + user: &str, + email: &str, + root: Option<&str>, +) -> anyhow::Result<()> { + let root = root.map(Uuid::parse_str).transpose()?; + sqlx::query( + "INSERT INTO v2_job (id, workspace_id, created_by, permissioned_as, permissioned_as_email, + kind, script_lang, runnable_path, tag, parent_job, root_job, flow_innermost_root_job) + VALUES ($1, 'test-workspace', $2, 'u/' || $2, $3, 'script', 'deno', 'u/x/y', 'deno', + $4, $4, $4)", + ) + .bind(Uuid::parse_str(id)?) + .bind(user) + .bind(email) + .bind(root) + .execute(db) + .await?; + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn test_run_lineage_must_be_the_callers_own(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + + sqlx::query( + "INSERT INTO script (workspace_id, created_by, content, schema, summary, description, + path, hash, language, lock, kind) + VALUES ('test-workspace', 'test-user-3', 'export function main() {}', '{}', '', '', + 'u/test-user-3/probe', 424242, 'deno', '', 'script')", + ) + .execute(&db) + .await?; + insert_job(&db, ADMIN_JOB, "test-user", "test@windmill.dev", None).await?; + insert_job(&db, ROOT_JOB, "test-user-3", "test3@windmill.dev", None).await?; + insert_job( + &db, + USER_JOB, + "test-user-3", + "test3@windmill.dev", + Some(ROOT_JOB), + ) + .await?; + + let server = ApiServer::start(db.clone()).await?; + set_jwt_secret().await; + let job_token = windmill_common::auth::create_token_for_owner( + &db, + "test-workspace", + "u/test-user-3", + "ephemeral-script", + 300, + "test3@windmill.dev", + &Uuid::parse_str(USER_JOB)?, + None, + None, + ) + .await?; + + let client = reqwest::Client::new(); + let cases = [ + // a user's own token claims neither someone else's job nor even its own + ( + "SECRET_TOKEN_3", + Some(ADMIN_JOB), + Some(ADMIN_JOB), + None, + None, + ), + ("SECRET_TOKEN_3", Some(USER_JOB), None, None, None), + // what the SDKs send from inside a job: that job and its root flow + ( + &job_token, + Some(USER_JOB), + Some(ROOT_JOB), + Some(USER_JOB), + Some(ROOT_JOB), + ), + ( + &job_token, + Some(USER_JOB), + Some(ADMIN_JOB), + Some(USER_JOB), + None, + ), + ("SECRET_TOKEN", Some(USER_JOB), None, Some(USER_JOB), None), + ]; + for (token, parent, root, kept_parent, kept_root) in cases { + let query = [("parent_job", parent), ("root_job", root)] + .into_iter() + .filter_map(|(k, v)| v.map(|v| format!("{k}={v}"))) + .collect::>() + .join("&"); + let resp = client + .post(format!( + "http://localhost:{}/api/w/test-workspace/jobs/run/p/u/test-user-3/probe?{query}", + server.addr.port() + )) + .bearer_auth(token) + .json(&serde_json::json!({})) + .send() + .await?; + let status = resp.status(); + let body = resp.text().await?; + assert_eq!(status, StatusCode::CREATED, "{query}: {body}"); + let (stored_parent, stored_root): (Option, Option) = + sqlx::query_as("SELECT parent_job, flow_innermost_root_job FROM v2_job WHERE id = $1") + .bind(Uuid::parse_str(&body)?) + .fetch_one(&db) + .await?; + let expected = |id: Option<&str>| id.map(Uuid::parse_str).transpose(); + assert_eq!( + stored_parent, + expected(kept_parent)?, + "parent_job for {query}" + ); + assert_eq!(stored_root, expected(kept_root)?, "root_job for {query}"); + } + Ok(()) +} diff --git a/backend/windmill-api-jobs/src/execution.rs b/backend/windmill-api-jobs/src/execution.rs index 9ad7189434..fc64a693c1 100644 --- a/backend/windmill-api-jobs/src/execution.rs +++ b/backend/windmill-api-jobs/src/execution.rs @@ -95,6 +95,68 @@ pub async fn check_tag_as_written_available_for_workspace( } } +/// Clears a `parent_job` / `root_job` the caller cannot claim. They identify the pushed job +/// (`WM_FLOW_JOB_ID`, `WM_ROOT_FLOW_JOB_ID`, the OIDC token's flow path), so only a job's own +/// `WM_TOKEN` may name that job or its ancestors, and a workspace admin any job of the workspace. +/// Cleared rather than refused: the SDKs send them from inside a job whatever token they hold. +pub async fn drop_unclaimable_run_lineage( + db: &DB, + w_id: &str, + run_query: &mut RunJobQuery, + authed: &ApiAuthed, +) -> error::Result<()> { + let mut referenced: Vec = [run_query.parent_job, run_query.root_job] + .into_iter() + .flatten() + .collect(); + referenced.dedup(); + if let Some(token_job) = authed.job_id { + referenced.retain(|id| *id != token_job); + if !referenced.is_empty() { + if let Some(lineage) = sqlx::query!( + "SELECT parent_job, root_job, flow_innermost_root_job FROM v2_job + WHERE id = $1 AND workspace_id = $2", + token_job, + w_id + ) + .fetch_optional(db) + .await? + { + let ancestors = [ + lineage.parent_job, + lineage.root_job, + lineage.flow_innermost_root_job, + ]; + referenced.retain(|id| !ancestors.contains(&Some(*id))); + } + } + } + if referenced.is_empty() { + return Ok(()); + } + let in_workspace = if authed.is_admin { + sqlx::query_scalar!( + "SELECT id FROM v2_job WHERE id = ANY($1) AND workspace_id = $2", + &referenced, + w_id, + ) + .fetch_all(db) + .await? + } else { + vec![] + }; + for field in [&mut run_query.parent_job, &mut run_query.root_job] { + if field.is_some_and(|id| referenced.contains(&id) && !in_workspace.contains(&id)) { + tracing::warn!( + "ignoring parent_job/root_job {field:?} that {} cannot claim in {w_id}", + authed.username + ); + *field = None; + } + } + Ok(()) +} + #[cfg(feature = "enterprise")] pub async fn check_license_key_valid() -> error::Result<()> { use windmill_common::ee_oss::LICENSE_KEY_VALID; @@ -776,6 +838,8 @@ pub async fn run_flow<'c>( bool, Option>, )> { + let mut run_query = run_query; + drop_unclaimable_run_lineage(db, w_id, &mut run_query, authed).await?; let on_behalf_of = flow_version_info.on_behalf_of(w_id, &db).await?; let FlowVersionInfo { version, @@ -1017,6 +1081,8 @@ pub async fn push_script_job_by_path_into_queue<'c>( let script_path = script_path.to_path(); check_scopes(&authed, || format!("jobs:run:scripts:{script_path}"))?; + let mut run_query = run_query; + drop_unclaimable_run_lineage(&db, &w_id, &mut run_query, &authed).await?; let userdb_authed = UserDbWithAuthed { db: user_db.clone(), authed: &authed.to_authed_ref() }; let (job_payload, tag, delete_after_use, delete_after_secs, timeout, on_behalf_of) = diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 60a8037330..e23169e50b 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -7147,6 +7147,8 @@ pub async fn restart_flow( .script_path .with_context(|| "No flow path set for completed flow job")?; check_scopes(&authed, || format!("jobs:run:flows:{flow_path}"))?; + let mut run_query = run_query; + drop_unclaimable_run_lineage(&db, &w_id, &mut run_query, &authed).await?; let ehm = HashMap::new(); let push_args = completed_job @@ -7345,6 +7347,13 @@ pub async fn run_workflow_as_code( let args = PushArgs { args: &task.args.unwrap_or_else(HashMap::new), extra: Some(extra) }; check_tag_available_for_workspace(&db, &w_id, &run_query.tag, &args, &authed).await?; check_scopes(&authed, || format!("jobs:run"))?; + // The task becomes a child of `job_id`, runs its code and writes into its flow status, so + // only that job itself (the SDK's `task` wrapper, on its `WM_TOKEN`) or an admin may push it. + if authed.job_id != Some(job_id) && !authed.is_admin { + return Err(error::Error::PermissionDenied(format!( + "only job {job_id}'s own WM_TOKEN can run its workflow tasks" + ))); + } if !is_valid_entrypoint_name(&entrypoint) { return Err(error::Error::BadRequest(format!( @@ -7761,6 +7770,8 @@ pub async fn run_wait_result_job_by_path_get( let tag = run_query.tag.clone().or(tag); let push_args = PushArgs { args: &args.args, extra: args.extra }; check_tag_available_for_workspace(&db, &w_id, &tag, &push_args, &authed).await?; + let mut run_query = run_query; + drop_unclaimable_run_lineage(&db, &w_id, &mut run_query, &authed).await?; let (email, permissioned_as, push_authed, tx) = if let Some(on_behalf_of) = on_behalf_authed.as_ref() { @@ -7907,6 +7918,8 @@ pub async fn run_wait_result_script_by_path_internal( let tag = run_query.tag.clone().or(tag); let push_args = PushArgs { args: &args.args, extra: args.extra }; check_tag_available_for_workspace(&db, &w_id, &tag, &push_args, &authed).await?; + let mut run_query = run_query; + drop_unclaimable_run_lineage(&db, &w_id, &mut run_query, &authed).await?; let (email, permissioned_as, push_authed, tx) = if let Some(on_behalf_of) = on_behalf_of.as_ref() { @@ -8021,6 +8034,8 @@ pub async fn run_wait_result_script_by_hash( let tag = run_query.tag.clone().or(tag); let push_args = PushArgs { args: &args.args, extra: args.extra }; check_tag_available_for_workspace(&db, &w_id, &tag, &push_args, &authed).await?; + let mut run_query = run_query; + drop_unclaimable_run_lineage(&db, &w_id, &mut run_query, &authed).await?; let (email, permissioned_as, push_authed, tx) = if let Some(obo) = on_behalf_of.as_ref() { ( @@ -10040,6 +10055,8 @@ pub async fn run_job_by_hash_inner( let push_args = PushArgs { args: &args.args, extra: args.extra }; check_tag_available_for_workspace(&db, &w_id, &tag, &push_args, &authed).await?; + let mut run_query = run_query; + drop_unclaimable_run_lineage(&db, &w_id, &mut run_query, &authed).await?; let (email, permissioned_as, push_authed, tx) = if let Some(obo) = on_behalf_of.as_ref() { (