mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-03 00:02:08 +00:00
fix: scope flow resume to its workspace and minting to the job's run (#11392)
* fix: scope flow resume to its workspace and minting to the job's run Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * docs: name the lineage columns resume minting checks 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:
co-authored by
Claude Opus 5.5
parent
14a2619ad2
commit
163a4ffa4e
+4
-3
@@ -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"
|
||||
}
|
||||
-23
@@ -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"
|
||||
}
|
||||
@@ -513,12 +513,12 @@ async fn test_single_job_read_authorization(db: Pool<Postgres>) -> 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!(
|
||||
|
||||
@@ -201,7 +201,10 @@ async fn wac_approval_urls_mint_guards(db: Pool<Postgres>) -> 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<Postgres>) -> 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<Postgres>) -> 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<Postgres>) -> anyhow::Result<()> {
|
||||
async fn wac_concurrent_colliding_mints_admit_exactly_one(
|
||||
db: Pool<Postgres>,
|
||||
) -> 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<Postgres>) ->
|
||||
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<Postgres>,
|
||||
) -> 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::<sha2::Sha256>::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(())
|
||||
}
|
||||
|
||||
@@ -105,10 +105,33 @@ pub async fn drop_unclaimable_run_lineage(
|
||||
run_query: &mut RunJobQuery,
|
||||
authed: &ApiAuthed,
|
||||
) -> error::Result<()> {
|
||||
let mut referenced: Vec<Uuid> = [run_query.parent_job, run_query.root_job]
|
||||
let referenced: Vec<Uuid> = [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<Uuid>,
|
||||
authed: &ApiAuthed,
|
||||
) -> error::Result<Vec<Uuid>> {
|
||||
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")]
|
||||
|
||||
@@ -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<QueryApprover>,
|
||||
) -> error::Result<String> {
|
||||
// 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<QueryApprover>,
|
||||
) -> error::JsonResult<ResumeUrls> {
|
||||
// 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:<path>`.
|
||||
/// 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<String> {
|
||||
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(),
|
||||
|
||||
Reference in New Issue
Block a user