feat: restart perpetual runs on the version a deploy makes runnable (#11200)

* feat: opt-in move of perpetual runs to a newly deployed script version

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: keep the perpetual-run opt-in across relocks and check the new version's tag

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* test: pin the tag check on a perpetual version switch

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: count perpetual runs past the first queue page in the deploy prompt

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* feat: restart perpetual runs on the version a deploy makes runnable

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: claim a perpetual run and queue its replacement in one transaction

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: honour a cancel that lands after the worker last read its queue row

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* refactor: leave the lost-cancel fix to its own PR and match the scale down to 0 wording

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: run a preprocessor the deployed version adds over the arguments carried over

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* style: shorten the wording of the modal's argument warning

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* refactor: never preprocess a restarted perpetual run, as every other restart does

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: preprocess for a replacement whose run had not been through one

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: push a perpetual replacement without the deployed debounce settings

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* feat: name the runs the deploy button restarts

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* perf: skip the perpetual restart lookups on a deploy that is not perpetual

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: move perpetual runs before anything that can fail after the deploy commits

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: report cancellation on the queue listing so the deploy prompt can skip it

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: pass the tag workspace to the availability check after the merge

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* style: say which arguments are defined differently and which values are kept

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix: resolve a dynamic tag before checking it for a restarted run

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* test: pin that a deployed dynamic tag is checked as it resolves

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
hugocasa
2026-09-29 18:03:43 +02:00
committed by GitHub
co-authored by Claude Opus 5
parent bda00c6cad
commit cf5c49c3dc
11 changed files with 1056 additions and 4 deletions
@@ -0,0 +1,23 @@
{
"db_name": "PostgreSQL",
"query": "SELECT restart_unless_cancelled FROM script WHERE hash = $1 AND workspace_id = $2",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "restart_unless_cancelled",
"type_info": "Bool"
}
],
"parameters": {
"Left": [
"Int8",
"Text"
]
},
"nullable": [
true
]
},
"hash": "1d27895aa42ccbb542479b19baefd62790205b529ab0d8af36f18c470e8bb838"
}
@@ -0,0 +1,25 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE v2_job_queue SET canceled_by = $1, canceled_reason = $2, scheduled_for = now(), suspend = 0 WHERE id = $3 AND workspace_id = $4 AND canceled_by IS NULL RETURNING id",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "id",
"type_info": "Uuid"
}
],
"parameters": {
"Left": [
"Varchar",
"Text",
"Uuid",
"Text"
]
},
"nullable": [
false
]
},
"hash": "611dd43eb1b629db860c1b477444cb3dadcbf8058a24883bc22f143762c5ea3d"
}
@@ -0,0 +1,95 @@
{
"db_name": "PostgreSQL",
"query": "SELECT q.id AS \"id!\", j.created_by, j.permissioned_as, j.permissioned_as_email, j.trigger, j.trigger_kind AS \"trigger_kind: TriggerKindLabel\", j.preprocessed, j.args AS \"args: sqlx::types::Json<HashMap<String, Box<RawValue>>>\" FROM v2_job_queue q JOIN v2_job j USING (id) JOIN script s ON s.workspace_id = j.workspace_id AND s.hash = j.runnable_id WHERE j.workspace_id = $1 AND j.runnable_path = $2 AND j.kind = 'script' AND j.flow_step_id IS NULL AND j.runnable_id != $3 AND q.canceled_by IS NULL AND s.restart_unless_cancelled",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "id!",
"type_info": "Uuid"
},
{
"ordinal": 1,
"name": "created_by",
"type_info": "Varchar"
},
{
"ordinal": 2,
"name": "permissioned_as",
"type_info": "Varchar"
},
{
"ordinal": 3,
"name": "permissioned_as_email",
"type_info": "Varchar"
},
{
"ordinal": 4,
"name": "trigger",
"type_info": "Varchar"
},
{
"ordinal": 5,
"name": "trigger_kind: TriggerKindLabel",
"type_info": {
"Custom": {
"name": "job_trigger_kind",
"kind": {
"Enum": [
"webhook",
"http",
"websocket",
"kafka",
"email",
"nats",
"schedule",
"app",
"ui",
"postgres",
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google",
"ci_test",
"github",
"azure",
"asset",
"freshness",
"amqp"
]
}
}
}
},
{
"ordinal": 6,
"name": "preprocessed",
"type_info": "Bool"
},
{
"ordinal": 7,
"name": "args: sqlx::types::Json<HashMap<String, Box<RawValue>>>",
"type_info": "Jsonb"
}
],
"parameters": {
"Left": [
"Text",
"Text",
"Int8"
]
},
"nullable": [
false,
false,
false,
false,
true,
true,
true,
true
]
},
"hash": "df8586283178b0e684b9a0e359efb16e38fd68a5fe2b01cd77f74a725bbf37a3"
}
+300
View File
@@ -0,0 +1,300 @@
//! A deploy stops the perpetual runs at its path and starts them again on the version it made
//! runnable, unless that version no longer loops.
use serde_json::{json, Value};
use sqlx::{Pool, Postgres};
use uuid::Uuid;
use windmill_queue::restart_perpetual_runs_on_new_version;
const W_ID: &str = "test-workspace";
// Paths are distinct across tests: a path resolves to its deployed version through a
// process-wide cache, while each test runs against its own database.
async fn insert_version(
db: &Pool<Postgres>,
path: &str,
hash: i64,
age_s: f64,
perpetual: bool,
) -> anyhow::Result<()> {
insert_tagged_version(db, path, hash, age_s, perpetual, None).await
}
async fn insert_tagged_version(
db: &Pool<Postgres>,
path: &str,
hash: i64,
age_s: f64,
perpetual: bool,
tag: Option<&str>,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO script (workspace_id, hash, path, summary, description, content, created_by, \
language, lock, restart_unless_cancelled, tag, created_at) \
VALUES ($1, $2, $3, '', '', 'echo', 'test-user', 'bash', '', $4, $5, \
now() - make_interval(secs => $6))",
)
.bind(W_ID)
.bind(hash)
.bind(path)
.bind(perpetual)
.bind(tag)
.bind(age_s)
.execute(db)
.await?;
Ok(())
}
/// A run of `hash` a worker has started.
async fn start_run(db: &Pool<Postgres>, path: &str, hash: i64) -> anyhow::Result<Uuid> {
start_run_as(db, path, hash, "test-user", "test@windmill.dev").await
}
/// `username` and `email` name one identity, as they do on a real run, so what the replacement
/// inherits can be told apart from what the deployed version names.
async fn start_run_as(
db: &Pool<Postgres>,
path: &str,
hash: i64,
username: &str,
email: &str,
) -> anyhow::Result<Uuid> {
let id = Uuid::new_v4();
sqlx::query(
"INSERT INTO v2_job (id, workspace_id, created_by, created_at, permissioned_as, \
permissioned_as_email, kind, runnable_id, runnable_path, script_lang, tag, args, \
visible_to_owner) \
VALUES ($1, $2, $6, now(), 'u/' || $6, $5, 'script', $3, \
$4, 'bash', 'bash', '{\"n\": 1}', true)",
)
.bind(id)
.bind(W_ID)
.bind(hash)
.bind(path)
.bind(email)
.bind(username)
.execute(db)
.await?;
sqlx::query(
"INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, started_at, tag) \
VALUES ($1, $2, now(), true, now(), 'bash')",
)
.bind(id)
.bind(W_ID)
.execute(db)
.await?;
Ok(id)
}
#[sqlx::test(fixtures("base"))]
async fn a_deploy_restarts_the_runs_of_earlier_versions(db: Pool<Postgres>) -> anyhow::Result<()> {
let path = "u/test-user/restarted";
insert_version(&db, path, 101, 60.0, true).await?;
let running = start_run(&db, path, 101).await?;
insert_version(&db, path, 102, 0.0, true).await?;
restart_perpetual_runs_on_new_version(&db, W_ID, path, "test-user").await;
let canceled: Option<String> =
sqlx::query_scalar("SELECT canceled_reason FROM v2_job_queue WHERE id = $1")
.bind(running)
.fetch_one(&db)
.await?;
assert!(
canceled.is_some_and(|reason| reason.contains(path)),
"the run of the earlier version is canceled"
);
// Joined on the queue: a row in `v2_job` that never reached it would run nothing.
let (hash, args): (i64, Value) = sqlx::query_as(
"SELECT j.runnable_id, j.args FROM v2_job j JOIN v2_job_queue q USING (id) \
WHERE j.workspace_id = $1 AND j.runnable_path = $2 AND j.id <> $3 \
AND q.canceled_by IS NULL",
)
.bind(W_ID)
.bind(path)
.bind(running)
.fetch_one(&db)
.await?;
assert_eq!(hash, 102, "the next run is queued on the deployed version");
assert_eq!(
args,
json!({ "n": 1 }),
"with the arguments of the run it replaces"
);
Ok(())
}
/// A run carries what started it until its own completion swaps in what the preprocessor returned,
/// so a replacement pushed without one would hand `main` the raw arguments, and every iteration
/// after it the same.
#[sqlx::test(fixtures("base"))]
async fn a_replacement_preprocesses_arguments_the_replaced_run_had_not(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
let path = "u/test-user/preprocessed";
insert_version(&db, path, 501, 60.0, true).await?;
let running = start_run(&db, path, 501).await?;
insert_version(&db, path, 502, 0.0, true).await?;
sqlx::query(
"UPDATE script SET has_preprocessor = true WHERE hash = ANY('{501,502}') AND workspace_id = $1",
)
.bind(W_ID)
.execute(&db)
.await?;
sqlx::query("UPDATE v2_job SET preprocessed = false WHERE id = $1")
.bind(running)
.execute(&db)
.await?;
restart_perpetual_runs_on_new_version(&db, W_ID, path, "test-user").await;
let preprocessed: Option<bool> = sqlx::query_scalar(
"SELECT preprocessed FROM v2_job \
WHERE workspace_id = $1 AND runnable_path = $2 AND id <> $3",
)
.bind(W_ID)
.bind(path)
.bind(running)
.fetch_one(&db)
.await?;
assert_eq!(
preprocessed,
Some(false),
"the replacement is queued to go through the preprocessor"
);
Ok(())
}
/// A version that names an identity runs as it, like any run of that version by path, rather than
/// as whoever started the loop.
#[sqlx::test(fixtures("base"))]
async fn a_replacement_runs_as_the_identity_the_deployed_version_names(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
let path = "u/test-user/on-behalf-of";
insert_version(&db, path, 601, 60.0, true).await?;
let running = start_run_as(&db, path, 601, "test-user-2", "test2@windmill.dev").await?;
insert_version(&db, path, 602, 0.0, true).await?;
sqlx::query(
"UPDATE script SET on_behalf_of = 'u/test-user' WHERE hash = 602 AND workspace_id = $1",
)
.bind(W_ID)
.execute(&db)
.await?;
restart_perpetual_runs_on_new_version(&db, W_ID, path, "test-user").await;
let (permissioned_as, email): (String, String) = sqlx::query_as(
"SELECT permissioned_as, permissioned_as_email FROM v2_job \
WHERE workspace_id = $1 AND runnable_path = $2 AND id <> $3",
)
.bind(W_ID)
.bind(path)
.bind(running)
.fetch_one(&db)
.await?;
assert_eq!(permissioned_as, "u/test-user");
assert_eq!(
email, "test@windmill.dev",
"resolved from the identity the version names, not the one the replaced run had"
);
Ok(())
}
#[sqlx::test(fixtures("base"))]
async fn a_run_stays_on_its_version_when_it_may_not_use_the_deployed_tag(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
let path = "u/test-user/tagged";
insert_version(&db, path, 301, 60.0, true).await?;
// Runs as the fixture's user who is no superadmin, and no workspace tag allows this one.
let running = start_run_as(&db, path, 301, "test-user-2", "test2@windmill.dev").await?;
insert_tagged_version(&db, path, 302, 0.0, true, Some("restricted")).await?;
restart_perpetual_runs_on_new_version(&db, W_ID, path, "test-user").await;
let queued: Vec<(Uuid, Option<String>)> = sqlx::query_as(
"SELECT q.id, q.canceled_reason FROM v2_job_queue q JOIN v2_job j USING (id) \
WHERE j.workspace_id = $1 AND j.runnable_path = $2",
)
.bind(W_ID)
.bind(path)
.fetch_all(&db)
.await?;
assert_eq!(
queued,
vec![(running, None)],
"a tag the run's identity may not use leaves it as it is"
);
Ok(())
}
/// A `$args[...]` tag names a worker group only once the run's arguments fill it in, so it is
/// checked the way a push checks one rather than as the literal the version carries.
#[sqlx::test(fixtures("base"))]
async fn a_deployed_dynamic_tag_is_checked_against_what_it_resolves_to(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
use windmill_common::worker::{CustomTags, CUSTOM_TAGS_PER_WORKSPACE};
let path = "u/test-user/dynamic-tag";
insert_version(&db, path, 701, 60.0, true).await?;
// Runs as the fixture's user who is no superadmin, so only an allowed tag moves it.
let running = start_run_as(&db, path, 701, "test-user-2", "test2@windmill.dev").await?;
sqlx::query("UPDATE v2_job SET args = '{\"region\": \"eu\"}'::jsonb WHERE id = $1")
.bind(running)
.execute(&db)
.await?;
insert_tagged_version(&db, path, 702, 0.0, true, Some("$args[region]")).await?;
CUSTOM_TAGS_PER_WORKSPACE.store(std::sync::Arc::new(CustomTags::from(
vec!["eu".to_string()],
)));
restart_perpetual_runs_on_new_version(&db, W_ID, path, "test-user").await;
let replacement: Option<String> = sqlx::query_scalar(
"SELECT j.tag FROM v2_job j JOIN v2_job_queue q USING (id) \
WHERE j.workspace_id = $1 AND j.runnable_path = $2 AND j.id <> $3 \
AND q.canceled_by IS NULL",
)
.bind(W_ID)
.bind(path)
.bind(running)
.fetch_optional(&db)
.await?;
CUSTOM_TAGS_PER_WORKSPACE.store(std::sync::Arc::new(CustomTags::default()));
assert_eq!(
replacement.as_deref(),
Some("eu"),
"the run moves, on the tag its arguments resolve to"
);
Ok(())
}
#[sqlx::test(fixtures("base"))]
async fn a_deploy_that_stops_looping_leaves_the_runs_alone(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
let path = "u/test-user/untouched";
insert_version(&db, path, 201, 60.0, true).await?;
let running = start_run(&db, path, 201).await?;
insert_version(&db, path, 202, 0.0, false).await?;
restart_perpetual_runs_on_new_version(&db, W_ID, path, "test-user").await;
let queued: Vec<(Uuid, Option<String>)> = sqlx::query_as(
"SELECT q.id, q.canceled_reason FROM v2_job_queue q JOIN v2_job j USING (id) \
WHERE j.workspace_id = $1 AND j.runnable_path = $2",
)
.bind(W_ID)
.bind(path)
.fetch_all(&db)
.await?;
assert_eq!(
queued,
vec![(running, None)],
"turning perpetual off leaves the run as it is"
);
Ok(())
}
@@ -551,6 +551,7 @@ async fn create_snapshot_script(
let mut handle_deployment_metadata = None;
let mut moved_native_triggers = Vec::new();
let mut deployed_path = None;
let mut deployed_perpetual = false;
while let Some(field) = multipart.next_field().await.unwrap() {
let name = field.name().unwrap().to_string();
let data = field.bytes().await.unwrap();
@@ -559,6 +560,7 @@ async fn create_snapshot_script(
let is_tar = ns.codebase.as_ref().is_some_and(|x| x.ends_with(".tar"));
let use_esm = ns.codebase.as_ref().is_some_and(|x| x.contains(".esm"));
deployed_path = Some(ns.path.clone());
deployed_perpetual = ns.restart_unless_cancelled == Some(true);
let (new_hash, ntx, hdm, moved) = create_script_internal(
ns,
w_id.clone(),
@@ -616,6 +618,19 @@ async fn create_snapshot_script(
}
reregister_moved_native_triggers(&db, &authed, &w_id, moved_native_triggers);
if let Some(hdm) = handle_deployment_metadata {
let runnable_now = matches!(hdm, PostCommitDeploy::Full { .. });
if let Some(script_path) = deployed_path
.as_deref()
.filter(|_| runnable_now && deployed_perpetual)
{
windmill_queue::restart_perpetual_runs_on_new_version(
&db,
&w_id,
script_path,
&authed.username,
)
.await;
}
hdm.handle(&db).await?;
}
return Ok((StatusCode::CREATED, format!("{}", script_hash.unwrap())));
@@ -734,6 +749,9 @@ async fn deploy_script(
return Err(Error::PermissionDenied(msg));
}
let script_path = ns.path.clone();
// Only a perpetual deploy can have runs to move, so every other one skips the lookups that
// would find that out.
let perpetual = ns.restart_unless_cancelled == Some(true);
let email = authed.email.clone();
let username = authed.username.clone();
let authed_for_triggers = authed.clone();
@@ -757,6 +775,19 @@ async fn deploy_script(
// they don't run against a version whose lock does not exist yet — and
// don't run twice.
let ready_to_test = matches!(hdm, PostCommitDeploy::Full { .. });
// The version is runnable, so the perpetual runs of earlier ones move to it here, before
// anything that can fail this deploy after its commit: a version nothing moved to would
// leave those runs on the old code with nothing left to notice. A deploy that needed lock
// generation hands this to its dependency job instead.
if ready_to_test && perpetual {
windmill_queue::restart_perpetual_runs_on_new_version(
&db,
&w_id,
&script_path,
&username,
)
.await;
}
hdm.handle(&db).await?;
let db2 = db.clone();
if ready_to_test {
@@ -2864,6 +2895,12 @@ async fn create_script_internal<'c>(
if let Some(dm) = ns.deployment_message {
args.insert("deployment_message".to_string(), to_raw_value(&dm));
}
// The version becomes runnable when this job writes its lock, which is where the
// perpetual runs of earlier versions can move to it. Only a deploy someone made carries
// this, so a relock triggered by an imported script changing leaves those runs alone.
if ns.restart_unless_cancelled.is_some_and(|x| x) {
args.insert("restart_perpetual_runs".to_string(), to_raw_value(&true));
}
if let Some(ref p_path) = p_path_opt {
args.insert("parent_path".to_string(), to_raw_value(&p_path));
}
+4
View File
@@ -4228,6 +4228,7 @@ async fn get_started_at_by_ids(
struct ListableQueuedJob {
pub id: Uuid,
pub running: bool,
pub canceled: bool,
pub created_by: String,
pub created_at: chrono::DateTime<chrono::Utc>,
pub started_at: Option<chrono::DateTime<chrono::Utc>>,
@@ -4273,6 +4274,9 @@ async fn list_queue_jobs(
&[
"v2_job.id",
"v2_job_queue.running",
// A canceled row stays in the queue until a worker picks it up and completes it, and
// the `QueuedJob` schema this answers with declares the field either way.
"v2_job_queue.canceled_by IS NOT NULL as canceled",
"v2_job.created_by",
"v2_job.created_at",
"v2_job_queue.started_at",
+287 -2
View File
@@ -41,7 +41,9 @@ use windmill_common::audit::AuditAuthor;
use windmill_common::auth::JobPerms;
#[cfg(feature = "benchmark")]
use windmill_common::bench::BenchmarkIter;
use windmill_common::jobs::{JobTriggerKind, TriggerKindLabel, EMAIL_ERROR_HANDLER_USER_EMAIL};
use windmill_common::jobs::{
script_path_to_payload, JobTriggerKind, TriggerKindLabel, EMAIL_ERROR_HANDLER_USER_EMAIL,
};
use windmill_common::min_version::{
MIN_VERSION_SUPPORTS_DEBOUNCING, MIN_VERSION_SUPPORTS_DEBOUNCING_V2,
};
@@ -598,8 +600,13 @@ async fn cancel_persistent_script_jobs_internal<'c>(
let mut tx = db.begin().await?;
// we could have retrieved the job IDs in the first query where we retrieve the hashes, but just in case a job was inserted in the queue right in-between the two above query, we re-do the fetch here
// Only the loops: a dependency job of this script shares its path, and a run of a version that
// does not restart itself ends on its own.
let jobs_to_cancel = sqlx::query_scalar::<_, Uuid>(
"SELECT j.id FROM v2_job_queue q JOIN v2_job j USING (id) WHERE j.workspace_id = $1 AND j.runnable_path = $2 AND q.canceled_by IS NULL",
"SELECT j.id FROM v2_job_queue q JOIN v2_job j USING (id) \
JOIN script s ON s.workspace_id = j.workspace_id AND s.hash = j.runnable_id \
WHERE j.workspace_id = $1 AND j.runnable_path = $2 AND j.kind = 'script' \
AND j.flow_step_id IS NULL AND q.canceled_by IS NULL AND s.restart_unless_cancelled",
)
.bind(w_id)
.bind(script_path)
@@ -626,6 +633,284 @@ async fn cancel_persistent_script_jobs_internal<'c>(
return Ok(jobs_to_cancel);
}
/// Moves the perpetual runs at `script_path` to the version a deploy just made runnable: each one
/// is canceled and pushed again on that version with the arguments it ran with. A deploy that
/// leaves the script non-perpetual moves nothing, so turning perpetual off keeps the runs going as
/// it does today.
///
/// Carries the authority of the deploy, which every caller has authorized: it cancels and pushes
/// runs at `script_path` without an `Authed` of its own. A replacement runs as the deployed
/// version's identity when it names one and otherwise as the identity of the run it replaces, and
/// a run only moves to a tag that identity may use.
///
/// Errors are logged, never returned: a deploy stands whatever happens to the runs of its earlier
/// versions.
pub async fn restart_perpetual_runs_on_new_version(
db: &Pool<Postgres>,
w_id: &str,
script_path: &str,
deployed_by: &str,
) {
let loops = match restart_perpetual_runs_at_path(db, w_id, script_path, deployed_by).await {
Ok(loops) => loops,
Err(e) => {
tracing::error!(
"Could not restart the perpetual runs of {script_path} on the deployed version: {e:#}"
);
// The second pass is the retry.
true
}
};
if !loops {
return;
}
// A run that ended just before its cancel restarts itself on its own version, and only a
// cancel this won is replaced, so that run is still on the earlier version. It is queued
// again within the 10s a perpetual restart is throttled to, which this second pass then
// catches.
let (db, w_id, script_path, deployed_by) = (
db.clone(),
w_id.to_string(),
script_path.to_string(),
deployed_by.to_string(),
);
tokio::spawn(async move {
sleep(std::time::Duration::from_secs(5)).await;
if let Err(e) = restart_perpetual_runs_at_path(&db, &w_id, &script_path, &deployed_by).await
{
tracing::error!(
"Could not restart the perpetual runs of {script_path} on the deployed version: {e:#}"
);
}
});
}
/// Whether the deployed version loops, which is what the second pass is for.
async fn restart_perpetual_runs_at_path(
db: &Pool<Postgres>,
w_id: &str,
script_path: &str,
deployed_by: &str,
) -> error::Result<bool> {
// Built the way a run of this path is built anywhere else, so the next run takes the deployed
// version's tag, timeout, language and identity. Whether its preprocessor runs is decided per
// run, below.
let (mut payload, tag, _, _, timeout, on_behalf_of) =
script_path_to_payload(script_path, None, db.clone(), w_id, None).await?;
// A replacement continues a loop rather than answering a trigger, and every perpetual restart
// is pushed without debouncing for that reason. Deployed settings here would debounce the
// loops at this path against each other and collapse those that share arguments into one.
if let JobPayload::ScriptHash { debouncing_settings, .. } = &mut payload {
*debouncing_settings = DebouncingSettings::default();
}
let JobPayload::ScriptHash { hash, dedicated_worker, .. } = &payload else {
return Ok(false);
};
let (hash, dedicated_worker) = (*hash, *dedicated_worker);
let perpetual = sqlx::query_scalar!(
"SELECT restart_unless_cancelled FROM script WHERE hash = $1 AND workspace_id = $2",
hash.0,
w_id
)
.fetch_optional(db)
.await?
.flatten()
.unwrap_or(false);
if !perpetual {
return Ok(false);
}
let runs = sqlx::query_as!(
PerpetualRunToRestart,
"SELECT q.id AS \"id!\", j.created_by, j.permissioned_as, j.permissioned_as_email, \
j.trigger, j.trigger_kind AS \"trigger_kind: TriggerKindLabel\", j.preprocessed, \
j.args AS \"args: sqlx::types::Json<HashMap<String, Box<RawValue>>>\" \
FROM v2_job_queue q JOIN v2_job j USING (id) \
JOIN script s ON s.workspace_id = j.workspace_id AND s.hash = j.runnable_id \
WHERE j.workspace_id = $1 AND j.runnable_path = $2 AND j.kind = 'script' \
AND j.flow_step_id IS NULL AND j.runnable_id != $3 AND q.canceled_by IS NULL \
AND s.restart_unless_cancelled",
w_id,
script_path,
hash.0
)
.fetch_all(db)
.await?;
for run in runs {
let id = run.id;
// Per run, so that a run this fails on leaves the others to move.
if let Err(e) = restart_perpetual_run(
db,
w_id,
script_path,
deployed_by,
RestartOnVersion {
payload: &payload,
tag: tag.as_deref(),
timeout,
dedicated_worker,
on_behalf_of: on_behalf_of.as_ref(),
},
run,
)
.await
{
tracing::error!(
"Could not restart perpetual run {id} on the version deployed at {script_path}: {e:#}"
);
}
}
Ok(true)
}
/// A run of an earlier version at the path, and what its replacement inherits from it.
struct PerpetualRunToRestart {
id: Uuid,
created_by: String,
permissioned_as: String,
permissioned_as_email: String,
trigger: Option<String>,
trigger_kind: Option<TriggerKindLabel>,
/// `Some(false)` while the run still carries the arguments it was started with: only its own
/// completion swaps in what a preprocessor returned.
preprocessed: Option<bool>,
args: Option<sqlx::types::Json<HashMap<String, Box<RawValue>>>>,
}
/// What every run at the path moves to.
struct RestartOnVersion<'a> {
payload: &'a JobPayload,
tag: Option<&'a str>,
timeout: Option<i32>,
dedicated_worker: Option<bool>,
on_behalf_of: Option<&'a windmill_common::jobs::OnBehalfOf>,
}
async fn restart_perpetual_run(
db: &Pool<Postgres>,
w_id: &str,
script_path: &str,
deployed_by: &str,
version: RestartOnVersion<'_>,
run: PerpetualRunToRestart,
) -> error::Result<()> {
let RestartOnVersion { payload, tag, timeout, dedicated_worker, on_behalf_of } = version;
let (email, permissioned_as) = match on_behalf_of {
Some(obo) => (obo.email.clone(), obo.permissioned_as.clone()),
None => (
run.permissioned_as_email.clone(),
run.permissioned_as.clone(),
),
};
let args = run.args.clone().map(|args| args.0).unwrap_or_default();
// The run's own tag was checked when the loop started; the deployed version's has not
// been checked against the identity that would run it. Checked the way a push checks one,
// so a `$args[...]` tag resolves from the arguments this run carries. A dedicated worker's
// tag is the script's own and names no worker group to gain access to.
if dedicated_worker != Some(true) {
if let Some(tag) = tag.filter(|tag| !tag.is_empty()) {
let is_super_admin = windmill_common::auth::is_super_admin_email(db, &email).await?;
if let Err(e) = check_tag_available_for_push(
db,
w_id,
tag,
&PushArgs::from(&args),
is_super_admin,
None,
)
.await
{
tracing::warn!(
"Perpetual run {} stays on its version: the deployed version of \
{script_path} has tag {tag}: {e}",
run.id
);
return Ok(());
}
}
}
// A run whose own preprocessor has not run yet carries what started it, so the replacement has
// to run one: pushed without, those arguments reach `main` and every iteration after it.
let mut payload = payload.clone();
if let JobPayload::ScriptHash { apply_preprocessor, .. } = &mut payload {
*apply_preprocessor = *apply_preprocessor && run.preprocessed == Some(false);
}
let mut tx = db.begin().await?;
// Claiming the run and queueing its replacement in one transaction: a push that fails
// leaves the run looping on its own version rather than canceled with nothing to follow
// it, and a concurrent deploy cannot claim a run this one already has. A worker completes
// a run canceled this way when it next pulls it.
let claimed = sqlx::query_scalar!(
"UPDATE v2_job_queue SET canceled_by = $1, canceled_reason = $2, scheduled_for = now(), \
suspend = 0 WHERE id = $3 AND workspace_id = $4 AND canceled_by IS NULL RETURNING id",
deployed_by,
format!("a new version of {script_path} was deployed"),
run.id,
w_id
)
.fetch_optional(&mut *tx)
.await?;
if claimed.is_none() {
// It ended or was canceled since the scan. Its own restart, if it had one, is a run of
// an earlier version the next pass picks up.
return Ok(());
}
let (_, tx) = push(
db,
PushIsolationLevel::Transaction(tx),
w_id,
payload,
PushArgs::from(&args),
&run.created_by,
&email,
permissioned_as,
Some(&format!("deploy.restart.{}", run.id)),
None,
None,
schedule_path(&run.trigger_kind, &run.trigger),
None,
None,
None,
None,
false,
false,
None,
true,
tag.map(str::to_string),
timeout,
None,
None,
None,
false,
None,
None,
None,
)
.await?;
tx.commit().await?;
// Now that the replacement is queued: the children the run left behind, and, for a run no
// worker would pull, its completion.
match cancel_job(
deployed_by,
Some(format!("a new version of {script_path} was deployed")),
run.id,
w_id,
db.begin().await?,
db,
false,
false,
)
.await
{
Ok((tx, _)) => tx.commit().await?,
Err(e) => {
tracing::error!("Could not finish canceling perpetual run {}: {e:#}", run.id)
}
}
Ok(())
}
#[derive(Serialize, Debug)]
pub struct WrappedError {
pub error: serde_json::Value,
@@ -589,6 +589,16 @@ pub async fn handle_dependency_job(
// hash whose content cache has not caught up.
windmill_common::invalidate_deployed_script_hash_cache(w_id, script_path);
if restart_perpetual_runs_from_args(job.args.as_ref()) {
windmill_queue::restart_perpetual_runs_on_new_version(
db,
w_id,
script_path,
&job.created_by,
)
.await;
}
if let Err(e) = handle_deployment_metadata(
&job.permissioned_as_email,
&job.created_by,
@@ -1319,6 +1329,14 @@ pub async fn handle_flow_dependency_job(
})))
}
/// Set by a deploy someone made of a perpetual script, and by nothing else: a relock this path
/// gets because an imported script changed leaves the runs of the version it replaces alone.
fn restart_perpetual_runs_from_args(args: Option<&Json<HashMap<String, Box<RawValue>>>>) -> bool {
args.and_then(|args| args.0.get("restart_perpetual_runs"))
.and_then(|value| serde_json::from_str::<bool>(value.get()).ok())
.unwrap_or(false)
}
fn get_deployment_msg_and_parent_path_from_args(
args: Option<Json<HashMap<String, Box<RawValue>>>>,
) -> (Option<String>, Option<String>) {
@@ -108,6 +108,12 @@
import DeployButton from './DeployButton.svelte'
import { type Trigger, deployTriggers, handleSelectTriggerFromKind } from './triggers/utils'
import DraftChangesConfirmationModal from './common/confirmationModal/DraftChangesConfirmationModal.svelte'
import PerpetualRunsDeployModal from './scripts/PerpetualRunsDeployModal.svelte'
import {
loadPerpetualRunsAtPath,
stopPerpetualRuns,
type PerpetualRunsAtPath
} from './scripts/perpetualRuns'
import { Triggers } from './triggers/triggers.svelte'
import type { ScriptBuilderProps } from './script_builder'
import WorkerTagSelect from './WorkerTagSelect.svelte'
@@ -287,6 +293,19 @@
let draftTriggersModalOpen = $state(false)
let confirmDeploymentCallback: (triggersToDeploy: Trigger[]) => void = () => {}
let perpetualRunsToConfirm: PerpetualRunsAtPath | undefined = $state(undefined)
let confirmPerpetualRunsCallback: () => void = () => {}
async function stopPerpetualRunsFromModal(): Promise<boolean> {
try {
await stopPerpetualRuns(opWorkspace!, initialPath)
return true
} catch (error) {
sendUserToast(`Could not stop the runs of this script: ${error.body ?? error.message}`, true)
return false
}
}
async function handleDraftTriggersConfirmed(event: CustomEvent<{ selectedTriggers: Trigger[] }>) {
const { selectedTriggers } = event.detail
// Continue with saving the flow
@@ -658,7 +677,8 @@
stay: boolean,
parentHash: string,
deploymentMsg?: string,
triggersToDeploy?: Trigger[]
triggersToDeploy?: Trigger[],
perpetualRunsConfirmed?: boolean
): Promise<void> {
if (!triggersToDeploy) {
// Check if there are draft triggers that need confirmation
@@ -666,12 +686,39 @@
if (draftTriggers.length > 0) {
draftTriggersModalOpen = true
confirmDeploymentCallback = async (triggersToDeploy: Trigger[]) => {
await editScript(stay, parentHash, deploymentMsg, triggersToDeploy)
await editScript(
stay,
parentHash,
deploymentMsg,
triggersToDeploy,
perpetualRunsConfirmed
)
}
return
}
}
// Runs are restarted on a newer version at their own path, so a deploy that renames the
// script leaves them running the version they have.
if (
!perpetualRunsConfirmed &&
script.restart_unless_cancelled &&
initialPath &&
script.path === initialPath
) {
loadingSave = true
const runs = await loadPerpetualRunsAtPath(opWorkspace!, initialPath, script.schema)
loadingSave = false
if (runs) {
confirmPerpetualRunsCallback = () => {
perpetualRunsToConfirm = undefined
editScript(stay, parentHash, deploymentMsg, triggersToDeploy, true)
}
perpetualRunsToConfirm = runs
return
}
}
loadingSave = true
try {
// Legacy drafts can carry `schema: {}` (no `properties`), which trips
@@ -1294,6 +1341,13 @@
on:confirmed={handleDraftTriggersConfirmed}
/>
<PerpetualRunsDeployModal
runs={perpetualRunsToConfirm}
onStop={stopPerpetualRunsFromModal}
onConfirmed={() => confirmPerpetualRunsCallback()}
onCanceled={() => (perpetualRunsToConfirm = undefined)}
/>
{#if !actingUser?.operator}
<Drawer
placement="right"
@@ -0,0 +1,116 @@
<script lang="ts">
import ConfirmationModal from '$lib/components/common/confirmationModal/ConfirmationModal.svelte'
import { Alert, Button } from '$lib/components/common'
import type { PerpetualRunsAtPath } from './perpetualRuns'
interface Props {
/** Open while set. */
runs: PerpetualRunsAtPath | undefined
/** Stops the runs without deploying, so the deploy that follows starts nothing in their place. */
onStop: () => Promise<boolean>
onConfirmed: () => void
onCanceled: () => void
}
let { runs, onStop, onConfirmed, onCanceled }: Props = $props()
let stopping = $state(false)
let stopped = $state(false)
const single = $derived(runs?.count === 1)
const subject = $derived(single ? 'it' : 'them')
const runsText = $derived(
single ? '1 run of this script is' : `${runs?.count} runs of this script are`
)
function reset() {
stopping = false
stopped = false
}
</script>
<ConfirmationModal
open={!!runs}
title="Perpetual runs on an earlier version"
confirmationText={stopped
? 'Deploy'
: runs?.count === undefined
? 'Deploy and restart runs'
: `Deploy and restart ${runs.count} ${single ? 'run' : 'runs'}`}
type="reload"
loading={stopping}
onConfirmed={() => {
reset()
onConfirmed()
}}
onCanceled={() => {
reset()
onCanceled()
}}
>
{#if runs}
<div class="flex flex-col gap-3">
{#if stopped}
<p>
{single ? 'The run was' : 'The runs were'} scaled down. Deploying starts nothing in {single
? 'its'
: 'their'} place.
</p>
{:else}
{#if runs.count === undefined}
<p>
Runs of this script could not be listed. Any queued or running on an earlier version
stop and start again on this version, with the values they have now.
</p>
{:else}
<p>
{runsText} queued or running on an earlier version. Deploying stops {subject} and starts
{subject}
again on this version, with the values {single ? 'it has' : 'they have'} now.
</p>
{/if}
{#if runs.count === undefined || runs.mismatchedArgs.length > 0}
<Alert
type="warning"
size="xs"
title={runs.mismatchedArgs.length > 0
? 'Arguments no longer match'
: 'Their arguments could not be checked'}
>
<div class="flex flex-col items-start gap-2">
<p>
{#if runs.mismatchedArgs.length > 0}
This version defines
{runs.mismatchedArgs.length === 1 ? 'this argument' : 'these arguments'} differently:
{#each runs.mismatchedArgs as arg, i (arg)}
<code>{arg}</code>{i < runs.mismatchedArgs.length - 1 ? ', ' : '.'}
{/each}
{:else}
The versions the runs are on could not be read, so arguments this version defines
differently would still be carried over.
{/if}
{single ? 'The restarted run keeps' : 'Restarted runs keep'} the values {single
? 'it has'
: 'they have'} now. To run this version with values that fit it, scale down to 0 here
and start a run yourself once deployed.
</p>
<Button
variant="default"
unifiedSize="xs"
btnClasses="bg-surface"
disabled={stopping}
onclick={async () => {
stopping = true
stopped = await onStop()
stopping = false
}}
>
{stopping ? 'Scaling down to 0…' : 'Scale down to 0'}
</Button>
</div>
</Alert>
{/if}
{/if}
</div>
{/if}
</ConfirmationModal>
@@ -0,0 +1,95 @@
import { JobService, ScriptService, type QueuedJob, type Script } from '$lib/gen'
import { computeDiff } from '$lib/components/schema/schemaUtils.svelte'
const QUEUE_PAGE_SIZE = 1000
export type PerpetualRunsAtPath = {
/** Undefined when the runs at the path could not be listed: the deploy restarts them either way. */
count: number | undefined
/** Arguments the deployed schema removes, retypes or newly requires, compared with the runs' versions. */
mismatchedArgs: string[]
}
const RUNS_UNKNOWN: PerpetualRunsAtPath = { count: undefined, mismatchedArgs: [] }
/** By path rather than by the ids listed: a loop restarts every 10s, so the run listed when the
* modal opened is usually gone, and this endpoint makes the two passes that window needs. */
export async function stopPerpetualRuns(workspace: string, path: string): Promise<void> {
await JobService.cancelPersistentQueuedJobs({
workspace,
path,
requestBody: { reason: 'stopped before a new version was deployed' }
})
}
// Every page: the deploy switches every perpetual run at the path, so one past the first page
// still needs the prompt and its version's arguments compared.
async function listQueuedAtPath(workspace: string, path: string): Promise<QueuedJob[]> {
const jobs = new Map<string, QueuedJob>()
for (let page = 1; ; page++) {
const batch = await JobService.listQueue({
workspace,
scriptPathExact: path,
jobKinds: 'script',
perPage: QUEUE_PAGE_SIZE,
page
})
for (const job of batch) jobs.set(job.id, job)
if (batch.length < QUEUE_PAGE_SIZE) return [...jobs.values()]
}
}
export async function loadPerpetualRunsAtPath(
workspace: string,
path: string,
schema: { [key: string]: any } | undefined
): Promise<PerpetualRunsAtPath | undefined> {
// A deploy restarts every perpetual run at the path whatever this finds, so anything it cannot
// read leaves the count unknown rather than reporting none and skipping the prompt.
let queued: QueuedJob[]
try {
queued = await listQueuedAtPath(workspace, path)
} catch (error) {
console.error('Could not list the runs of this perpetual script', error)
return RUNS_UNKNOWN
}
// Only what the backend restarts: never a flow step or a run already canceled, and only a run
// of a perpetual version.
const candidates = queued.filter((job) => !job.is_flow_step && !job.canceled && job.script_hash)
const hashes = [...new Set(candidates.map((job) => job.script_hash!))]
let versions: Map<string, Script>
try {
versions = new Map(
await Promise.all(
hashes.map(
async (hash) => [hash, await ScriptService.getScriptByHash({ workspace, hash })] as const
)
)
)
} catch (error) {
console.error('Could not read the versions the runs of this script are on', error)
return RUNS_UNKNOWN
}
const runs = candidates.filter((job) => versions.get(job.script_hash!)?.restart_unless_cancelled)
if (runs.length === 0) {
return undefined
}
const mismatchedArgs = new Set<string>()
for (const hash of new Set(runs.map((job) => job.script_hash!))) {
const previous = versions.get(hash)?.schema
for (const [arg, { diff }] of Object.entries(computeDiff(schema, previous))) {
// An added argument only breaks a reused run when it is required, checked below.
if (diff !== 'same' && diff !== 'added') mismatchedArgs.add(arg)
}
const previouslyRequired: unknown[] = Array.isArray(previous?.required) ? previous.required : []
for (const arg of schema?.required ?? []) {
if (!previouslyRequired.includes(arg)) mismatchedArgs.add(arg)
}
}
return {
count: runs.length,
mismatchedArgs: [...mismatchedArgs]
}
}