mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-04 08:02:23 +00:00
fix: complete a canceled flow whose worker died between two steps (#11366)
* fix: complete a canceled flow whose worker died between two steps Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UzGsxKey5g3kycKwGmpKNF * fix: complete only the stranded canceled flow and let its parent process it Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UzGsxKey5g3kycKwGmpKNF * fix: requeue a stranded canceled flow for a worker to complete its cancel Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UzGsxKey5g3kycKwGmpKNF * fix: keep a requeued canceled flow's start time Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UzGsxKey5g3kycKwGmpKNF --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
65cba2dbb7
commit
c2d8997549
+14
@@ -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"
|
||||
}
|
||||
-14
@@ -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"
|
||||
}
|
||||
-15
@@ -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"
|
||||
}
|
||||
+3
-9
@@ -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"
|
||||
}
|
||||
+9
-3
@@ -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<str>\", 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<str>\", 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"
|
||||
}
|
||||
+107
-12
@@ -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::<FlowStatus>(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<Uuid>,
|
||||
step: i32,
|
||||
stranded: bool,
|
||||
canceled_by: Option<&str>,
|
||||
) -> anyhow::Result<Uuid> {
|
||||
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<String>, 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(())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<Postgres>,
|
||||
parent_canceled: bool,
|
||||
canceled: bool,
|
||||
) -> anyhow::Result<f64> {
|
||||
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<Postgres>, id: Uuid) -> anyhow::Result<f64> {
|
||||
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<Postgres>, 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<Postgres>,
|
||||
) -> 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(())
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user