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) <noreply@anthropic.com>

* fix: drop an unclaimable run lineage instead of refusing the run

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: only let a job's own token run its workflow-as-code tasks

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-09-27 14:09:59 +00:00
committed by GitHub
co-authored by Claude Opus 5.5
parent 47525b211a
commit 90f9e59321
5 changed files with 277 additions and 0 deletions
@@ -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"
}
@@ -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"
}
+136
View File
@@ -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<Postgres>,
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<Postgres>) -> 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::<Vec<_>>()
.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<Uuid>, Option<Uuid>) =
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(())
}
@@ -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<Uuid> = [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<sqlx::Transaction<'c, sqlx::Postgres>>,
)> {
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) =
+17
View File
@@ -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() {
(