From 59d731c35d6fcb2f793893ed47d999929c8a349a Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Thu, 20 Jun 2024 21:22:44 +0200 Subject: [PATCH] add timeout for cancel jobs (#3941) --- ...9cfe9a2ec315d6589e940be157a0563f81af.json} | 4 +- ...a4eafb2456c123f9b6b96bb4ba2409166e5a.json} | 4 +- backend/src/monitor.rs | 23 +++---- backend/windmill-api/src/jobs.rs | 60 ++++++++++++------- backend/windmill-queue/src/jobs.rs | 60 +++++++++++-------- 5 files changed, 85 insertions(+), 66 deletions(-) rename backend/.sqlx/{query-45c9ecf8b1f8cbca7c75dab24a1eb6da8ceb45258ee5817ec71e73bebbe415bd.json => query-099e7c7a66968575f896e0c11ecd9cfe9a2ec315d6589e940be157a0563f81af.json} (66%) rename backend/.sqlx/{query-c05be905e46c5b0a2186ba859a725495e55df7f2ad839aea22d0286525eb823e.json => query-215e0d320a304c8cb9ef12e7ea98a4eafb2456c123f9b6b96bb4ba2409166e5a.json} (65%) diff --git a/backend/.sqlx/query-45c9ecf8b1f8cbca7c75dab24a1eb6da8ceb45258ee5817ec71e73bebbe415bd.json b/backend/.sqlx/query-099e7c7a66968575f896e0c11ecd9cfe9a2ec315d6589e940be157a0563f81af.json similarity index 66% rename from backend/.sqlx/query-45c9ecf8b1f8cbca7c75dab24a1eb6da8ceb45258ee5817ec71e73bebbe415bd.json rename to backend/.sqlx/query-099e7c7a66968575f896e0c11ecd9cfe9a2ec315d6589e940be157a0563f81af.json index 02677837d1..59fb4a5dba 100644 --- a/backend/.sqlx/query-45c9ecf8b1f8cbca7c75dab24a1eb6da8ceb45258ee5817ec71e73bebbe415bd.json +++ b/backend/.sqlx/query-099e7c7a66968575f896e0c11ecd9cfe9a2ec315d6589e940be157a0563f81af.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE queue SET last_ping = now() WHERE id = $1 AND workspace_id = $2", + "query": "UPDATE queue SET last_ping = now() WHERE id = $1 AND workspace_id = $2 AND canceled = false", "describe": { "columns": [], "parameters": { @@ -11,5 +11,5 @@ }, "nullable": [] }, - "hash": "45c9ecf8b1f8cbca7c75dab24a1eb6da8ceb45258ee5817ec71e73bebbe415bd" + "hash": "099e7c7a66968575f896e0c11ecd9cfe9a2ec315d6589e940be157a0563f81af" } diff --git a/backend/.sqlx/query-c05be905e46c5b0a2186ba859a725495e55df7f2ad839aea22d0286525eb823e.json b/backend/.sqlx/query-215e0d320a304c8cb9ef12e7ea98a4eafb2456c123f9b6b96bb4ba2409166e5a.json similarity index 65% rename from backend/.sqlx/query-c05be905e46c5b0a2186ba859a725495e55df7f2ad839aea22d0286525eb823e.json rename to backend/.sqlx/query-215e0d320a304c8cb9ef12e7ea98a4eafb2456c123f9b6b96bb4ba2409166e5a.json index 6679b7849f..0889a223e6 100644 --- a/backend/.sqlx/query-c05be905e46c5b0a2186ba859a725495e55df7f2ad839aea22d0286525eb823e.json +++ b/backend/.sqlx/query-215e0d320a304c8cb9ef12e7ea98a4eafb2456c123f9b6b96bb4ba2409166e5a.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "UPDATE queue SET running = false, started_at = null WHERE id = $1", + "query": "UPDATE queue SET running = false, started_at = null WHERE id = $1 AND canceled = false", "describe": { "columns": [], "parameters": { @@ -10,5 +10,5 @@ }, "nullable": [] }, - "hash": "c05be905e46c5b0a2186ba859a725495e55df7f2ad839aea22d0286525eb823e" + "hash": "215e0d320a304c8cb9ef12e7ea98a4eafb2456c123f9b6b96bb4ba2409166e5a" } diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 542d40f1c5..8542eb818b 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1185,7 +1185,7 @@ async fn handle_zombie_flows( SELECT * FROM queue WHERE running = true AND suspend = 0 AND suspend_until IS null AND scheduled_for <= now() AND (job_kind = 'flow' OR job_kind = 'flowpreview') - AND last_ping IS NOT NULL AND last_ping < NOW() - ($1 || ' seconds')::interval + AND last_ping IS NOT NULL AND last_ping < NOW() - ($1 || ' seconds')::interval AND canceled = false "#, ).bind(FLOW_ZOMBIE_TRANSITION_TIMEOUT.as_str()) .fetch_all(db) @@ -1205,7 +1205,7 @@ async fn handle_zombie_flows( ); // if the flow hasn't started and is a zombie, we can simply restart it sqlx::query!( - "UPDATE queue SET running = false, started_at = null WHERE id = $1", + "UPDATE queue SET running = false, started_at = null WHERE id = $1 AND canceled = false", flow.id ) .execute(db) @@ -1237,11 +1237,12 @@ async fn handle_zombie_flows( .await?; for flow in flows2 { - let in_queue = - sqlx::query_as::<_, QueuedJob>("SELECT * FROM queue WHERE id = $1 AND running = true") - .bind(flow.parent_flow_id) - .fetch_optional(db) - .await?; + let in_queue = sqlx::query_as::<_, QueuedJob>( + "SELECT * FROM queue WHERE id = $1 AND running = true AND canceled = false", + ) + .bind(flow.parent_flow_id) + .fetch_optional(db) + .await?; if let Some(job) = in_queue { tracing::error!( "parallel Zombie flow detected: {} in workspace {}. Last ping was: {:?}.", @@ -1271,7 +1272,7 @@ async fn cancel_zombie_flow_job( flow.id, flow.workspace_id ); - let (mut ntx, _) = cancel_job( + let (ntx, _) = cancel_job( "monitor", Some(message), flow.id, @@ -1283,12 +1284,6 @@ async fn cancel_zombie_flow_job( false, ) .await?; - sqlx::query!( - "UPDATE queue SET running = false, started_at = null WHERE id = $1", - flow.id - ) - .execute(&mut *ntx) - .await?; ntx.commit().await?; Ok(()) } diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 3d00244fc9..61ae464c59 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -346,18 +346,26 @@ async fn cancel_job_api( }, }; - let (mut tx, job_option) = windmill_queue::cancel_job( - &audit_author.username, - reason, - id, - &w_id, - tx, - &db, - rsmq, - false, - opt_authed.is_none(), + let (mut tx, job_option) = tokio::time::timeout( + std::time::Duration::from_secs(120), + windmill_queue::cancel_job( + &audit_author.username, + reason, + id, + &w_id, + tx, + &db, + rsmq, + false, + opt_authed.is_none(), + ), ) - .await?; + .await + .map_err(|e| { + Error::InternalErr(format!( + "timeout after 120s while cancelling job {id} in {w_id}: {e:#}" + )) + })??; if let Some(id) = job_option { audit_log( @@ -453,18 +461,26 @@ async fn force_cancel( }, }; - let (mut tx, job_option) = windmill_queue::cancel_job( - &audit_author.username, - reason, - id, - &w_id, - tx, - &db, - rsmq, - true, - opt_authed.is_none(), + let (mut tx, job_option) = tokio::time::timeout( + std::time::Duration::from_secs(120), + windmill_queue::cancel_job( + &audit_author.username, + reason, + id, + &w_id, + tx, + &db, + rsmq, + true, + opt_authed.is_none(), + ), ) - .await?; + .await + .map_err(|e| { + Error::InternalErr(format!( + "timeout after 120s while cancelling job {id} in {w_id}: {e:#}" + )) + })??; if let Some(id) = job_option { audit_log( diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 183aab57eb..e305d704c2 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -160,32 +160,39 @@ pub async fn cancel_single_job<'c>( tracing::info!("Soft cancelling job {}", id); } } else { - let reason: String = reason - .clone() - .unwrap_or_else(|| "unexplicited reasons".to_string()); - let e = serde_json::json!({"message": format!("Job canceled: {reason} by {username}"), "name": "Canceled", "reason": reason, "canceler": username}); - append_logs( - &job_running.id, - w_id.to_string(), - format!("canceled by {username}: (force cancel: {force_cancel})"), - db, - ) - .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(reason) }), - e, - rsmq.clone(), - "server", - false, - ) - .await; + let username = username.to_string(); + let job_running = job_running.clone(); + let w_id = w_id.to_string(); + let db = db.clone(); + let rsmq = rsmq.clone(); + tokio::task::spawn(async move { + let reason: String = reason + .clone() + .unwrap_or_else(|| "unexplicited reasons".to_string()); + let e = serde_json::json!({"message": format!("Job canceled: {reason} by {username}"), "name": "Canceled", "reason": reason, "canceler": username}); + append_logs( + &job_running.id, + w_id.to_string(), + format!("canceled by {username}: (force cancel: {force_cancel})"), + &db, + ) + .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(reason) }), + e, + rsmq.clone(), + "server", + false, + ) + .await; - if let Err(e) = add_job { - tracing::error!("Failed to add canceled job: {}", e); - } + if let Err(e) = add_job { + tracing::error!("Failed to add canceled job: {}", e); + } + }); } if let Some(mut rsmq) = rsmq.clone() { rsmq.change_message_visibility(&job_running.tag, &job_running.id.to_string(), 0) @@ -235,6 +242,7 @@ pub async fn cancel_job<'c>( jobs.extend(new_jobs.clone()); jobs_to_cancel.extend(new_jobs); } + jobs.reverse(); let (ntx, _) = cancel_single_job( username, @@ -661,7 +669,7 @@ pub async fn add_completed_job< parent_job ); sqlx::query!( - "UPDATE queue SET last_ping = now() WHERE id = $1 AND workspace_id = $2", + "UPDATE queue SET last_ping = now() WHERE id = $1 AND workspace_id = $2 AND canceled = false", parent_job, &queued_job.workspace_id )