Files
windmill/backend/tests/restart_flow_version_scope.rs
e7fc1b2e2e feat: restricted job tokens per script and flow (#11484)
* feat: restricted job tokens (job_token_scopes on scripts and flows)

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

* fix: admit flow-run reads, skip dedicated workers, gate on worker version

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

* fix: keep restricted jobs off flow runners, preserve scopes on rename and promotion

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

* fix: keep restricted jobs off every dedicated handoff, confine progress flow id

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

* fix: exclude restricted runnables from dedicated worker startup

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

* fix: gate restrictions on the release after 1.821.0

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

* fix: store per-job scopes on job_perms instead of v2_job, pin inline runs to the checked version

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

* feat: step-level job_token_scopes for flow steps and agent tools

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: fail closed on perms read errors, refuse restricted queue imports, gate step scopes in previews

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: carry a job's scopes on its completion so a re-run keeps the caller's cap

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: carry a zombie job's scopes into its completion

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: leave a zombie for the next sweep when its scopes cannot be read

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* docs: correct the QueuedJobV2 completion comment

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: validate step scopes in batch flows, fail closed on unvalidated step scopes

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: refuse flows with step or tool restrictions at push while an older worker is live

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: apply the step-scope worker gate to flow restarts

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: list the job token toggle with the other step and flow settings

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* chore: pin the EE companion merged with EE main

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* perf: skip scope lookups for unrestricted jobs; list job token setting last

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* style: rustfmt scopes tests

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* fix: confine restricted job tokens to their own run lineage; drop remaining extra lookups

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014FHGE3ynoAu6yrgeg4kLwL

* chore: update ee-repo-ref to 259ad3bfeef5285ba80eedc86309b11dca001220

This commit updates the EE repository reference after PR #843 was merged in windmill-ee-private.

Previous ee-repo-ref: 2b77c0225dca441235daf7bf0a06ba968df0c927

New ee-repo-ref: 259ad3bfeef5285ba80eedc86309b11dca001220

Automated by sync-ee-ref workflow.

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
2026-10-03 09:33:44 +02:00

137 lines
4.3 KiB
Rust

//! 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,
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(())
}