diff --git a/backend/.sqlx/query-3f9dd35d19e6badc92397fdb476d39ffadf1ff97931a9506772f76c3f36bd2c7.json b/backend/.sqlx/query-3f9dd35d19e6badc92397fdb476d39ffadf1ff97931a9506772f76c3f36bd2c7.json new file mode 100644 index 0000000000..e51ceff529 --- /dev/null +++ b/backend/.sqlx/query-3f9dd35d19e6badc92397fdb476d39ffadf1ff97931a9506772f76c3f36bd2c7.json @@ -0,0 +1,41 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT COUNT(*) AS count,\n MIN(scheduled_for) AS oldest_job,\n COUNT(*) FILTER (WHERE canceled_by IS NOT NULL) AS canceled_count,\n MIN(scheduled_for) FILTER (WHERE canceled_by IS NOT NULL) AS canceled_oldest\n FROM v2_job_queue\n WHERE tag = $1\n AND scheduled_for <= NOW() - $2::interval\n AND running = false\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "count", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "oldest_job", + "type_info": "Timestamptz" + }, + { + "ordinal": 2, + "name": "canceled_count", + "type_info": "Int8" + }, + { + "ordinal": 3, + "name": "canceled_oldest", + "type_info": "Timestamptz" + } + ], + "parameters": { + "Left": [ + "Text", + "Interval" + ] + }, + "nullable": [ + null, + null, + null, + null + ] + }, + "hash": "3f9dd35d19e6badc92397fdb476d39ffadf1ff97931a9506772f76c3f36bd2c7" +} diff --git a/backend/.sqlx/query-953566f05e2ada95ad440c42aede50a8265917b5c5a8eecbcfe139102aa48348.json b/backend/.sqlx/query-953566f05e2ada95ad440c42aede50a8265917b5c5a8eecbcfe139102aa48348.json deleted file mode 100644 index 0954e506dc..0000000000 --- a/backend/.sqlx/query-953566f05e2ada95ad440c42aede50a8265917b5c5a8eecbcfe139102aa48348.json +++ /dev/null @@ -1,35 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n WITH waiting AS (\n SELECT ck.key AS concurrency_key, cs.concurrent_limit,\n cs.concurrency_time_window_s,\n COUNT(*) AS waiting_count, MIN(q.scheduled_for) AS oldest_job\n FROM v2_job_queue q\n LEFT JOIN concurrency_key ck ON ck.job_id = q.id\n LEFT JOIN runnable_settings rs ON rs.hash = q.runnable_settings_handle\n LEFT JOIN concurrency_settings cs ON cs.hash = rs.concurrency_settings\n WHERE q.tag = $1\n AND q.scheduled_for <= NOW() - $2::interval\n AND q.running = false\n GROUP BY ck.key, cs.concurrent_limit, cs.concurrency_time_window_s\n ),\n classified AS MATERIALIZED (\n SELECT w.waiting_count, w.oldest_job,\n w.concurrent_limit > 0 AND w.concurrency_key <> '' AND (\n (SELECT COUNT(*) FROM jsonb_object_keys(cc.job_uuids))\n + (SELECT COUNT(*) FROM concurrency_key done\n WHERE done.key = w.concurrency_key\n AND done.ended_at >= NOW() - INTERVAL '1 second'\n * COALESCE(w.concurrency_time_window_s, 0))\n ) >= w.concurrent_limit AS gated\n FROM waiting w\n LEFT JOIN concurrency_counter cc ON cc.concurrency_id = w.concurrency_key\n )\n SELECT COALESCE(SUM(waiting_count) FILTER (WHERE gated IS NOT TRUE), 0)::bigint AS count,\n COALESCE(SUM(waiting_count) FILTER (WHERE gated), 0)::bigint AS gated_count,\n MIN(oldest_job) FILTER (WHERE gated IS NOT TRUE) AS oldest_job\n FROM classified\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "count", - "type_info": "Int8" - }, - { - "ordinal": 1, - "name": "gated_count", - "type_info": "Int8" - }, - { - "ordinal": 2, - "name": "oldest_job", - "type_info": "Timestamptz" - } - ], - "parameters": { - "Left": [ - "Text", - "Interval" - ] - }, - "nullable": [ - null, - null, - null - ] - }, - "hash": "953566f05e2ada95ad440c42aede50a8265917b5c5a8eecbcfe139102aa48348" -} diff --git a/backend/.sqlx/query-dbacbbbe27d1589dd26d6bcbd131a4ce32944cb2b36323117bbe92460e88b767.json b/backend/.sqlx/query-dbacbbbe27d1589dd26d6bcbd131a4ce32944cb2b36323117bbe92460e88b767.json new file mode 100644 index 0000000000..3f0b84d9ac --- /dev/null +++ b/backend/.sqlx/query-dbacbbbe27d1589dd26d6bcbd131a4ce32944cb2b36323117bbe92460e88b767.json @@ -0,0 +1,35 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH waiting AS (\n SELECT ck.key AS concurrency_key, cs.concurrent_limit,\n cs.concurrency_time_window_s,\n COUNT(*) AS waiting_count, MIN(q.scheduled_for) AS oldest_job\n FROM v2_job_queue q\n LEFT JOIN concurrency_key ck ON ck.job_id = q.id\n LEFT JOIN runnable_settings rs ON rs.hash = q.runnable_settings_handle\n LEFT JOIN concurrency_settings cs ON cs.hash = rs.concurrency_settings\n WHERE q.tag = $1\n AND q.scheduled_for <= NOW() - $2::interval\n AND q.running = false\n AND q.canceled_by IS NULL\n GROUP BY ck.key, cs.concurrent_limit, cs.concurrency_time_window_s\n ),\n classified AS MATERIALIZED (\n SELECT w.waiting_count, w.oldest_job,\n w.concurrent_limit > 0 AND w.concurrency_key <> '' AND (\n (SELECT COUNT(*) FROM jsonb_object_keys(cc.job_uuids))\n + (SELECT COUNT(*) FROM concurrency_key done\n WHERE done.key = w.concurrency_key\n AND done.ended_at >= NOW() - INTERVAL '1 second'\n * COALESCE(w.concurrency_time_window_s, 0))\n ) >= w.concurrent_limit AS gated\n FROM waiting w\n LEFT JOIN concurrency_counter cc ON cc.concurrency_id = w.concurrency_key\n )\n SELECT COALESCE(SUM(waiting_count) FILTER (WHERE gated IS NOT TRUE), 0)::bigint AS count,\n COALESCE(SUM(waiting_count) FILTER (WHERE gated), 0)::bigint AS gated_count,\n MIN(oldest_job) FILTER (WHERE gated IS NOT TRUE) AS oldest_job\n FROM classified\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "count", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "gated_count", + "type_info": "Int8" + }, + { + "ordinal": 2, + "name": "oldest_job", + "type_info": "Timestamptz" + } + ], + "parameters": { + "Left": [ + "Text", + "Interval" + ] + }, + "nullable": [ + null, + null, + null + ] + }, + "hash": "dbacbbbe27d1589dd26d6bcbd131a4ce32944cb2b36323117bbe92460e88b767" +} diff --git a/backend/.sqlx/query-dd2752edc8245a636baa187ac6ee05cc57be86a79be09b35223aa8170fb51cd4.json b/backend/.sqlx/query-dd2752edc8245a636baa187ac6ee05cc57be86a79be09b35223aa8170fb51cd4.json deleted file mode 100644 index fa71794c74..0000000000 --- a/backend/.sqlx/query-dd2752edc8245a636baa187ac6ee05cc57be86a79be09b35223aa8170fb51cd4.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n SELECT COUNT(*) AS count, MIN(scheduled_for) AS oldest_job\n FROM v2_job_queue\n WHERE tag = $1\n AND scheduled_for <= NOW() - $2::interval\n AND running = false\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "count", - "type_info": "Int8" - }, - { - "ordinal": 1, - "name": "oldest_job", - "type_info": "Timestamptz" - } - ], - "parameters": { - "Left": [ - "Text", - "Interval" - ] - }, - "nullable": [ - null, - null - ] - }, - "hash": "dd2752edc8245a636baa187ac6ee05cc57be86a79be09b35223aa8170fb51cd4" -} diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 93d7ffc9ab..af1bf852b3 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -f01fc53cad5ec0c598b6eac904e53443f683aae5 \ No newline at end of file +d54657fe4641a910cb2396de9cf83e850f118d2c \ No newline at end of file diff --git a/backend/windmill-common/tests/jobs_waiting_alerts.rs b/backend/windmill-common/tests/jobs_waiting_alerts.rs index 290faee269..d727a3272c 100644 --- a/backend/windmill-common/tests/jobs_waiting_alerts.rs +++ b/backend/windmill-common/tests/jobs_waiting_alerts.rs @@ -101,14 +101,30 @@ mod tests { /// `insert_concurrency_key` writes at push time for any runnable carrying a /// concurrent_limit, plus the settings handle that carries the limit. async fn queue_jobs(db: &Pool, n: usize, gate: Option<(&str, i64)>) { + queue_jobs_inner(db, n, gate, None).await + } + + /// Soft-cancelled jobs stay queued until a worker picks them up, and the + /// pull path skips concurrency limiting for them entirely. + async fn queue_canceled_jobs(db: &Pool, n: usize, gate: Option<(&str, i64)>) { + queue_jobs_inner(db, n, gate, Some("canceller")).await + } + + async fn queue_jobs_inner( + db: &Pool, + n: usize, + gate: Option<(&str, i64)>, + canceled_by: Option<&str>, + ) { for _ in 0..n { let id = Uuid::new_v4(); sqlx::query!( - "INSERT INTO v2_job_queue (id, workspace_id, tag, running, scheduled_for, runnable_settings_handle) - VALUES ($1, 'test-workspace', $2, false, NOW() - INTERVAL '5 minutes', $3)", + "INSERT INTO v2_job_queue (id, workspace_id, tag, running, scheduled_for, runnable_settings_handle, canceled_by) + VALUES ($1, 'test-workspace', $2, false, NOW() - INTERVAL '5 minutes', $3, $4)", id, TAG, gate.map(|(_, handle)| handle), + canceled_by, ) .execute(db) .await @@ -193,6 +209,29 @@ mod tests { ); } + /// Cancellation is a limiter bypass too: the pull path hands a job with + /// `canceled_by` set straight to a worker without consulting the counter. + /// Such jobs stay queued until a worker processes them, so a worker outage + /// piles them up behind whatever key they carry -- and classifying them off + /// that key hides an unbounded cancelled backlog from the alert. + #[ignore = "requires database setup - run with --ignored flag"] + #[sqlx::test(migrations = "../migrations")] + async fn canceled_jobs_behind_saturated_gate_still_alert(db: Pool) { + free_the_lock(&db).await; + configure_alert(&db, 100).await; + seed_gate(&db, GATE_KEY, 1, 1, 1).await; + queue_canceled_jobs(&db, 200, Some((GATE_KEY, 1))).await; + + jobs_waiting_alerts(&db).await; + + let messages = alert_messages(&db).await; + assert_eq!( + messages.len(), + 1, + "soft-cancelled jobs bypass the limiter and must still alert, got {messages:?}" + ); + } + /// An empty concurrency key is a limiter bypass: `apply_concurrency_limit` /// admits unconditionally instead of consulting the counter, so these jobs /// are runnable and any backlog of them is genuine capacity starvation.