From 893e64f63056ba42b0fe0bcf07c6078830d3669c Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sun, 27 Sep 2026 16:10:21 +0200 Subject: [PATCH] fix: only restart a flow on a version of its own path and workspace (#11376) * fix: only restart a flow on a version of its own path and workspace Co-Authored-By: Claude Opus 5.5 (1M context) * test: pin cross-workspace restart version rejection Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) --- ...3e413174d3ae4d639d50fc4a3344c02af3f95.json | 24 ++++ backend/tests/restart_flow_version_scope.rs | 135 ++++++++++++++++++ backend/windmill-api/src/jobs.rs | 9 +- backend/windmill-queue/src/jobs.rs | 37 ++++- 4 files changed, 198 insertions(+), 7 deletions(-) create mode 100644 backend/.sqlx/query-cd8a933c10eefc76ad5be2316c83e413174d3ae4d639d50fc4a3344c02af3f95.json create mode 100644 backend/tests/restart_flow_version_scope.rs diff --git a/backend/.sqlx/query-cd8a933c10eefc76ad5be2316c83e413174d3ae4d639d50fc4a3344c02af3f95.json b/backend/.sqlx/query-cd8a933c10eefc76ad5be2316c83e413174d3ae4d639d50fc4a3344c02af3f95.json new file mode 100644 index 0000000000..02e74bd0fb --- /dev/null +++ b/backend/.sqlx/query-cd8a933c10eefc76ad5be2316c83e413174d3ae4d639d50fc4a3344c02af3f95.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT EXISTS(SELECT 1 FROM flow_version WHERE id = $1 AND workspace_id = $2 AND path = $3) AS \"exists!\"", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "exists!", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Int8", + "Text", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "cd8a933c10eefc76ad5be2316c83e413174d3ae4d639d50fc4a3344c02af3f95" +} diff --git a/backend/tests/restart_flow_version_scope.rs b/backend/tests/restart_flow_version_scope.rs new file mode 100644 index 0000000000..82e8f8e4e7 --- /dev/null +++ b/backend/tests/restart_flow_version_scope.rs @@ -0,0 +1,135 @@ +//! A restart's `flow_version` is caller-supplied while the restarted job keeps the original +//! job's path, so only versions of that same flow may be restarted on. + +use sqlx::{Pool, Postgres}; +use windmill_common::cache; +use windmill_common::error::Error; +use windmill_common::flow_status::FlowStatus; +use windmill_common::jobs::JobPayload; +use windmill_queue::PushIsolationLevel; + +const WS: &str = "test-workspace"; + +async fn push_restart( + db: &Pool, + completed_job_id: uuid::Uuid, + flow_version: i64, +) -> Result { + let args = std::collections::HashMap::new(); + let (uuid, tx) = windmill_queue::push( + db, + PushIsolationLevel::IsolatedRoot(db.clone()), + WS, + JobPayload::RestartedFlow { + completed_job_id, + step_id: "a".into(), + branch_or_iteration_n: None, + flow_version: Some(flow_version), + branch_chosen: None, + nested: None, + }, + windmill_queue::PushArgs::from(&args), + "test-user", + "test@windmill.dev", + "u/test-user".to_string(), + None, + None, + None, + None, + None, + None, + None, + None, + false, + false, + None, + true, + None, + None, + None, + None, + None, + false, + None, + None, + None, + ) + .await?; + tx.commit().await?; + Ok(uuid) +} + +#[sqlx::test(fixtures("base", "hello"))] +async fn test_restart_rejects_flow_version_of_another_path( + db: Pool, +) -> anyhow::Result<()> { + // A completed run of f/system/hello_flow (version 1443253234253453). + let version = 1443253234253453; + let flow = cache::flow::fetch_version(&db, version).await?; + let completed_job_id = uuid::Uuid::new_v4(); + sqlx::query( + "INSERT INTO v2_job (id, workspace_id, kind, runnable_path, runnable_id, created_by, + permissioned_as, permissioned_as_email, tag) + VALUES ($1, $2, 'flow', 'f/system/hello_flow', $3, 'test-user', 'u/test-user', + 'test@windmill.dev', 'flow')", + ) + .bind(completed_job_id) + .bind(WS) + .bind(version) + .execute(&db) + .await?; + sqlx::query( + "INSERT INTO v2_job_completed (id, workspace_id, duration_ms, status, flow_status) + VALUES ($1, $2, 0, 'success', $3)", + ) + .bind(completed_job_id) + .bind(WS) + .bind(serde_json::to_value(FlowStatus::new(flow.value()))?) + .execute(&db) + .await?; + + // f/system/hello_with_nodes_flow's version, in the same workspace. + let err = push_restart(&db, completed_job_id, 1443253234253454) + .await + .expect_err("a version of another flow must be rejected"); + assert!(matches!(err, Error::BadRequest(_)), "{err:?}"); + + // A version at the same path in another workspace. + sqlx::query( + "INSERT INTO workspace (id, name, owner) VALUES ('other-workspace', 'other-workspace', 'test-user')", + ) + .execute(&db) + .await?; + sqlx::query( + "INSERT INTO flow (workspace_id, path, summary, description, value, edited_by) + SELECT 'other-workspace', path, summary, description, value, edited_by FROM flow + WHERE workspace_id = $1 AND path = 'f/system/hello_flow'", + ) + .bind(WS) + .execute(&db) + .await?; + let other_ws_version: i64 = sqlx::query_scalar( + "INSERT INTO flow_version (workspace_id, path, schema, value, created_by) + SELECT 'other-workspace', path, schema, value, created_by FROM flow_version WHERE id = $1 + RETURNING id", + ) + .bind(version) + .fetch_one(&db) + .await?; + let err = push_restart(&db, completed_job_id, other_ws_version) + .await + .expect_err("a version from another workspace must be rejected"); + assert!(matches!(err, Error::BadRequest(_)), "{err:?}"); + + // Another version of the restarted flow itself stays allowed. + let same_path_version: i64 = sqlx::query_scalar( + "INSERT INTO flow_version (workspace_id, path, schema, value, created_by) + SELECT workspace_id, path, schema, value, created_by FROM flow_version WHERE id = $1 + RETURNING id", + ) + .bind(version) + .fetch_one(&db) + .await?; + push_restart(&db, completed_job_id, same_path_version).await?; + Ok(()) +} diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index e23169e50b..ce0fedc577 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -6941,7 +6941,8 @@ async fn resolve_nested_restart( parent_step_id: &str, parent_branch_or_iteration_n: Option, nested_path: Vec, - parent_flow_version: Option, + // (flow path of the restarted job, version to restart on) + parent_flow_version: Option<(&str, i64)>, ) -> error::Result<( Option, Option>, @@ -6969,8 +6970,8 @@ async fn resolve_nested_restart( ) })?; - let flow_data = if let Some(version) = parent_flow_version { - cache::flow::fetch_version(db, version).await? + let flow_data = if let Some((flow_path, version)) = parent_flow_version { + windmill_queue::fetch_restart_flow_version(db, workspace_id, flow_path, version).await? } else { cache::job::fetch_flow(db, &row.job_kind, row.runnable_id) .or_else(|_| cache::job::fetch_preview_flow(db, &parent_job_id, row.raw_flow)) @@ -7166,7 +7167,7 @@ pub async fn restart_flow( &step_id, branch_or_iteration_n, nested_path, - flow_version, + flow_version.map(|v| (flow_path.as_str(), v)), ) .await?; diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 3041d89245..880b72cc0e 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -7689,6 +7689,34 @@ fn reuse_completed_zombie_module(module: FlowStatusModule) -> FlowStatusModule { } } +/// Loads the flow version a restart switches to. The version id is caller-supplied and the +/// restarted job keeps the original job's path, so it must be a version of that same flow in +/// that same workspace: any other id would run foreign code under the original path. +/// This only ties the version to the flow; it does not authorize the caller. `workspace_id` +/// and `flow_path` must come from the original job, which the caller is already allowed to +/// restart. +pub async fn fetch_restart_flow_version( + db: &Pool, + workspace_id: &str, + flow_path: &str, + version: i64, +) -> Result, Error> { + let belongs = sqlx::query_scalar!( + "SELECT EXISTS(SELECT 1 FROM flow_version WHERE id = $1 AND workspace_id = $2 AND path = $3) AS \"exists!\"", + version, + workspace_id, + flow_path, + ) + .fetch_one(db) + .await?; + if !belongs { + return Err(Error::BadRequest(format!( + "flow version {version} is not a version of flow {flow_path} in workspace {workspace_id}" + ))); + } + cache::flow::fetch_version(db, version).await +} + async fn restarted_flows_resolution( db: &Pool, workspace_id: &str, @@ -7750,9 +7778,12 @@ async fn restarted_flows_resolution( && row.job_kind == JobKind::Flow; let flow_data = if is_version_change { - // Fetch the new flow version - let new_version = flow_version.unwrap(); - cache::flow::fetch_version(db, new_version).await? + let flow_path = row.script_path.as_deref().ok_or_else(|| { + Error::BadRequest(format!( + "completed flow {completed_flow_id} has no path to restart a version of" + )) + })?; + fetch_restart_flow_version(db, workspace_id, flow_path, flow_version.unwrap()).await? } else { cache::job::fetch_flow(db, &row.job_kind, row.script_hash) .or_else(|_| {