diff --git a/backend/.sqlx/query-1d27895aa42ccbb542479b19baefd62790205b529ab0d8af36f18c470e8bb838.json b/backend/.sqlx/query-1d27895aa42ccbb542479b19baefd62790205b529ab0d8af36f18c470e8bb838.json new file mode 100644 index 0000000000..be045d9e75 --- /dev/null +++ b/backend/.sqlx/query-1d27895aa42ccbb542479b19baefd62790205b529ab0d8af36f18c470e8bb838.json @@ -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" +} diff --git a/backend/.sqlx/query-611dd43eb1b629db860c1b477444cb3dadcbf8058a24883bc22f143762c5ea3d.json b/backend/.sqlx/query-611dd43eb1b629db860c1b477444cb3dadcbf8058a24883bc22f143762c5ea3d.json new file mode 100644 index 0000000000..6e13ae6be1 --- /dev/null +++ b/backend/.sqlx/query-611dd43eb1b629db860c1b477444cb3dadcbf8058a24883bc22f143762c5ea3d.json @@ -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" +} diff --git a/backend/.sqlx/query-df8586283178b0e684b9a0e359efb16e38fd68a5fe2b01cd77f74a725bbf37a3.json b/backend/.sqlx/query-df8586283178b0e684b9a0e359efb16e38fd68a5fe2b01cd77f74a725bbf37a3.json new file mode 100644 index 0000000000..3383951c13 --- /dev/null +++ b/backend/.sqlx/query-df8586283178b0e684b9a0e359efb16e38fd68a5fe2b01cd77f74a725bbf37a3.json @@ -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>>\" 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>>", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [ + "Text", + "Text", + "Int8" + ] + }, + "nullable": [ + false, + false, + false, + false, + true, + true, + true, + true + ] + }, + "hash": "df8586283178b0e684b9a0e359efb16e38fd68a5fe2b01cd77f74a725bbf37a3" +} diff --git a/backend/tests/perpetual_deploy_restart.rs b/backend/tests/perpetual_deploy_restart.rs new file mode 100644 index 0000000000..6f9e82e5c6 --- /dev/null +++ b/backend/tests/perpetual_deploy_restart.rs @@ -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, + 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, + 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, path: &str, hash: i64) -> anyhow::Result { + 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, + path: &str, + hash: i64, + username: &str, + email: &str, +) -> anyhow::Result { + 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) -> 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 = + 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, +) -> 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 = 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, +) -> 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, +) -> 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)> = 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, +) -> 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 = 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, +) -> 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)> = 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(()) +} diff --git a/backend/windmill-api-scripts/src/scripts.rs b/backend/windmill-api-scripts/src/scripts.rs index 42617b249b..80a175dd99 100644 --- a/backend/windmill-api-scripts/src/scripts.rs +++ b/backend/windmill-api-scripts/src/scripts.rs @@ -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)); } diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index a92d0ee144..fe73e216bf 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -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, pub started_at: Option>, @@ -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", diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 06e5c6003a..990a2a220f 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -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, + 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, + w_id: &str, + script_path: &str, + deployed_by: &str, +) -> error::Result { + // 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>>\" \ + 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, + trigger_kind: Option, + /// `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, + args: Option>>>, +} + +/// What every run at the path moves to. +struct RestartOnVersion<'a> { + payload: &'a JobPayload, + tag: Option<&'a str>, + timeout: Option, + dedicated_worker: Option, + on_behalf_of: Option<&'a windmill_common::jobs::OnBehalfOf>, +} + +async fn restart_perpetual_run( + db: &Pool, + 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, diff --git a/backend/windmill-worker/src/worker_lockfiles.rs b/backend/windmill-worker/src/worker_lockfiles.rs index 03150beb00..861a9625e9 100644 --- a/backend/windmill-worker/src/worker_lockfiles.rs +++ b/backend/windmill-worker/src/worker_lockfiles.rs @@ -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>>>) -> bool { + args.and_then(|args| args.0.get("restart_perpetual_runs")) + .and_then(|value| serde_json::from_str::(value.get()).ok()) + .unwrap_or(false) +} + fn get_deployment_msg_and_parent_path_from_args( args: Option>>>, ) -> (Option, Option) { diff --git a/frontend/src/lib/components/ScriptBuilder.svelte b/frontend/src/lib/components/ScriptBuilder.svelte index 990f02de54..25942b9a6d 100644 --- a/frontend/src/lib/components/ScriptBuilder.svelte +++ b/frontend/src/lib/components/ScriptBuilder.svelte @@ -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 { + 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 { 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} /> + confirmPerpetualRunsCallback()} + onCanceled={() => (perpetualRunsToConfirm = undefined)} +/> + {#if !actingUser?.operator} + 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 + 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 + } + + + { + reset() + onConfirmed() + }} + onCanceled={() => { + reset() + onCanceled() + }} +> + {#if runs} +
+ {#if stopped} +

+ {single ? 'The run was' : 'The runs were'} scaled down. Deploying starts nothing in {single + ? 'its' + : 'their'} place. +

+ {:else} + {#if runs.count === undefined} +

+ 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. +

+ {:else} +

+ {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. +

+ {/if} + {#if runs.count === undefined || runs.mismatchedArgs.length > 0} + 0 + ? 'Arguments no longer match' + : 'Their arguments could not be checked'} + > +
+

+ {#if runs.mismatchedArgs.length > 0} + This version defines + {runs.mismatchedArgs.length === 1 ? 'this argument' : 'these arguments'} differently: + {#each runs.mismatchedArgs as arg, i (arg)} + {arg}{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. +

+ +
+
+ {/if} + {/if} +
+ {/if} +
diff --git a/frontend/src/lib/components/scripts/perpetualRuns.ts b/frontend/src/lib/components/scripts/perpetualRuns.ts new file mode 100644 index 0000000000..2230476a33 --- /dev/null +++ b/frontend/src/lib/components/scripts/perpetualRuns.ts @@ -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 { + 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 { + const jobs = new Map() + 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 { + // 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 + 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() + 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] + } +}