diff --git a/backend/.sqlx/query-3e68b73a1c2fdcb2ef538de1addbb096921e26fed398d5c0966b60503ca97fdb.json b/backend/.sqlx/query-11b5d006acf040b2461b9539863b1e511c23e1a5f9be97da066a6808a49f2092.json similarity index 53% rename from backend/.sqlx/query-3e68b73a1c2fdcb2ef538de1addbb096921e26fed398d5c0966b60503ca97fdb.json rename to backend/.sqlx/query-11b5d006acf040b2461b9539863b1e511c23e1a5f9be97da066a6808a49f2092.json index a40cc12148..bc821f2b7b 100644 --- a/backend/.sqlx/query-3e68b73a1c2fdcb2ef538de1addbb096921e26fed398d5c0966b60503ca97fdb.json +++ b/backend/.sqlx/query-11b5d006acf040b2461b9539863b1e511c23e1a5f9be97da066a6808a49f2092.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH job_info AS (\n SELECT id, kind::text AS kind, parent_job\n FROM v2_job\n WHERE id = $1\n )\n SELECT\n q.id AS \"id!\",\n s.flow_status,\n q.suspend AS \"suspend!\",\n j.runnable_path AS script_path,\n j.permissioned_as_email AS email,\n (ji.kind IN ('flow', 'flowpreview', 'singlestepflow')) AS \"is_flow_level!\",\n (ji.kind NOT IN ('flow', 'flowpreview', 'singlestepflow') AND q.id = ji.id) AS \"is_wac!\"\n FROM job_info ji\n JOIN v2_job_queue q ON q.id = CASE\n WHEN ji.kind IN ('flow', 'flowpreview', 'singlestepflow') THEN ji.id\n ELSE COALESCE(ji.parent_job, ji.id)\n END\n JOIN v2_job j ON j.id = q.id\n LEFT JOIN v2_job_status s ON s.id = q.id\n FOR UPDATE OF q\n ", + "query": "\n WITH job_info AS (\n SELECT id, kind::text AS kind, parent_job\n FROM v2_job\n WHERE id = $1 AND workspace_id = $2\n )\n SELECT\n q.id AS \"id!\",\n s.flow_status,\n q.suspend AS \"suspend!\",\n j.runnable_path AS script_path,\n j.permissioned_as_email AS email,\n (ji.kind IN ('flow', 'flowpreview', 'singlestepflow')) AS \"is_flow_level!\",\n (ji.kind NOT IN ('flow', 'flowpreview', 'singlestepflow') AND q.id = ji.id) AS \"is_wac!\"\n FROM job_info ji\n JOIN v2_job_queue q ON q.id = CASE\n WHEN ji.kind IN ('flow', 'flowpreview', 'singlestepflow') THEN ji.id\n ELSE COALESCE(ji.parent_job, ji.id)\n END\n JOIN v2_job j ON j.id = q.id\n LEFT JOIN v2_job_status s ON s.id = q.id\n FOR UPDATE OF q\n ", "describe": { "columns": [ { @@ -41,7 +41,8 @@ ], "parameters": { "Left": [ - "Uuid" + "Uuid", + "Text" ] }, "nullable": [ @@ -54,5 +55,5 @@ null ] }, - "hash": "3e68b73a1c2fdcb2ef538de1addbb096921e26fed398d5c0966b60503ca97fdb" + "hash": "11b5d006acf040b2461b9539863b1e511c23e1a5f9be97da066a6808a49f2092" } diff --git a/backend/.sqlx/query-f3bae157b92cbc1b3741593a620c9203012a7d3754f56af665e2f34a969eabd7.json b/backend/.sqlx/query-f3bae157b92cbc1b3741593a620c9203012a7d3754f56af665e2f34a969eabd7.json deleted file mode 100644 index f787574b3a..0000000000 --- a/backend/.sqlx/query-f3bae157b92cbc1b3741593a620c9203012a7d3754f56af665e2f34a969eabd7.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT EXISTS (SELECT 1 FROM v2_job WHERE id = $1 AND workspace_id = $2)", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "exists", - "type_info": "Bool" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Text" - ] - }, - "nullable": [ - null - ] - }, - "hash": "f3bae157b92cbc1b3741593a620c9203012a7d3754f56af665e2f34a969eabd7" -} diff --git a/backend/tests/jobs_read_auth.rs b/backend/tests/jobs_read_auth.rs index 8c5bfee0eb..9325f25652 100644 --- a/backend/tests/jobs_read_auth.rs +++ b/backend/tests/jobs_read_auth.rs @@ -513,12 +513,12 @@ async fn test_single_job_read_authorization(db: Pool) -> anyhow::Resul let (status, secret) = get( &authed_base, &format!("job_signature/{STEP_JOB}/0"), - Some("SECRET_TOKEN_2"), + Some("SECRET_TOKEN"), ) .await; assert!( status.is_success(), - "owner must mint a resume secret: {secret}" + "an admin must mint a resume secret: {secret}" ); let secret = secret.trim().trim_matches('"').to_string(); let approval_result = format!( diff --git a/backend/tests/wac_approval_urls.rs b/backend/tests/wac_approval_urls.rs index 0cb33e67c0..7883c69b17 100644 --- a/backend/tests/wac_approval_urls.rs +++ b/backend/tests/wac_approval_urls.rs @@ -201,7 +201,10 @@ async fn wac_approval_urls_mint_guards(db: Pool) -> anyhow::Result<()> let client = reqwest::Client::new(); let mint = |ws: &str, job: &str, key: &str| { let url = format!("{base}/w/{ws}/jobs/wac_approval_urls/{job}/{key}"); - client.get(url).header("Authorization", "Bearer SECRET_TOKEN").send() + client + .get(url) + .header("Authorization", "Bearer SECRET_TOKEN") + .send() }; assert_eq!( @@ -227,8 +230,15 @@ async fn wac_approval_urls_mint_guards(db: Pool) -> anyhow::Result<()> let resp = mint("test-workspace", WAC_JOB, APPROVAL_COLLISION_B).await?; let status = resp.status(); let body = resp.text().await.unwrap_or_default(); - assert_eq!(status, reqwest::StatusCode::BAD_REQUEST, "colliding key: {body}"); - assert!(body.contains(APPROVAL_COLLISION_A), "error must name the other key: {body}"); + assert_eq!( + status, + reqwest::StatusCode::BAD_REQUEST, + "colliding key: {body}" + ); + assert!( + body.contains(APPROVAL_COLLISION_A), + "error must name the other key: {body}" + ); Ok(()) } @@ -236,7 +246,9 @@ async fn wac_approval_urls_mint_guards(db: Pool) -> anyhow::Result<()> /// Colliding keys share one resume_job row and one capability, so the mint that /// records them must let exactly one through however the requests interleave. #[sqlx::test(fixtures("base", "wac_approval_urls"))] -async fn wac_concurrent_colliding_mints_admit_exactly_one(db: Pool) -> anyhow::Result<()> { +async fn wac_concurrent_colliding_mints_admit_exactly_one( + db: Pool, +) -> anyhow::Result<()> { initialize_tracing().await; let server = ApiServer::start(db.clone()).await?; let base = format!("http://localhost:{}/api", server.addr.port()); @@ -262,6 +274,92 @@ async fn wac_concurrent_colliding_mints_admit_exactly_one(db: Pool) -> other => anyhow::bail!("unexpected status {other}"), } } - assert_eq!((ok, rejected), (1, 1), "exactly one colliding key may be minted"); + assert_eq!( + (ok, rejected), + (1, 1), + "exactly one colliding key may be minted" + ); + Ok(()) +} + +/// A resume signature skips the step's approval_conditions, so only the suspended job's own +/// run (its `WM_TOKEN`) or a workspace admin may mint one, and a signature made with another +/// workspace's key must not reach a job outside that workspace. +#[sqlx::test(fixtures("base", "wac_approval_urls"))] +async fn resume_signatures_are_bound_to_the_run_and_its_workspace( + db: Pool, +) -> anyhow::Result<()> { + use hmac::Mac; + initialize_tracing().await; + + let server = ApiServer::start(db.clone()).await?; + let base = format!("http://localhost:{}/api", server.addr.port()); + let client = reqwest::Client::new(); + set_jwt_secret().await; + let job_token = windmill_common::auth::create_token_for_owner( + &db, + "test-workspace", + "u/test-user-2", + "ephemeral-script", + 300, + "test2@windmill.dev", + &uuid::Uuid::parse_str(WAC_JOB)?, + None, + None, + ) + .await?; + + for route in ["job_signature", "resume_urls"] { + for (token, allowed) in [ + ("SECRET_TOKEN_2", false), + (&job_token, true), + ("SECRET_TOKEN", true), + ] { + let status = client + .get(format!("{base}/w/test-workspace/jobs/{route}/{WAC_JOB}/7")) + .bearer_auth(token) + .send() + .await? + .status(); + assert_eq!( + status.is_success(), + allowed, + "{route} with a {} token: {status}", + if token == job_token { "job" } else { token } + ); + } + } + + let mut mac = hmac::Hmac::::new_from_slice(b"test-key-2")?; + mac.update(uuid::Uuid::parse_str(WAC_JOB)?.as_bytes()); + mac.update(&5u32.to_be_bytes()); + let foreign = hex::encode(mac.finalize().into_bytes()); + for op in ["resume", "cancel"] { + let status = client + .post(format!( + "{base}/w/test-workspace-2/jobs_u/{op}/{WAC_JOB}/5/{foreign}" + )) + .json(&serde_json::json!({})) + .send() + .await? + .status(); + assert_eq!( + status, + reqwest::StatusCode::NOT_FOUND, + "{op} from test-workspace-2" + ); + } + let (suspend, resumes): (i32, i64) = sqlx::query_as( + "SELECT suspend, (SELECT count(*) FROM resume_job WHERE job = $1::uuid) + FROM v2_job_queue WHERE id = $1::uuid", + ) + .bind(WAC_JOB) + .fetch_one(&db) + .await?; + assert_eq!( + (suspend, resumes), + (1, 0), + "the job must still be suspended" + ); Ok(()) } diff --git a/backend/windmill-api-jobs/src/execution.rs b/backend/windmill-api-jobs/src/execution.rs index fc64a693c1..1a750c38ba 100644 --- a/backend/windmill-api-jobs/src/execution.rs +++ b/backend/windmill-api-jobs/src/execution.rs @@ -105,10 +105,33 @@ pub async fn drop_unclaimable_run_lineage( run_query: &mut RunJobQuery, authed: &ApiAuthed, ) -> error::Result<()> { - let mut referenced: Vec = [run_query.parent_job, run_query.root_job] + let referenced: Vec = [run_query.parent_job, run_query.root_job] .into_iter() .flatten() .collect(); + let unclaimable = unclaimable_run_lineage(db, w_id, referenced, authed).await?; + for field in [&mut run_query.parent_job, &mut run_query.root_job] { + if field.is_some_and(|id| unclaimable.contains(&id)) { + tracing::warn!( + "ignoring parent_job/root_job {field:?} that {} cannot claim in {w_id}", + authed.username + ); + *field = None; + } + } + Ok(()) +} + +/// The jobs of `referenced` that `authed` cannot claim as its own run lineage: anything but the +/// token's own job and that job's `parent_job`, `root_job` and `flow_innermost_root_job`, or for +/// a workspace admin anything outside the workspace. +pub async fn unclaimable_run_lineage( + db: &DB, + w_id: &str, + mut referenced: Vec, + authed: &ApiAuthed, +) -> error::Result> { + referenced.sort(); referenced.dedup(); if let Some(token_job) = authed.job_id { referenced.retain(|id| *id != token_job); @@ -132,7 +155,7 @@ pub async fn drop_unclaimable_run_lineage( } } if referenced.is_empty() { - return Ok(()); + return Ok(referenced); } let in_workspace = if authed.is_admin { sqlx::query_scalar!( @@ -145,16 +168,8 @@ pub async fn drop_unclaimable_run_lineage( } 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(()) + referenced.retain(|id| !in_workspace.contains(id)); + Ok(referenced) } #[cfg(feature = "enterprise")] diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 633cc341d2..27e34b6ef2 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -5413,7 +5413,7 @@ async fn resume_suspended_job_internal( verify_suspended_secret(&w_id, &db, job_id, resume_id, &approver, secret).await?; // Get flow info - works for step-level, flow-level, and WAC approval - let (flow_info, is_flow_level, is_wac) = get_flow_info_for_resume(job_id, &db).await?; + let (flow_info, is_flow_level, is_wac) = get_flow_info_for_resume(job_id, &w_id, &db).await?; // HMAC secret = full capability. Skip approval_conditions checks: possession of the full // resume URL is the authorization (it is only disclosed to intended approvers, e.g. when a @@ -5645,7 +5645,11 @@ struct FlowInfo { /// Returns (FlowInfo, is_flow_level, is_wac) where: /// - is_flow_level: job_id was a flow job (pre-approval) /// - is_wac: job_id is a WAC workflow suspended for approval (target is itself) -async fn get_flow_info_for_resume(job_id: Uuid, db: &DB) -> error::Result<(FlowInfo, bool, bool)> { +async fn get_flow_info_for_resume( + job_id: Uuid, + w_id: &str, + db: &DB, +) -> error::Result<(FlowInfo, bool, bool)> { // Single query that determines if job_id is a flow, step, or WAC job, // and fetches the appropriate suspended job info. // For WAC jobs (no parent, not a flow), the job itself is the suspended target. @@ -5654,7 +5658,7 @@ async fn get_flow_info_for_resume(job_id: Uuid, db: &DB) -> error::Result<(FlowI WITH job_info AS ( SELECT id, kind::text AS kind, parent_job FROM v2_job - WHERE id = $1 + WHERE id = $1 AND workspace_id = $2 ) SELECT q.id AS "id!", @@ -5674,10 +5678,15 @@ async fn get_flow_info_for_resume(job_id: Uuid, db: &DB) -> error::Result<(FlowI FOR UPDATE OF q "#, job_id, + w_id, ) .fetch_optional(db) .await? - .ok_or_else(|| anyhow::anyhow!("job not found or parent flow not in queue: {}", job_id))?; + .ok_or_else(|| { + Error::NotFound(format!( + "job not found or parent flow not in queue: {job_id}" + )) + })?; let flow_info = FlowInfo { id: result.id, @@ -5981,12 +5990,7 @@ pub async fn create_job_signature( Path((w_id, job_id, resume_id)): Path<(String, Uuid, u32)>, Query(approver): Query, ) -> error::Result { - // The HMAC is treated as full authority by the resume endpoints, so minting - // it requires run scope on the suspended job's flow — not merely any - // jobs:run scope. No-op for unscoped tokens (incl. the in-flow substep token - // used by wmill.get_resume_urls()). - let flow_path = resume_target_flow_path(&db, &w_id, job_id).await?; - check_scopes(&authed, || format!("jobs:run:flows:{}", flow_path))?; + require_resume_mint_authority(&db, &w_id, &authed, job_id).await?; let key = get_workspace_key(&w_id, &db).await?; create_signature(key, job_id, resume_id, approver.approver) } @@ -6079,12 +6083,7 @@ pub async fn get_resume_urls( Path((w_id, job_id, resume_id)): Path<(String, Uuid, u32)>, Query(approver): Query, ) -> error::JsonResult { - // These URLs embed a resume signature (full resume capability), so a scoped - // token must hold run scope on the suspended job's flow. No-op for unscoped - // tokens (incl. the in-flow substep token). Trusted internal callers use - // get_resume_urls_internal directly and are unaffected. - let flow_path = resume_target_flow_path(&db, &w_id, job_id).await?; - check_scopes(&authed, || format!("jobs:run:flows:{}", flow_path))?; + require_resume_mint_authority(&db, &w_id, &authed, job_id).await?; get_resume_urls_internal( Extension(db), Path((w_id, job_id, resume_id)), @@ -6109,23 +6108,9 @@ pub async fn get_wac_approval_urls( "step_key must be the key of a wait_for_approval step".to_string(), )); } - let flow_path = resume_target_flow_path(&db, &w_id, job_id).await?; - check_scopes(&authed, || format!("jobs:run:flows:{}", flow_path))?; - - // This handler writes to the job's status row, so the job must actually be in - // the caller's workspace — `v2_job_status` is keyed by job id alone and would - // otherwise take a write aimed at another workspace's job. - let in_workspace = sqlx::query_scalar!( - "SELECT EXISTS (SELECT 1 FROM v2_job WHERE id = $1 AND workspace_id = $2)", - job_id, - w_id - ) - .fetch_one(&db) - .await? - .unwrap_or(false); - if !in_workspace { - return Err(Error::NotFound(format!("job {job_id} not found"))); - } + // Also the workspace check for the `v2_job_status` write below, which is keyed by job id + // alone and would otherwise take a write aimed at another workspace's job. + require_resume_mint_authority(&db, &w_id, &authed, job_id).await?; // The approval belongs to the WAC parent, but WM_JOB_ID is the child job when // this is called from inside a task() rather than a step(). Resolve up so the @@ -6325,12 +6310,34 @@ pub async fn get_resume_urls_internal( Ok(Json(res)) } -/// Resolve the runnable path of the flow a (possibly step) job belongs to, used +/// A resume signature is full approval authority: `jobs_u/resume|cancel` skip the step's +/// approval_conditions and accept any resume_id. So only a workspace admin, or a job token +/// whose run lineage claims the suspended job ([`unclaimable_run_lineage`]), may mint one, and +/// a scoped token also needs run scope on the flow. The job must be in `w_id`, whose key signs +/// the result. +async fn require_resume_mint_authority( + db: &DB, + w_id: &str, + authed: &ApiAuthed, + job_id: Uuid, +) -> error::Result<()> { + let flow_path = resume_target_flow_path(db, w_id, job_id).await?; + check_scopes(authed, || format!("jobs:run:flows:{}", flow_path))?; + if !unclaimable_run_lineage(db, w_id, vec![job_id], authed) + .await? + .is_empty() + { + return Err(Error::PermissionDenied(format!( + "only job {job_id}'s own run or a workspace admin can sign resume or approval urls for it" + ))); + } + Ok(()) +} + +/// Resolve the runnable path of the flow a (possibly step) job of `w_id` belongs to, used /// to scope-check resume-signature minting against `jobs:run:flows:`. -/// Returns an empty string when the path can't be resolved (e.g. previews or an -/// unknown job); an empty path only matters for path-restricted tokens, which -/// would not be running such a flow. Never hard-fails, so it can't break resume -/// for unscoped tokens (the in-flow `get_resume_urls()` path). +/// Returns an empty string when the path can't be resolved (e.g. previews); an empty +/// path only matters for path-restricted tokens, which would not be running such a flow. async fn resume_target_flow_path(db: &DB, w_id: &str, job_id: Uuid) -> error::Result { let job = sqlx::query!( r#"SELECT kind::text as "kind!", parent_job, runnable_path @@ -6339,10 +6346,8 @@ async fn resume_target_flow_path(db: &DB, w_id: &str, job_id: Uuid) -> error::Re w_id ) .fetch_optional(db) - .await?; - let Some(job) = job else { - return Ok(String::new()); - }; + .await? + .ok_or_else(|| Error::NotFound(format!("job {job_id} not found")))?; // All flow kinds: the job itself is the flow whose path scopes the resume. if matches!( job.kind.as_str(),