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

* test: pin cross-workspace restart version rejection

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:10:21 +00:00
committed by GitHub
co-authored by Claude Opus 5.5
parent 90f9e59321
commit 893e64f630
4 changed files with 198 additions and 7 deletions
@@ -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"
}
+135
View File
@@ -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<Postgres>,
completed_job_id: uuid::Uuid,
flow_version: i64,
) -> Result<uuid::Uuid, Error> {
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<Postgres>,
) -> 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(())
}
+5 -4
View File
@@ -6941,7 +6941,8 @@ async fn resolve_nested_restart(
parent_step_id: &str,
parent_branch_or_iteration_n: Option<usize>,
nested_path: Vec<NestedRestartStep>,
parent_flow_version: Option<i64>,
// (flow path of the restarted job, version to restart on)
parent_flow_version: Option<(&str, i64)>,
) -> error::Result<(
Option<windmill_common::flow_status::BranchChosen>,
Option<Box<RestartedFrom>>,
@@ -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?;
+34 -3
View File
@@ -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<Postgres>,
workspace_id: &str,
flow_path: &str,
version: i64,
) -> Result<Arc<FlowData>, 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<Postgres>,
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(|_| {