diff --git a/backend/.sqlx/query-0ad65482d8282f778f4123b4133f9f230fffd4ea1f403dd37b28d4a105c7ef58.json b/backend/.sqlx/query-0ad65482d8282f778f4123b4133f9f230fffd4ea1f403dd37b28d4a105c7ef58.json new file mode 100644 index 0000000000..a6bf4b92d5 --- /dev/null +++ b/backend/.sqlx/query-0ad65482d8282f778f4123b4133f9f230fffd4ea1f403dd37b28d4a105c7ef58.json @@ -0,0 +1,14 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE v2_job_queue SET running = false,\n started_at = CASE WHEN canceled_by IS NULL THEN NULL ELSE started_at END\n WHERE id = $1", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "0ad65482d8282f778f4123b4133f9f230fffd4ea1f403dd37b28d4a105c7ef58" +} diff --git a/backend/.sqlx/query-3e55d027327bd3c76810fbe22d3ccb1bbbf83c8cff69d8f5907d1417a2522e69.json b/backend/.sqlx/query-3e55d027327bd3c76810fbe22d3ccb1bbbf83c8cff69d8f5907d1417a2522e69.json deleted file mode 100644 index 2603a95877..0000000000 --- a/backend/.sqlx/query-3e55d027327bd3c76810fbe22d3ccb1bbbf83c8cff69d8f5907d1417a2522e69.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE v2_job_queue SET running = false, started_at = null\n WHERE id = $1 AND canceled_by IS NULL", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "3e55d027327bd3c76810fbe22d3ccb1bbbf83c8cff69d8f5907d1417a2522e69" -} diff --git a/backend/.sqlx/query-7d07a717533bfcaf581f6655bc387095542490fbb4aae30ec7fa75c2dae98ec8.json b/backend/.sqlx/query-7d07a717533bfcaf581f6655bc387095542490fbb4aae30ec7fa75c2dae98ec8.json deleted file mode 100644 index 4e90a2a0ec..0000000000 --- a/backend/.sqlx/query-7d07a717533bfcaf581f6655bc387095542490fbb4aae30ec7fa75c2dae98ec8.json +++ /dev/null @@ -1,15 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE v2_job_runtime r SET\n ping = now()\n FROM v2_job_queue q\n WHERE r.id = $1 AND q.id = r.id\n AND q.workspace_id = $2\n AND canceled_by IS NULL", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Uuid", - "Text" - ] - }, - "nullable": [] - }, - "hash": "7d07a717533bfcaf581f6655bc387095542490fbb4aae30ec7fa75c2dae98ec8" -} diff --git a/backend/.sqlx/query-700cb5b65f46da7b3e34566a4fd4ea27005eb98f5d0d41999161082152df0f40.json b/backend/.sqlx/query-975158ae72ec12f3f933fba1d9623bf1c925feebe746299cd1f0eda3f8b949d5.json similarity index 81% rename from backend/.sqlx/query-700cb5b65f46da7b3e34566a4fd4ea27005eb98f5d0d41999161082152df0f40.json rename to backend/.sqlx/query-975158ae72ec12f3f933fba1d9623bf1c925feebe746299cd1f0eda3f8b949d5.json index 1badb9cc53..0b166b32a4 100644 --- a/backend/.sqlx/query-700cb5b65f46da7b3e34566a4fd4ea27005eb98f5d0d41999161082152df0f40.json +++ b/backend/.sqlx/query-975158ae72ec12f3f933fba1d9623bf1c925feebe746299cd1f0eda3f8b949d5.json @@ -1,17 +1,12 @@ { "db_name": "PostgreSQL", - "query": "WITH deleted AS (\n DELETE FROM v2_job_queue WHERE id = $1\n RETURNING id, workspace_id, started_at, worker, canceled_by, canceled_reason\n ), completed AS (\n INSERT INTO v2_job_completed AS cj\n ( workspace_id\n , id\n , started_at\n , duration_ms\n , result\n , result_columns\n , canceled_by\n , canceled_reason\n , flow_status\n , workflow_as_code_status\n , memory_peak\n , status\n , worker\n )\n SELECT d.workspace_id, d.id, d.started_at,\n COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(d.started_at, now()))))*1000),\n $3::text::jsonb, $10,\n CASE WHEN $4::BOOL THEN $5 ELSE d.canceled_by END,\n CASE WHEN $4::BOOL THEN $6 WHEN d.canceled_by IS NOT NULL THEN d.canceled_reason END,\n s.flow_status, s.workflow_as_code_status, $8,\n CASE WHEN $4::BOOL OR d.canceled_by IS NOT NULL THEN 'canceled'::job_status\n WHEN $7::BOOL THEN 'skipped'::job_status\n WHEN $2::BOOL THEN 'success'::job_status\n ELSE 'failure'::job_status END,\n d.worker\n FROM deleted d LEFT JOIN v2_job_status s ON s.id = d.id\n ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3::text::jsonb,\n canceled_by = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL\n THEN EXCLUDED.canceled_by ELSE cj.canceled_by END,\n canceled_reason = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL\n THEN EXCLUDED.canceled_reason ELSE cj.canceled_reason END\n RETURNING duration_ms\n ), parent_ping AS (\n UPDATE v2_job_runtime r SET ping = now()\n FROM v2_job_queue q\n WHERE r.id = $11 AND q.id = r.id AND q.workspace_id = $12 AND q.canceled_by IS NULL\n AND EXISTS (SELECT 1 FROM completed)\n AND NOT EXISTS (SELECT 1 FROM deleted WHERE canceled_by IS NOT NULL)\n )\n SELECT c.duration_ms AS \"duration_ms!\",\n EXISTS (SELECT 1 FROM deleted WHERE canceled_by IS NOT NULL) AS \"carried_cancel!\"\n FROM completed c", + "query": "WITH deleted AS (\n DELETE FROM v2_job_queue WHERE id = $1\n RETURNING id, workspace_id, started_at, worker, canceled_by, canceled_reason\n ), completed AS (\n INSERT INTO v2_job_completed AS cj\n ( workspace_id\n , id\n , started_at\n , duration_ms\n , result\n , result_columns\n , canceled_by\n , canceled_reason\n , flow_status\n , workflow_as_code_status\n , memory_peak\n , status\n , worker\n )\n SELECT d.workspace_id, d.id, d.started_at,\n COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(d.started_at, now()))))*1000),\n $3::text::jsonb, $10,\n CASE WHEN $4::BOOL THEN $5 ELSE d.canceled_by END,\n CASE WHEN $4::BOOL THEN $6 WHEN d.canceled_by IS NOT NULL THEN d.canceled_reason END,\n s.flow_status, s.workflow_as_code_status, $8,\n CASE WHEN $4::BOOL OR d.canceled_by IS NOT NULL THEN 'canceled'::job_status\n WHEN $7::BOOL THEN 'skipped'::job_status\n WHEN $2::BOOL THEN 'success'::job_status\n ELSE 'failure'::job_status END,\n d.worker\n FROM deleted d LEFT JOIN v2_job_status s ON s.id = d.id\n ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3::text::jsonb,\n canceled_by = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL\n THEN EXCLUDED.canceled_by ELSE cj.canceled_by END,\n canceled_reason = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL\n THEN EXCLUDED.canceled_reason ELSE cj.canceled_reason END\n RETURNING duration_ms\n ), parent_ping AS (\n UPDATE v2_job_runtime r SET ping = now()\n FROM v2_job_queue q\n WHERE r.id = $11 AND q.id = r.id AND q.workspace_id = $12\n AND EXISTS (SELECT 1 FROM completed)\n )\n SELECT duration_ms AS \"duration_ms!\" FROM completed", "describe": { "columns": [ { "ordinal": 0, "name": "duration_ms!", "type_info": "Int8" - }, - { - "ordinal": 1, - "name": "carried_cancel!", - "type_info": "Bool" } ], "parameters": { @@ -31,9 +26,8 @@ ] }, "nullable": [ - false, - null + false ] }, - "hash": "700cb5b65f46da7b3e34566a4fd4ea27005eb98f5d0d41999161082152df0f40" + "hash": "975158ae72ec12f3f933fba1d9623bf1c925feebe746299cd1f0eda3f8b949d5" } diff --git a/backend/.sqlx/query-1bf87fe9667f860c6089bffb803291f4ff5c979be956369fa4406d789cca92fc.json b/backend/.sqlx/query-ca0ca0d93e09c0a6d75a996366902324a89374725d65e9d1f6667d13143be363.json similarity index 77% rename from backend/.sqlx/query-1bf87fe9667f860c6089bffb803291f4ff5c979be956369fa4406d789cca92fc.json rename to backend/.sqlx/query-ca0ca0d93e09c0a6d75a996366902324a89374725d65e9d1f6667d13143be363.json index 1ee774a993..8aece1b2cd 100644 --- a/backend/.sqlx/query-1bf87fe9667f860c6089bffb803291f4ff5c979be956369fa4406d789cca92fc.json +++ b/backend/.sqlx/query-ca0ca0d93e09c0a6d75a996366902324a89374725d65e9d1f6667d13143be363.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT\n j.id AS \"id!\", j.workspace_id AS \"workspace_id!\", j.parent_job, j.flow_step_id IS NOT NULL AS \"is_flow_step?\",\n COALESCE(s.flow_status, s.workflow_as_code_status)::text AS \"flow_status: Box\", r.ping AS last_ping, j.same_worker AS \"same_worker?\",\n q.worker AS \"worker?\",\n wp.ping_at AS \"worker_last_ping?\",\n wp.memory_usage AS \"worker_memory_usage?\",\n wp.wm_memory_usage AS \"worker_wm_memory_usage?\",\n wp.memory AS \"worker_memory_total?\",\n wp.worker_group AS \"worker_group?\",\n wp.wm_version AS \"worker_version?\",\n wp.current_job_id AS \"worker_current_job_id?\",\n wp.worker_instance AS \"worker_instance?\"\n FROM v2_job_queue q JOIN v2_job j USING (id) LEFT JOIN v2_job_runtime r USING (id) LEFT JOIN v2_job_status s USING (id)\n LEFT JOIN worker_ping wp ON wp.worker = q.worker\n WHERE q.running = true AND q.suspend = 0 AND q.suspend_until IS null AND q.scheduled_for <= now()\n AND (j.kind = 'flow' OR j.kind = 'flowpreview' OR j.kind = 'flownode' OR j.kind = 'singlestepflow')\n AND r.ping IS NOT NULL AND r.ping < NOW() - ($1 || ' seconds')::interval\n AND q.canceled_by IS NULL\n\n ", + "query": "\n SELECT\n j.id AS \"id!\", j.workspace_id AS \"workspace_id!\", j.parent_job, j.flow_step_id IS NOT NULL AS \"is_flow_step?\",\n COALESCE(s.flow_status, s.workflow_as_code_status)::text AS \"flow_status: Box\", r.ping AS last_ping, j.same_worker AS \"same_worker?\",\n q.worker AS \"worker?\",\n wp.ping_at AS \"worker_last_ping?\",\n wp.memory_usage AS \"worker_memory_usage?\",\n wp.wm_memory_usage AS \"worker_wm_memory_usage?\",\n wp.memory AS \"worker_memory_total?\",\n wp.worker_group AS \"worker_group?\",\n wp.wm_version AS \"worker_version?\",\n wp.current_job_id AS \"worker_current_job_id?\",\n wp.worker_instance AS \"worker_instance?\",\n q.canceled_by AS \"canceled_by?\"\n FROM v2_job_queue q JOIN v2_job j USING (id) LEFT JOIN v2_job_runtime r USING (id) LEFT JOIN v2_job_status s USING (id)\n LEFT JOIN worker_ping wp ON wp.worker = q.worker\n WHERE q.running = true AND q.suspend = 0 AND q.suspend_until IS null AND q.scheduled_for <= now()\n AND (j.kind = 'flow' OR j.kind = 'flowpreview' OR j.kind = 'flownode' OR j.kind = 'singlestepflow')\n AND r.ping IS NOT NULL AND r.ping < NOW() - ($1 || ' seconds')::interval\n\n ", "describe": { "columns": [ { @@ -82,6 +82,11 @@ "ordinal": 15, "name": "worker_instance?", "type_info": "Varchar" + }, + { + "ordinal": 16, + "name": "canceled_by?", + "type_info": "Varchar" } ], "parameters": { @@ -105,8 +110,9 @@ false, false, true, - false + false, + true ] }, - "hash": "1bf87fe9667f860c6089bffb803291f4ff5c979be956369fa4406d789cca92fc" + "hash": "ca0ca0d93e09c0a6d75a996366902324a89374725d65e9d1f6667d13143be363" } diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index f8b5ab7303..ad638c8afa 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -6516,13 +6516,13 @@ async fn handle_zombie_flows(db: &DB) -> error::Result<()> { wp.worker_group AS "worker_group?", wp.wm_version AS "worker_version?", wp.current_job_id AS "worker_current_job_id?", - wp.worker_instance AS "worker_instance?" + wp.worker_instance AS "worker_instance?", + q.canceled_by AS "canceled_by?" FROM v2_job_queue q JOIN v2_job j USING (id) LEFT JOIN v2_job_runtime r USING (id) LEFT JOIN v2_job_status s USING (id) LEFT JOIN worker_ping wp ON wp.worker = q.worker WHERE q.running = true AND q.suspend = 0 AND q.suspend_until IS null AND q.scheduled_for <= now() AND (j.kind = 'flow' OR j.kind = 'flowpreview' OR j.kind = 'flownode' OR j.kind = 'singlestepflow') AND r.ping IS NOT NULL AND r.ping < NOW() - ($1 || ' seconds')::interval - AND q.canceled_by IS NULL "#, FLOW_ZOMBIE_TRANSITION_TIMEOUT.as_str() @@ -6535,19 +6535,30 @@ async fn handle_zombie_flows(db: &DB) -> error::Result<()> { .flow_status .as_deref() .and_then(|x| serde_json::from_str::(x).ok()); - if !flow.same_worker.unwrap_or(false) - && status.as_ref().is_some_and(|s| s.is_not_yet_started()) + // A worker that pulls a canceled flow completes it as canceled, and hands that to its + // parent like any step's cancel, so a canceled flow whose transition was lost goes back + // to the queue exactly like a flow that never started. + if flow.canceled_by.is_some() + || (!flow.same_worker.unwrap_or(false) + && status.as_ref().is_some_and(|s| s.is_not_yet_started())) { - let error_message = format!( - "Zombie flow detected: {} in workspace {}. It hasn't started yet, restarting it.", - flow.id, flow.workspace_id - ); + let error_message = match flow.canceled_by.as_deref() { + Some(canceler) => format!( + "Zombie flow detected: {} in workspace {}. It was canceled by {canceler} but its worker stopped between two steps, queuing it again to complete the cancel.", + flow.id, flow.workspace_id + ), + None => format!( + "Zombie flow detected: {} in workspace {}. It hasn't started yet, restarting it.", + flow.id, flow.workspace_id + ), + }; tracing::error!(error_message); - if !CRITICAL_ALERT_MUTE_ZOMBIE_JOB_RESTART.load(Ordering::Relaxed) { + if flow.canceled_by.is_some() + || !CRITICAL_ALERT_MUTE_ZOMBIE_JOB_RESTART.load(Ordering::Relaxed) + { report_critical_error(error_message, db.clone(), Some(&flow.workspace_id), None) .await; } - // if the flow hasn't started and is a zombie, we can simply restart it let mut tx = db.begin().await?; let concurrency_key = @@ -6569,9 +6580,12 @@ async fn handle_zombie_flows(db: &DB) -> error::Result<()> { } } + // A canceled flow keeps its start: the pull only sets a missing one, and the canceled + // run's duration is measured from it. sqlx::query!( - "UPDATE v2_job_queue SET running = false, started_at = null - WHERE id = $1 AND canceled_by IS NULL", + "UPDATE v2_job_queue SET running = false, + started_at = CASE WHEN canceled_by IS NULL THEN NULL ELSE started_at END + WHERE id = $1", flow.id ) .execute(&mut *tx) @@ -7797,3 +7811,84 @@ mod log_file_listing_tests { assert_eq!(files[0].0.to_string(), "2026-08-29 06:46:00"); } } + +#[cfg(test)] +mod canceled_zombie_flow_tests { + use super::{handle_zombie_flows, DB}; + use serde_json::json; + use uuid::Uuid; + + /// A running flow at `step`; `stranded` gives it the stale ping of a lost transition. + async fn insert_flow( + db: &DB, + parent: Option, + step: i32, + stranded: bool, + canceled_by: Option<&str>, + ) -> anyhow::Result { + let id = Uuid::new_v4(); + sqlx::query( + "INSERT INTO v2_job (id, workspace_id, created_by, permissioned_as, permissioned_as_email, + kind, tag, parent_job, flow_step_id) + VALUES ($1, 'admins', 'admin', 'u/admin', 'admin@windmill.dev', 'flowpreview', 'flow', + $2, CASE WHEN $2 IS NOT NULL THEN 'sf' END)", + ) + .bind(id) + .bind(parent) + .execute(db) + .await?; + sqlx::query( + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, started_at, tag, + canceled_by, canceled_reason) + VALUES ($1, 'admins', now(), true, now() - interval '1 hour', 'flow', $2, $2)", + ) + .bind(id) + .bind(canceled_by) + .execute(db) + .await?; + sqlx::query( + "INSERT INTO v2_job_runtime (id, ping) + VALUES ($1, CASE WHEN $2 THEN now() - interval '1 hour' END)", + ) + .bind(id) + .bind(stranded) + .execute(db) + .await?; + sqlx::query("INSERT INTO v2_job_status (id, flow_status) VALUES ($1, $2)") + .bind(id) + .bind(json!({"step": step, "modules": [{"type": "InProgress", "id": "a", "job": Uuid::nil()}], + "failure_module": {"type": "WaitingForPriorSteps", "id": "failure"}})) + .execute(db) + .await?; + Ok(id) + } + + /// `(running, canceled_by, still started an hour ago)` + async fn queue_row(db: &DB, id: Uuid) -> anyhow::Result<(bool, Option, bool)> { + Ok(sqlx::query_as( + "SELECT running, canceled_by, started_at < now() - interval '30 minutes' + FROM v2_job_queue WHERE id = $1", + ) + .bind(id) + .fetch_one(db) + .await?) + } + + /// A subflow canceled on its own whose worker died between two steps goes back to the queue + /// with its cancel and its start, for a worker to complete it and hand it to its parent. The + /// parent is left alone: the cancel must not reach flows the user never canceled. + #[sqlx::test(migrations = "./migrations")] + async fn requeues_a_stranded_canceled_flow_and_nothing_else(db: DB) -> anyhow::Result<()> { + let root = insert_flow(&db, None, 0, false, None).await?; + let child = insert_flow(&db, Some(root), 1, true, Some("admin")).await?; + + handle_zombie_flows(&db).await?; + + assert_eq!( + queue_row(&db, child).await?, + (false, Some("admin".to_string()), true) + ); + assert_eq!(queue_row(&db, root).await?, (true, None, true)); + Ok(()) + } +} diff --git a/backend/tests/flow_step_completion_ping.rs b/backend/tests/flow_step_completion_ping.rs index e4e685ddac..89dab1ceef 100644 --- a/backend/tests/flow_step_completion_ping.rs +++ b/backend/tests/flow_step_completion_ping.rs @@ -1,5 +1,6 @@ //! A flow step's completion is the progress that keeps its flow from being reaped as a zombie: it -//! must refresh the parent's ping, unless the parent is being canceled. +//! must refresh the parent's ping, canceled or not, since a canceled flow still needs its next +//! transition to complete. use serde_json::json; use sqlx::{types::Json, Pool, Postgres}; @@ -46,16 +47,17 @@ async fn insert_running( async fn complete_step_and_read_parent_ping_age( db: &Pool, - parent_canceled: bool, + canceled: bool, ) -> anyhow::Result { let parent = insert_running(db, "flow", None, None).await?; - if parent_canceled { - sqlx::query("UPDATE v2_job_queue SET canceled_by = 'test-user' WHERE id = $1") - .bind(parent) + let step = insert_running(db, "script", Some(parent), Some("a")).await?; + if canceled { + // a cancel of the flow marks the flow and each of its steps + sqlx::query("UPDATE v2_job_queue SET canceled_by = 'test-user' WHERE id = ANY($1)") + .bind(vec![parent, step]) .execute(db) .await?; } - let step = insert_running(db, "script", Some(parent), Some("a")).await?; let job = get_mini_completed_job(&step, W_ID, db).await?.unwrap(); add_completed_job( db, @@ -79,7 +81,7 @@ async fn complete_step_and_read_parent_ping_age( .bind(step) .fetch_one(db) .await?; - assert_eq!(status, "success"); + assert_eq!(status, if canceled { "canceled" } else { "success" }); assert!(!queued, "the step left the queue"); Ok(sqlx::query_scalar( @@ -99,97 +101,8 @@ async fn a_step_completion_refreshes_its_running_flow_ping( let age = complete_step_and_read_parent_ping_age(&db, true).await?; assert!( - age > 3000.0, - "a canceled parent's ping is left alone: {age}s old" - ); - Ok(()) -} - -async fn ping_age(db: &Pool, id: Uuid) -> anyhow::Result { - Ok(sqlx::query_scalar( - "SELECT EXTRACT(EPOCH FROM now() - ping)::float8 FROM v2_job_runtime WHERE id = $1", - ) - .bind(id) - .fetch_one(db) - .await?) -} - -/// Waits for the completion to be stuck on the queue row of `id`. -async fn wait_until_blocked_on(db: &Pool, id: Uuid) -> anyhow::Result<()> { - for _ in 0..100 { - let blocked: bool = sqlx::query_scalar( - "SELECT EXISTS (SELECT 1 FROM pg_locks l JOIN v2_job_queue q \ - ON q.ctid = ('(' || l.page || ',' || l.tuple || ')')::tid \ - WHERE NOT l.granted AND l.locktype = 'tuple' AND q.id = $1 \ - AND l.database = (SELECT oid FROM pg_database WHERE datname = current_database())) \ - OR EXISTS (SELECT 1 FROM pg_stat_activity \ - WHERE datname = current_database() AND wait_event_type = 'Lock' \ - AND query LIKE '%DELETE FROM v2_job_queue WHERE id = $1%')", - ) - .bind(id) - .fetch_one(db) - .await?; - if blocked { - return Ok(()); - } - tokio::time::sleep(std::time::Duration::from_millis(50)).await; - } - anyhow::bail!("the completion never blocked on the step's queue row") -} - -/// A cancel of the flow marks the flow and then its steps in one transaction. A step completing -/// while that is uncommitted waits on its own row, and must then leave the canceled flow's ping -/// alone although it began before the cancel was visible. -#[sqlx::test(fixtures("base"))] -async fn a_flow_cancel_landing_during_a_step_completion_keeps_the_flow_ping( - db: Pool, -) -> anyhow::Result<()> { - let parent = insert_running(&db, "flow", None, None).await?; - let step = insert_running(&db, "script", Some(parent), Some("a")).await?; - let job = get_mini_completed_job(&step, W_ID, &db).await?.unwrap(); - - let mut cancel = db.begin().await?; - sqlx::query( - "UPDATE v2_job_queue SET canceled_by = 'test-user', canceled_reason = 'stop' \ - WHERE id = ANY($1)", - ) - .bind(vec![parent, step]) - .execute(&mut *cancel) - .await?; - - let completing = tokio::spawn({ - let db = db.clone(); - async move { - add_completed_job( - &db, - &job, - true, - false, - Json(&json!("done")), - None, - 0, - None, - false, - None, - false, - ) - .await - } - }); - wait_until_blocked_on(&db, step).await?; - cancel.commit().await?; - completing.await??; - - let status: String = - sqlx::query_scalar("SELECT status::text FROM v2_job_completed WHERE id = $1") - .bind(step) - .fetch_one(&db) - .await?; - assert_eq!(status, "canceled"); - let age = ping_age(&db, parent).await?; - assert!( - age > 3000.0, - "the canceled flow's ping is left alone: {age}s old" + age < 60.0, + "a canceled parent's ping was refreshed: {age}s old" ); Ok(()) } diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 349684c590..3041d89245 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -1572,10 +1572,10 @@ impl Completion<'_> { .await .map_err(err); }; - // A cancel of the flow marks the flow and then its steps. When the delete waited on it, - // the parent row read by this statement's snapshot predates it, so the ping moves to a - // statement of its own that sees the cancel. - let completed = sqlx::query!( + // A canceled flow is pinged too: it is completed by its next transition like any other + // flow, and the zombie flow monitor needs the ping to finish the cancel if that + // transition is lost. + sqlx::query_scalar!( "WITH deleted AS ( DELETE FROM v2_job_queue WHERE id = $1 RETURNING id, workspace_id, started_at, worker, canceled_by, canceled_reason @@ -1616,13 +1616,10 @@ impl Completion<'_> { ), parent_ping AS ( UPDATE v2_job_runtime r SET ping = now() FROM v2_job_queue q - WHERE r.id = $11 AND q.id = r.id AND q.workspace_id = $12 AND q.canceled_by IS NULL + WHERE r.id = $11 AND q.id = r.id AND q.workspace_id = $12 AND EXISTS (SELECT 1 FROM completed) - AND NOT EXISTS (SELECT 1 FROM deleted WHERE canceled_by IS NOT NULL) ) - SELECT c.duration_ms AS \"duration_ms!\", - EXISTS (SELECT 1 FROM deleted WHERE canceled_by IS NOT NULL) AS \"carried_cancel!\" - FROM completed c", + SELECT duration_ms AS \"duration_ms!\" FROM completed", /* $1 */ completed_job.id, /* $2 */ success, /* $3 */ result, @@ -1639,26 +1636,7 @@ impl Completion<'_> { .fetch_optional(&mut *conn) .warn_after_seconds(10) .await - .map_err(err)?; - let Some(completed) = completed else { - return Ok(None); - }; - if completed.carried_cancel { - sqlx::query!( - "UPDATE v2_job_runtime r SET - ping = now() - FROM v2_job_queue q - WHERE r.id = $1 AND q.id = r.id - AND q.workspace_id = $2 - AND canceled_by IS NULL", - parent_to_ping, - &completed_job.workspace_id - ) - .execute(&mut *conn) - .warn_after_seconds(10) - .await?; - } - Ok(Some(completed.duration_ms)) + .map_err(err) } }