From 2eea00a2cdc15b3ba2159c909cc5e092328ffd61 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Wed, 17 Apr 2024 21:58:20 +0200 Subject: [PATCH] fix: improve cancel_all to never deadlock --- ...e2d2c3bc3fdbde530a0a6b6fd68a22b867c3.json} | 7 +- ...feb5220ab367e339510ee0e83bb7b5abfd184.json | 23 ++++++ backend/windmill-api/src/jobs.rs | 75 +++++++++++-------- 3 files changed, 70 insertions(+), 35 deletions(-) rename backend/.sqlx/{query-18699cb0eca25b6bde05d81571dfdea8cafd0043634f61b0f652a93767c9c30a.json => query-caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3.json} (59%) create mode 100644 backend/.sqlx/query-d4fb94ee8198592c24e85d29078feb5220ab367e339510ee0e83bb7b5abfd184.json diff --git a/backend/.sqlx/query-18699cb0eca25b6bde05d81571dfdea8cafd0043634f61b0f652a93767c9c30a.json b/backend/.sqlx/query-caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3.json similarity index 59% rename from backend/.sqlx/query-18699cb0eca25b6bde05d81571dfdea8cafd0043634f61b0f652a93767c9c30a.json rename to backend/.sqlx/query-caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3.json index 6ad0bd2efd..95315372d0 100644 --- a/backend/.sqlx/query-18699cb0eca25b6bde05d81571dfdea8cafd0043634f61b0f652a93767c9c30a.json +++ b/backend/.sqlx/query-caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE queue SET canceled = true, canceled_by = $2, scheduled_for = now(), suspend = 0 WHERE scheduled_for < now() AND workspace_id = $1 AND schedule_path IS NULL RETURNING id, running, is_flow_step", + "query": "SELECT id, running, is_flow_step FROM queue WHERE scheduled_for < now() AND workspace_id = $1 AND schedule_path IS NULL", "describe": { "columns": [ { @@ -21,8 +21,7 @@ ], "parameters": { "Left": [ - "Text", - "Varchar" + "Text" ] }, "nullable": [ @@ -31,5 +30,5 @@ true ] }, - "hash": "18699cb0eca25b6bde05d81571dfdea8cafd0043634f61b0f652a93767c9c30a" + "hash": "caeb49629b8673c1f1c84a6e40c3e2d2c3bc3fdbde530a0a6b6fd68a22b867c3" } diff --git a/backend/.sqlx/query-d4fb94ee8198592c24e85d29078feb5220ab367e339510ee0e83bb7b5abfd184.json b/backend/.sqlx/query-d4fb94ee8198592c24e85d29078feb5220ab367e339510ee0e83bb7b5abfd184.json new file mode 100644 index 0000000000..a5ed70f19e --- /dev/null +++ b/backend/.sqlx/query-d4fb94ee8198592c24e85d29078feb5220ab367e339510ee0e83bb7b5abfd184.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE queue SET canceled = true, canceled_by = $1, scheduled_for = now(), suspend = 0 WHERE id = $2 RETURNING 1 as one", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "one", + "type_info": "Int4" + } + ], + "parameters": { + "Left": [ + "Varchar", + "Uuid" + ] + }, + "nullable": [ + null + ] + }, + "hash": "d4fb94ee8198592c24e85d29078feb5220ab367e339510ee0e83bb7b5abfd184" +} diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index b9851bc99a..cf3a0189e2 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -1013,49 +1013,62 @@ async fn cancel_all( ) -> error::JsonResult> { require_admin(authed.is_admin, &authed.username)?; - let mut jobs = sqlx::query!( - "UPDATE queue SET canceled = true, canceled_by = $2, scheduled_for = now(), suspend = 0 WHERE scheduled_for < now() AND workspace_id = $1 AND schedule_path IS NULL RETURNING id, running, is_flow_step", + let jobs = sqlx::query!( + "SELECT id, running, is_flow_step FROM queue WHERE scheduled_for < now() AND workspace_id = $1 AND schedule_path IS NULL", w_id, - authed.username ) .fetch_all(&db) .await?; let username = authed.username; + let mut uuids = vec![]; for j in jobs.iter() { - if !j.running && !j.is_flow_step.unwrap_or(false) { - let e = serde_json::json!({"message": format!("Job canceled: cancel_all by {username}"), "name": "Canceled", "reason": "cancel_all", "canceler": username}); - let job_running = get_queued_job(&j.id, &w_id, &db).await?; + let r = sqlx::query!( + "UPDATE queue SET canceled = true, canceled_by = $1, scheduled_for = now(), suspend = 0 WHERE id = $2 RETURNING 1 as one", + username, + j.id, + ) + .fetch_optional(&db) + .await; - if let Some(job_running) = job_running { - append_logs( - j.id, - w_id.clone(), - format!("canceled by {username}: cancel_all"), - db.clone(), - ) - .await; - let add_job = add_completed_job_error( - &db, - &job_running, - job_running.mem_peak.unwrap_or(0), - Some(CanceledBy { - username: Some(username.to_string()), - reason: Some("cancel_all".to_string()), - }), - e, - rsmq.clone(), - "server", - true, - ) - .await; - if let Err(e) = add_job { - tracing::error!("Failed to add canceled job: {}", e); + if r.as_ref().is_ok_and(|x| x.is_some()) { + uuids.push(j.id); + + if !j.running && !j.is_flow_step.unwrap_or(false) { + let e = serde_json::json!({"message": format!("Job canceled: cancel_all by {username}"), "name": "Canceled", "reason": "cancel_all", "canceler": username}); + let job_running = get_queued_job(&j.id, &w_id, &db).await?; + + if let Some(job_running) = job_running { + append_logs( + j.id, + w_id.clone(), + format!("canceled by {username}: cancel_all"), + db.clone(), + ) + .await; + let add_job = add_completed_job_error( + &db, + &job_running, + job_running.mem_peak.unwrap_or(0), + Some(CanceledBy { + username: Some(username.to_string()), + reason: Some("cancel_all".to_string()), + }), + e, + rsmq.clone(), + "server", + true, + ) + .await; + if let Err(e) = add_job { + tracing::error!("Failed to add canceled job: {}", e); + } } } + } else { + tracing::error!("Failed to cancel job: {:?} {:?}", j.id, r.err()); } } - let uuids = jobs.iter_mut().map(|j| j.id).collect::>(); Ok(Json(uuids)) }