diff --git a/backend/.sqlx/query-08f288d2781d823e109a9e5b8848234ca7d1efeee9661f3901f298da375e73f7.json b/backend/.sqlx/query-08f288d2781d823e109a9e5b8848234ca7d1efeee9661f3901f298da375e73f7.json index 2be39fce26..4bcf3c6ce3 100644 --- a/backend/.sqlx/query-08f288d2781d823e109a9e5b8848234ca7d1efeee9661f3901f298da375e73f7.json +++ b/backend/.sqlx/query-08f288d2781d823e109a9e5b8848234ca7d1efeee9661f3901f298da375e73f7.json @@ -130,28 +130,28 @@ }, { "ordinal": 25, - "name": "teams_command_script", - "type_info": "Text" - }, - { - "ordinal": 26, - "name": "teams_team_id", - "type_info": "Text" - }, - { - "ordinal": 27, - "name": "teams_team_name", - "type_info": "Text" - }, - { - "ordinal": 28, "name": "ai_models", "type_info": "VarcharArray" }, { - "ordinal": 29, + "ordinal": 26, "name": "code_completion_model", "type_info": "Varchar" + }, + { + "ordinal": 27, + "name": "teams_command_script", + "type_info": "Text" + }, + { + "ordinal": 28, + "name": "teams_team_id", + "type_info": "Text" + }, + { + "ordinal": 29, + "name": "teams_team_name", + "type_info": "Text" } ], "parameters": { @@ -185,10 +185,10 @@ true, true, true, - true, - true, - true, false, + true, + true, + true, true ] }, diff --git a/backend/.sqlx/query-1ab0d1ba1fbfad31ffb28a01a6c9640d0ac142aabee8d288a4f9c56ad9dbeac4.json b/backend/.sqlx/query-1ab0d1ba1fbfad31ffb28a01a6c9640d0ac142aabee8d288a4f9c56ad9dbeac4.json new file mode 100644 index 0000000000..1b5ddc6ac6 --- /dev/null +++ b/backend/.sqlx/query-1ab0d1ba1fbfad31ffb28a01a6c9640d0ac142aabee8d288a4f9c56ad9dbeac4.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO job_logs (job_id, logs)\n VALUES ($1, $2)\n ON CONFLICT (job_id) DO UPDATE SET logs = job_logs.logs || '\n' || EXCLUDED.logs\n WHERE job_logs.job_id = $1", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [] + }, + "hash": "1ab0d1ba1fbfad31ffb28a01a6c9640d0ac142aabee8d288a4f9c56ad9dbeac4" +} diff --git a/backend/.sqlx/query-55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2.json b/backend/.sqlx/query-55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2.json index 9243288c9d..14685a8bfa 100644 --- a/backend/.sqlx/query-55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2.json +++ b/backend/.sqlx/query-55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2.json @@ -130,28 +130,28 @@ }, { "ordinal": 25, - "name": "teams_command_script", - "type_info": "Text" - }, - { - "ordinal": 26, - "name": "teams_team_id", - "type_info": "Text" - }, - { - "ordinal": 27, - "name": "teams_team_name", - "type_info": "Text" - }, - { - "ordinal": 28, "name": "ai_models", "type_info": "VarcharArray" }, { - "ordinal": 29, + "ordinal": 26, "name": "code_completion_model", "type_info": "Varchar" + }, + { + "ordinal": 27, + "name": "teams_command_script", + "type_info": "Text" + }, + { + "ordinal": 28, + "name": "teams_team_id", + "type_info": "Text" + }, + { + "ordinal": 29, + "name": "teams_team_name", + "type_info": "Text" } ], "parameters": { @@ -185,10 +185,10 @@ true, true, true, - true, - true, - true, false, + true, + true, + true, true ] }, diff --git a/backend/.sqlx/query-653574b381a31548d82c1f6f3f44ec826597c42310ab78f9aacd7d9448206c6a.json b/backend/.sqlx/query-653574b381a31548d82c1f6f3f44ec826597c42310ab78f9aacd7d9448206c6a.json deleted file mode 100644 index 3d3309cb5c..0000000000 --- a/backend/.sqlx/query-653574b381a31548d82c1f6f3f44ec826597c42310ab78f9aacd7d9448206c6a.json +++ /dev/null @@ -1,34 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "WITH zombie_jobs AS (\n UPDATE v2_job_queue q SET running = false, started_at = null\n FROM v2_job j, v2_job_runtime r\n WHERE j.id = q.id AND j.id = r.id\n AND ping < now() - ($1 || ' seconds')::interval\n AND running = true\n AND kind NOT IN ('flow', 'flowpreview', 'flownode', 'singlescriptflow')\n AND same_worker = false\n RETURNING q.id, q.workspace_id, ping\n ),\n update_concurrency AS (\n UPDATE concurrency_counter cc\n SET job_uuids = job_uuids - zj.id::text\n FROM zombie_jobs zj\n INNER JOIN concurrency_key ck ON ck.job_id = zj.id\n WHERE cc.concurrency_id = ck.key\n )\n SELECT id, workspace_id, ping FROM zombie_jobs", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "workspace_id", - "type_info": "Varchar" - }, - { - "ordinal": 2, - "name": "ping", - "type_info": "Timestamptz" - } - ], - "parameters": { - "Left": [ - "Text" - ] - }, - "nullable": [ - false, - false, - true - ] - }, - "hash": "653574b381a31548d82c1f6f3f44ec826597c42310ab78f9aacd7d9448206c6a" -} diff --git a/backend/.sqlx/query-67afe352fc26dda9107c90e50e954642d877178ce2c0e73b72c3824135ef86f4.json b/backend/.sqlx/query-67afe352fc26dda9107c90e50e954642d877178ce2c0e73b72c3824135ef86f4.json deleted file mode 100644 index a09be9741e..0000000000 --- a/backend/.sqlx/query-67afe352fc26dda9107c90e50e954642d877178ce2c0e73b72c3824135ef86f4.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n INSERT INTO job_logs (job_id, logs)\n VALUES ($1, 'Restarted job after not receiving job''s ping for too long the ' || now() || '\n\n')\n ON CONFLICT (job_id) DO UPDATE SET logs = job_logs.logs || '\n' || EXCLUDED.logs\n WHERE job_logs.job_id = $1", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "67afe352fc26dda9107c90e50e954642d877178ce2c0e73b72c3824135ef86f4" -} diff --git a/backend/.sqlx/query-b45dc2baa48df2272dbac6e3537fc59ad8fab54027a9993cb61cdde68df3cbe6.json b/backend/.sqlx/query-b45dc2baa48df2272dbac6e3537fc59ad8fab54027a9993cb61cdde68df3cbe6.json new file mode 100644 index 0000000000..93dd534634 --- /dev/null +++ b/backend/.sqlx/query-b45dc2baa48df2272dbac6e3537fc59ad8fab54027a9993cb61cdde68df3cbe6.json @@ -0,0 +1,41 @@ +{ + "db_name": "PostgreSQL", + "query": "WITH to_update AS (\n SELECT q.id, q.workspace_id, r.ping, COALESCE(zjc.counter, 0) as counter\n FROM v2_job_queue q\n JOIN v2_job j ON j.id = q.id\n JOIN v2_job_runtime r ON r.id = j.id\n LEFT JOIN zombie_job_counter zjc ON zjc.job_id = q.id\n WHERE ping < now() - ($1 || ' seconds')::interval\n AND running = true\n AND kind NOT IN ('flow', 'flowpreview', 'flownode', 'singlescriptflow')\n AND same_worker = false\n AND (zjc.counter IS NULL OR zjc.counter <= $2)\n FOR UPDATE of q SKIP LOCKED\n ),\n zombie_jobs AS (\n UPDATE v2_job_queue q\n SET running = false, started_at = null\n FROM to_update tu\n WHERE q.id = tu.id AND (tu.counter IS NULL OR tu.counter < $2)\n RETURNING q.id, q.workspace_id, ping, tu.counter\n ),\n increment_counter AS (\n INSERT INTO zombie_job_counter (job_id, counter)\n SELECT id, 1 FROM to_update WHERE counter < $2\n ON CONFLICT (job_id) DO UPDATE \n SET counter = zombie_job_counter.counter + 1\n ),\n update_concurrency AS (\n UPDATE concurrency_counter cc\n SET job_uuids = job_uuids - zj.id::text\n FROM zombie_jobs zj\n INNER JOIN concurrency_key ck ON ck.job_id = zj.id\n WHERE cc.concurrency_id = ck.key\n )\n SELECT id AS \"id!\", workspace_id AS \"workspace_id!\", ping, counter + 1 AS counter FROM to_update", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "workspace_id!", + "type_info": "Varchar" + }, + { + "ordinal": 2, + "name": "ping", + "type_info": "Timestamptz" + }, + { + "ordinal": 3, + "name": "counter", + "type_info": "Int4" + } + ], + "parameters": { + "Left": [ + "Text", + "Int4" + ] + }, + "nullable": [ + false, + false, + true, + null + ] + }, + "hash": "b45dc2baa48df2272dbac6e3537fc59ad8fab54027a9993cb61cdde68df3cbe6" +} diff --git a/backend/migrations/20250205131522_add_zombie_job_counter.down.sql b/backend/migrations/20250205131522_add_zombie_job_counter.down.sql new file mode 100644 index 0000000000..e8686738c5 --- /dev/null +++ b/backend/migrations/20250205131522_add_zombie_job_counter.down.sql @@ -0,0 +1,3 @@ +-- Add down migration script here + +DROP TABLE zombie_job_counter; \ No newline at end of file diff --git a/backend/migrations/20250205131522_add_zombie_job_counter.up.sql b/backend/migrations/20250205131522_add_zombie_job_counter.up.sql new file mode 100644 index 0000000000..cc8c7457e3 --- /dev/null +++ b/backend/migrations/20250205131522_add_zombie_job_counter.up.sql @@ -0,0 +1,7 @@ +-- Add up migration script here + +CREATE TABLE IF NOT EXISTS zombie_job_counter ( + job_id UUID PRIMARY KEY REFERENCES v2_job (id) ON DELETE CASCADE, + counter INTEGER NOT NULL DEFAULT 0 +); + diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index b421054c2d..7aa0fe2c16 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1533,18 +1533,38 @@ pub async fn reload_base_url_setting(db: &DB) -> error::Result<()> { Ok(()) } +const RESTART_LIMIT: i32 = 3; + async fn handle_zombie_jobs(db: &Pool, base_internal_url: &str, worker_name: &str) { + let mut zombie_jobs_uuid_restart_limit_reached = vec![]; + if *RESTART_ZOMBIE_JOBS { let restarted = sqlx::query!( - "WITH zombie_jobs AS ( - UPDATE v2_job_queue q SET running = false, started_at = null - FROM v2_job j, v2_job_runtime r - WHERE j.id = q.id AND j.id = r.id - AND ping < now() - ($1 || ' seconds')::interval + "WITH to_update AS ( + SELECT q.id, q.workspace_id, r.ping, COALESCE(zjc.counter, 0) as counter + FROM v2_job_queue q + JOIN v2_job j ON j.id = q.id + JOIN v2_job_runtime r ON r.id = j.id + LEFT JOIN zombie_job_counter zjc ON zjc.job_id = q.id + WHERE ping < now() - ($1 || ' seconds')::interval AND running = true AND kind NOT IN ('flow', 'flowpreview', 'flownode', 'singlescriptflow') AND same_worker = false - RETURNING q.id, q.workspace_id, ping + AND (zjc.counter IS NULL OR zjc.counter <= $2) + FOR UPDATE of q SKIP LOCKED + ), + zombie_jobs AS ( + UPDATE v2_job_queue q + SET running = false, started_at = null + FROM to_update tu + WHERE q.id = tu.id AND (tu.counter IS NULL OR tu.counter < $2) + RETURNING q.id, q.workspace_id, ping, tu.counter + ), + increment_counter AS ( + INSERT INTO zombie_job_counter (job_id, counter) + SELECT id, 1 FROM to_update WHERE counter < $2 + ON CONFLICT (job_id) DO UPDATE + SET counter = zombie_job_counter.counter + 1 ), update_concurrency AS ( UPDATE concurrency_counter cc @@ -1553,8 +1573,9 @@ async fn handle_zombie_jobs(db: &Pool, base_internal_url: &str, worker INNER JOIN concurrency_key ck ON ck.job_id = zj.id WHERE cc.concurrency_id = ck.key ) - SELECT id, workspace_id, ping FROM zombie_jobs", + SELECT id AS \"id!\", workspace_id AS \"workspace_id!\", ping, counter + 1 AS counter FROM to_update", *ZOMBIE_JOB_TIMEOUT, + RESTART_LIMIT ) .fetch_all(db) .await @@ -1574,22 +1595,61 @@ async fn handle_zombie_jobs(db: &Pool, base_internal_url: &str, worker "no last ping".to_string() }; let url = format!("{}/run/{}?workspace={}", base_url, r.id, r.workspace_id,); - let error_message = format!( - "Zombie job {} on {} ({}) detected, restarting it, {}", - r.id, r.workspace_id, url, last_ping - ); + let restart = r.counter.is_none_or(|x| x < RESTART_LIMIT); + let (critical_error_message, restart_message) = if restart { + ( + format!( + "Zombie job {} on {} ({}) detected, restarting it ({}/{} attempts), last ping: {}", + r.id, + r.workspace_id, + url, + r.counter.unwrap_or(0) + 1, + RESTART_LIMIT, + last_ping + ), + format!( + "Restarted job after not receiving job's ping for too long the {} ({}/{} attempts)\n\n", + last_ping, + r.counter.unwrap_or(0) + 1, + RESTART_LIMIT + ) + ) + } else { + ( + format!( + "Zombie job {} on {} ({}) detected, but restart limit ({}) reached, job will be processed as an error, last ping: {}", + r.id, r.workspace_id, url, RESTART_LIMIT, last_ping + ), + format!( + "job's ping was received last at {}, job will be processed as an error since all {} restart attempts failed", + last_ping, RESTART_LIMIT + ) + ) + }; - let _ = sqlx::query!(" + let _ = sqlx::query!( + " INSERT INTO job_logs (job_id, logs) - VALUES ($1, 'Restarted job after not receiving job''s ping for too long the ' || now() || '\n\n') + VALUES ($1, $2) ON CONFLICT (job_id) DO UPDATE SET logs = job_logs.logs || '\n' || EXCLUDED.logs WHERE job_logs.job_id = $1", - r.id + r.id, + restart_message ) .execute(db) .await; - tracing::error!(error_message); - report_critical_error(error_message, db.clone(), Some(&r.workspace_id), None).await; + tracing::error!(critical_error_message); + report_critical_error( + critical_error_message, + db.clone(), + Some(&r.workspace_id), + None, + ) + .await; + + if !restart { + zombie_jobs_uuid_restart_limit_reached.push(r.id); + } } } @@ -1664,9 +1724,43 @@ async fn handle_zombie_jobs(db: &Pool, base_internal_url: &str, worker .unwrap_or_else(|| vec![]) }; + enum ErrorMessage { + RestartLimit, + SameWorker, + RestartDisabled, + } + + impl ErrorMessage { + fn to_string(&self) -> String { + match self { + ErrorMessage::RestartLimit => format!("RestartLimit ({})", RESTART_LIMIT), + ErrorMessage::SameWorker => "SameWorker".to_string(), + ErrorMessage::RestartDisabled => "RestartDisabled".to_string(), + } + } + } + + let zombie_jobs_restart_limit_reached = + sqlx::query_as::<_, QueuedJob>("SELECT * FROM v2_as_queue WHERE id = ANY($1)") + .bind(&zombie_jobs_uuid_restart_limit_reached[..]) + .fetch_all(db) + .await + .ok() + .unwrap_or_else(|| vec![]); + let timeouts = non_restartable_jobs .into_iter() - .chain(same_worker_timeout_jobs) + .map(|x| (x, ErrorMessage::RestartDisabled)) + .chain( + same_worker_timeout_jobs + .into_iter() + .map(|x| (x, ErrorMessage::SameWorker)), + ) + .chain( + zombie_jobs_restart_limit_reached + .into_iter() + .map(|x| (x, ErrorMessage::RestartLimit)), + ) .collect::>(); #[cfg(feature = "prometheus")] @@ -1674,7 +1768,7 @@ async fn handle_zombie_jobs(db: &Pool, base_internal_url: &str, worker QUEUE_ZOMBIE_DELETE_COUNT.inc_by(timeouts.len() as _); } - for job in timeouts { + for (job, error_kind) in timeouts { // since the job is unrecoverable, the same worker queue should never be sent anything let (same_worker_tx_never_used, _same_worker_rx_never_used) = mpsc::channel::(1); @@ -1709,20 +1803,21 @@ async fn handle_zombie_jobs(db: &Pool, base_internal_url: &str, worker }; let last_ping = job.last_ping.clone(); + let error_message = format!( + "Job timed out after no ping from job since {} (ZOMBIE_JOB_TIMEOUT: {}, reason: {:?})", + last_ping + .map(|x| x.to_string()) + .unwrap_or_else(|| "no ping".to_string()), + *ZOMBIE_JOB_TIMEOUT, + error_kind.to_string() + ); let _ = handle_job_error( db, &client, &job, 0, None, - error::Error::ExecutionErr(format!( - "Job timed out after no ping from job since {} (ZOMBIE_JOB_TIMEOUT: {}, same_worker: {})", - last_ping - .map(|x| x.to_string()) - .unwrap_or_else(|| "no ping".to_string()), - *ZOMBIE_JOB_TIMEOUT, - job.same_worker - )), + error::Error::ExecutionErr(error_message), true, same_worker_tx_never_used, "", diff --git a/backend/windmill-worker/src/result_processor.rs b/backend/windmill-worker/src/result_processor.rs index 4b433c4546..bb4c1937ae 100644 --- a/backend/windmill-worker/src/result_processor.rs +++ b/backend/windmill-worker/src/result_processor.rs @@ -653,7 +653,6 @@ pub async fn handle_job_error( if let Some(f) = update_job_future { let _ = f().await; } - tracing::error!(job_id = %job.id, "error handling job: {err:?} {} {} {}", job.id, job.workspace_id, job.created_by); } #[derive(Debug, Serialize)]