diff --git a/backend/tests/batch_pull.rs b/backend/tests/batch_pull.rs new file mode 100644 index 0000000000..184b85dd06 --- /dev/null +++ b/backend/tests/batch_pull.rs @@ -0,0 +1,278 @@ +//! Pins what `pull_batch` promises a caller serving many waiting workers from one claim: every +//! claimed job is marked running under exactly one of those workers, concurrent batches never +//! claim the same job, the queues are walked in the same order a single pull walks them, and +//! the batch peek keeps the ordered index scan the single pull relies on. + +use std::collections::HashSet; + +use serde_json::Value; +use sqlx::{Pool, Postgres}; +use uuid::Uuid; +use windmill_common::worker::{make_batch_pull_query, PullQueue}; +use windmill_queue::{pull_batch, PulledJobResult}; + +const W_ID: &str = "test-workspace"; + +async fn queue_job(db: &Pool, tag: &str, suspend: Option) -> anyhow::Result { + let id = Uuid::new_v4(); + sqlx::query( + "INSERT INTO v2_job (id, workspace_id, created_by, created_at, permissioned_as, \ + permissioned_as_email, kind, tag, args, visible_to_owner) \ + VALUES ($1, $2, 'test-user', now(), 'u/test-user', 'test@windmill.dev', 'noop', $3, \ + '{}', true)", + ) + .bind(id) + .bind(W_ID) + .bind(tag) + .execute(db) + .await?; + // A suspended flow sits running with `suspend_until` set; `suspend = 0` makes it resumable. + sqlx::query( + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, tag, suspend, \ + suspend_until) \ + VALUES ($1, $2, now() - interval '1 second', $3, $4, COALESCE($5, 0), \ + CASE WHEN $5::int IS NULL THEN NULL ELSE now() + interval '1 day' END)", + ) + .bind(id) + .bind(W_ID) + .bind(suspend.is_some()) + .bind(tag) + .bind(suspend) + .execute(db) + .await?; + sqlx::query("INSERT INTO v2_job_runtime (id) VALUES ($1)") + .bind(id) + .execute(db) + .await?; + Ok(id) +} + +fn names(prefix: &str, n: usize) -> Vec { + (0..n).map(|i| format!("{prefix}-{i}")).collect() +} + +fn job_id(res: &PulledJobResult) -> Uuid { + res.job + .as_ref() + .expect("an admitted result carries its job") + .id +} + +#[sqlx::test(fixtures("base"))] +async fn batch_pull_assigns_each_job_to_one_waiting_worker( + db: Pool, +) -> anyhow::Result<()> { + for _ in 0..30 { + queue_job(&db, "batch", None).await?; + } + let groups = vec![vec!["batch".to_string()]]; + + // Three batches race over one queue; a name listed twice still gets one job. + let mut a = names("a", 8); + a.push("a-0".to_string()); + let (b, c) = (names("b", 8), names("c", 8)); + let (ra, rb, rc) = tokio::join!( + pull_batch(&db, &groups, &a, false), + pull_batch(&db, &groups, &b, false), + pull_batch(&db, &groups, &c, false), + ); + let claimed: Vec<(String, PulledJobResult)> = [ra?, rb?, rc?].into_iter().flatten().collect(); + assert_eq!( + claimed.len(), + 24, + "30 queued jobs cover all 24 distinct workers" + ); + let workers: HashSet<&String> = claimed.iter().map(|(w, _)| w).collect(); + let jobs: HashSet = claimed.iter().map(|(_, r)| job_id(r)).collect(); + assert_eq!(workers.len(), 24, "no worker is handed two jobs"); + assert_eq!(jobs.len(), 24, "no job is handed to two workers"); + + for (worker, res) in &claimed { + let (running, row_worker): (bool, Option) = + sqlx::query_as("SELECT running, worker FROM v2_job_queue WHERE id = $1") + .bind(job_id(res)) + .fetch_one(&db) + .await?; + assert!(running); + assert_eq!( + row_worker.as_ref(), + Some(worker), + "job is running under its worker" + ); + } + + // The result is the wire type a remote caller hands to its worker. + let (_, first) = &claimed[0]; + let round_trip: PulledJobResult = serde_json::from_str(&serde_json::to_string(first)?)?; + assert_eq!(job_id(&round_trip), job_id(first)); + + // Six jobs remain for ten workers: four go without. + let rest = pull_batch(&db, &groups, &names("d", 10), false).await?; + assert_eq!(rest.len(), 6); + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn batch_pull_walks_suspended_then_priority_groups(db: Pool) -> anyhow::Result<()> { + let resumable = queue_job(&db, "low", Some(0)).await?; + let mut high = HashSet::new(); + for _ in 0..3 { + high.insert(queue_job(&db, "high", None).await?); + queue_job(&db, "low", None).await?; + } + let groups = vec![vec!["high".to_string()], vec!["low".to_string()]]; + + let claimed = pull_batch(&db, &groups, &names("w", 5), true).await?; + assert_eq!(claimed.len(), 5); + let suspended: Vec = claimed + .iter() + .filter(|(_, r)| r.suspended) + .map(|(_, r)| job_id(r)) + .collect(); + assert_eq!( + suspended, + vec![resumable], + "the resumable flow is claimed first" + ); + let ready: HashSet = claimed + .iter() + .filter(|(_, r)| !r.suspended) + .map(|(_, r)| job_id(r)) + .collect(); + assert!( + high.is_subset(&ready), + "the high group drains before the low one" + ); + assert_eq!(ready.len(), 4); + Ok(()) +} + +// A job over its concurrency limit must go back to the queue rather than sit claimed, and the +// worker it was claimed for is pulled for again in the same order: a resumable flow still wins +// over a ready job on the re-pull, as it does on every loop of a single pull. +#[cfg(feature = "enterprise")] +#[sqlx::test(fixtures("base"))] +async fn batch_pull_requeues_over_limit_job_and_serves_its_worker( + db: Pool, +) -> anyhow::Result<()> { + let limited = queue_job(&db, "batch", Some(0)).await?; + sqlx::query( + "UPDATE v2_job SET kind = 'script', runnable_path = 'f/test/limited', \ + concurrent_limit = 1, concurrency_time_window_s = 0 WHERE id = $1", + ) + .bind(limited) + .execute(&db) + .await?; + // Claimed first, so the bounce happens before the other jobs are reached. + sqlx::query("UPDATE v2_job_queue SET priority = 10 WHERE id = $1") + .bind(limited) + .execute(&db) + .await?; + sqlx::query("INSERT INTO concurrency_key (key, job_id) VALUES ('limited-key', $1)") + .bind(limited) + .execute(&db) + .await?; + sqlx::query( + "INSERT INTO concurrency_counter (concurrency_id, job_uuids) \ + VALUES ('limited-key', jsonb_build_object(gen_random_uuid()::text, '{}'::jsonb))", + ) + .execute(&db) + .await?; + let resumable = queue_job(&db, "batch", Some(0)).await?; + queue_job(&db, "batch", None).await?; + + let groups = vec![vec!["batch".to_string()]]; + let claimed = pull_batch(&db, &groups, &names("w", 1), true).await?; + assert_eq!(claimed.len(), 1); + assert_eq!(claimed[0].0, "w-0"); + assert_eq!( + job_id(&claimed[0].1), + resumable, + "the worker is served the other resumable flow" + ); + assert!(claimed[0].1.suspended); + + let (running, rescheduled): (bool, bool) = + sqlx::query_as("SELECT running, scheduled_for > now() FROM v2_job_queue WHERE id = $1") + .bind(limited) + .fetch_one(&db) + .await?; + assert!(!running, "the over-limit job is back in the queue"); + assert!(rescheduled, "the over-limit job waits for a free slot"); + Ok(()) +} + +// Admission failing for one claimed job must hand it back to the queue: it is already marked +// running under its worker, which will never be told about it. +#[sqlx::test(fixtures("base"))] +async fn batch_pull_requeues_job_whose_admission_fails(db: Pool) -> anyhow::Result<()> { + let broken = queue_job(&db, "batch", None).await?; + // A settings handle with no `runnable_settings` row makes the settings lookup fail. + sqlx::query("UPDATE v2_job_queue SET runnable_settings_handle = -42 WHERE id = $1") + .bind(broken) + .execute(&db) + .await?; + + let groups = vec![vec!["batch".to_string()]]; + let claimed = pull_batch(&db, &groups, &names("w", 1), false).await?; + assert!(claimed.is_empty()); + + let (running, started): (bool, bool) = + sqlx::query_as("SELECT running, started_at IS NOT NULL FROM v2_job_queue WHERE id = $1") + .bind(broken) + .fetch_one(&db) + .await?; + assert!(!running && !started, "the job is back in the queue, not stranded"); + Ok(()) +} + +/// Depth-first walk of an `EXPLAIN (FORMAT JSON)` plan tree. +fn nodes(plan: &Value, out: &mut Vec) { + out.push(plan.clone()); + for child in plan["Plans"].as_array().unwrap_or(&vec![]) { + nodes(child, out); + } +} + +// The batch LIMIT is a bind parameter, so a cached generic plan cannot see it and the planner +// falls back to guessing a fraction of the matching rows. On a deep backlog, a guess that tips +// it into a bitmap scan plus sort reads every queued job of the tag on each pull. +#[sqlx::test(fixtures("base"))] +async fn batch_pull_generic_plan_walks_the_queue_index(db: Pool) -> anyhow::Result<()> { + sqlx::query( + "INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, tag) + SELECT gen_random_uuid(), 'test-workspace', now() - make_interval(secs => i % 1000), + false, CASE WHEN i % 10 = 0 THEN 'python3' ELSE 'deno' END + FROM generate_series(1, 200000) i", + ) + .execute(&db) + .await?; + sqlx::query("ANALYZE v2_job_queue").execute(&db).await?; + + let mut conn = db.acquire().await?; + sqlx::query("SET plan_cache_mode = force_generic_plan") + .execute(&mut *conn) + .await?; + let query = make_batch_pull_query(&["python3".to_string()], PullQueue::Ready); + sqlx::query(&format!("PREPARE batch_pull(text[]) AS {query}")) + .execute(&mut *conn) + .await?; + let explained: Value = + sqlx::query_scalar("EXPLAIN (FORMAT JSON) EXECUTE batch_pull(ARRAY['w-0', 'w-1'])") + .fetch_one(&mut *conn) + .await?; + + let mut all = vec![]; + nodes(&explained[0]["Plan"], &mut all); + let pretty = serde_json::to_string_pretty(&explained)?; + assert!( + all.iter() + .any(|n| n["Node Type"] == "Index Scan" && n["Index Name"] == "queue_sort_v2"), + "batch peek does not walk queue_sort_v2 in order:\n{pretty}" + ); + assert!( + !all.iter().any(|n| n["Node Type"] == "Sort"), + "batch peek sorts the backlog:\n{pretty}" + ); + Ok(()) +} diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index d480a89851..449fbd7606 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -774,6 +774,60 @@ impl From<&Pool> for Connection { } } +/// The columns every pull returns, read by name into `windmill_queue::PulledJob`. Expects `q` +/// (the claimed `v2_job_queue` rows), `j` (their `v2_job` rows), `f`, `p` and `pj` in scope. +const PULLED_JOB_COLUMNS: &str = "j.id, j.workspace_id, j.parent_job, j.created_by, q.started_at, q.scheduled_for, + j.runnable_id, j.runnable_path, j.args, q.canceled_by, + q.canceled_reason, j.kind, j.trigger, j.trigger_kind, j.permissioned_as, + f.flow_status, j.script_lang, + j.same_worker, j.pre_run_error, j.visible_to_owner, + j.tag, j.concurrent_limit, j.concurrency_time_window_s, j.flow_innermost_root_job, j.root_job, + j.timeout, j.flow_step_id, j.cache_ttl, q.cache_ignore_s3_path, q.runnable_settings_handle, j.priority, j.raw_code, j.raw_lock, j.raw_flow, + j.script_entrypoint_override, j.preprocessed, COALESCE(pj.runnable_path, j.args->>'_FLOW_PATH') as parent_runnable_path, + COALESCE(p.email, j.permissioned_as_email) as permissioned_as_email, p.username as permissioned_as_username, p.is_admin as permissioned_as_is_admin, + p.is_operator as permissioned_as_is_operator, p.groups as permissioned_as_groups, p.folders as permissioned_as_folders, p.end_user_email as permissioned_as_end_user_email"; + +fn sql_tag_list(tags: &[String]) -> String { + tags.iter() + .map(|x| format!("'{}'", x.replace('\'', "''"))) + .join(", ") +} + +/// Selects the next runnable jobs for `tags`, highest priority first. `exclude_overloaded` names +/// the text[] bind parameter holding workspaces that workspace fairness currently caps. +fn ready_peek(tags: &[String], exclude_overloaded: Option<&str>, limit: &str) -> String { + let fairness = exclude_overloaded + .map(|param| format!("\n AND workspace_id <> ALL({param}::text[])")) + .unwrap_or_default(); + format!( + "SELECT id + FROM v2_job_queue + WHERE running = false AND tag IN ({}) AND scheduled_for <= now(){fairness} + ORDER BY priority DESC NULLS LAST, scheduled_for + FOR UPDATE SKIP LOCKED + LIMIT {limit}", + sql_tag_list(tags) + ) +} + +// The `CASE` is `suspend <= 0 OR suspend_until <= now()` written as one indexable +// expression, equivalent only under the `suspend_until IS NOT NULL` guard. It must stay in +// sync with `queue_suspended_v2` (migration 20260826202939): if it no longer matches, the +// test silently reverts to a heap filter over every suspended row on every worker poll. +fn suspended_peek(tags: &[String], limit: &str) -> String { + format!( + "SELECT id + FROM v2_job_queue + WHERE suspend_until IS NOT NULL + AND (CASE WHEN suspend <= 0 THEN '-infinity'::timestamptz ELSE suspend_until END) <= now() + AND tag IN ({}) + ORDER BY priority DESC NULLS LAST, created_at + FOR UPDATE SKIP LOCKED + LIMIT {limit}", + sql_tag_list(tags) + ) +} + fn format_pull_query(peek: String) -> String { let r = format!( "WITH peek AS ( @@ -803,16 +857,7 @@ fn format_pull_query(peek: String) -> String { raw_flow, script_entrypoint_override, preprocessed FROM v2_job WHERE id = (SELECT id FROM peek) - ) SELECT j.id, j.workspace_id, j.parent_job, j.created_by, started_at, scheduled_for, - j.runnable_id, j.runnable_path, j.args, canceled_by, - canceled_reason, j.kind, j.trigger, j.trigger_kind, j.permissioned_as, - flow_status, j.script_lang, - j.same_worker, j.pre_run_error, j.visible_to_owner, - j.tag, j.concurrent_limit, j.concurrency_time_window_s, j.flow_innermost_root_job, j.root_job, - j.timeout, j.flow_step_id, j.cache_ttl, q.cache_ignore_s3_path, q.runnable_settings_handle, j.priority, j.raw_code, j.raw_lock, j.raw_flow, - j.script_entrypoint_override, j.preprocessed, COALESCE(pj.runnable_path, j.args->>'_FLOW_PATH') as parent_runnable_path, - COALESCE(p.email, j.permissioned_as_email) as permissioned_as_email, p.username as permissioned_as_username, p.is_admin as permissioned_as_is_admin, - p.is_operator as permissioned_as_is_operator, p.groups as permissioned_as_groups, p.folders as permissioned_as_folders, p.end_user_email as permissioned_as_end_user_email + ) SELECT {PULLED_JOB_COLUMNS} FROM q, j LEFT JOIN v2_job_status f USING (id) LEFT JOIN job_perms p ON p.job_id = j.id @@ -824,22 +869,63 @@ fn format_pull_query(peek: String) -> String { r } -// The `CASE` is `suspend <= 0 OR suspend_until <= now()` written as one indexable -// expression, equivalent only under the `suspend_until IS NOT NULL` guard. It must stay in -// sync with `queue_suspended_v2` (migration 20260826202939): if it no longer matches, the -// test silently reverts to a heap filter over every suspended row on every worker poll. +/// Claims up to `cardinality($1)` jobs in one statement and assigns the i-th claimed job to the +/// worker named `$1[i]`, so each row is marked running under the worker that will run it. +/// Returns the `PULLED_JOB_COLUMNS` plus `assigned_worker`. +fn format_batch_pull_query(peek: String) -> String { + format!( + "WITH peek AS MATERIALIZED ( + {peek} + ), assign AS ( + SELECT id, ($1::text[])[row_number() OVER ()] AS assigned_worker FROM peek + ), q AS ( + UPDATE v2_job_queue SET + running = true, + started_at = coalesce(started_at, now()), + suspend_until = null, + worker = assign.assigned_worker + FROM assign + WHERE v2_job_queue.id = assign.id + RETURNING + v2_job_queue.id, started_at, scheduled_for, + canceled_by, canceled_reason, worker, cache_ignore_s3_path, runnable_settings_handle + ), r AS ( + UPDATE v2_job_runtime SET + ping = now() + WHERE id = ANY(ARRAY(SELECT id FROM q)) + ) SELECT {PULLED_JOB_COLUMNS}, q.worker AS assigned_worker + FROM q + JOIN v2_job j ON j.id = q.id + LEFT JOIN v2_job_status f ON f.id = q.id + LEFT JOIN job_perms p ON p.job_id = q.id + LEFT JOIN v2_job pj ON j.parent_job = pj.id" + ) +} + +/// Which part of the queue a pull claims from. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum PullQueue { + /// Jobs ready to start. + Ready, + /// Jobs ready to start, excluding the workspaces bound as `$2::text[]`. + ReadyFairness, + /// Suspended flows whose resume condition is met. + Suspended, +} + +/// Batch variant of the pull queries: binds the claiming workers' names as `$1::text[]`, which +/// must be distinct, and claims at most one job per name. +pub fn make_batch_pull_query(tags: &[String], queue: PullQueue) -> String { + const LIMIT: &str = "cardinality($1::text[])"; + format_batch_pull_query(match queue { + PullQueue::Ready => ready_peek(tags, None, LIMIT), + PullQueue::ReadyFairness => ready_peek(tags, Some("$2"), LIMIT), + PullQueue::Suspended => suspended_peek(tags, LIMIT), + }) +} + pub fn make_suspended_pull_query(tags: &[String]) -> String { - format_pull_query(format!( - "SELECT id - FROM v2_job_queue - WHERE suspend_until IS NOT NULL - AND (CASE WHEN suspend <= 0 THEN '-infinity'::timestamptz ELSE suspend_until END) <= now() - AND tag IN ({}) - ORDER BY priority DESC NULLS LAST, created_at - FOR UPDATE SKIP LOCKED - LIMIT 1", - tags.iter().map(|x| format!("'{x}'")).join(", ") - )) + format_pull_query(suspended_peek(tags, "1")) } // pub async fn make_suspended pub async fn store_suspended_pull_query(wc: &WorkerConfig) { @@ -852,16 +938,7 @@ pub async fn store_suspended_pull_query(wc: &WorkerConfig) { } pub fn make_pull_query(tags: &[String]) -> String { - let query = format_pull_query(format!( - "SELECT id - FROM v2_job_queue - WHERE running = false AND tag IN ({}) AND scheduled_for <= now() - ORDER BY priority DESC NULLS LAST, scheduled_for - FOR UPDATE SKIP LOCKED - LIMIT 1", - tags.iter().map(|x| format!("'{x}'")).join(", ") - )); - query + format_pull_query(ready_peek(tags, None, "1")) } // Variant of `make_pull_query` that additionally excludes jobs whose workspace_id is in the @@ -872,17 +949,7 @@ pub fn make_pull_query(tags: &[String]) -> String { // `pub(crate)` because only `store_pull_query` consumes it; the resulting query string is what // crosses crate boundaries via `WORKER_PULL_QUERIES_FAIRNESS`. pub(crate) fn make_pull_query_fairness(tags: &[String]) -> String { - let query = format_pull_query(format!( - "SELECT id - FROM v2_job_queue - WHERE running = false AND tag IN ({}) AND scheduled_for <= now() - AND workspace_id <> ALL($2::text[]) - ORDER BY priority DESC NULLS LAST, scheduled_for - FOR UPDATE SKIP LOCKED - LIMIT 1", - tags.iter().map(|x| format!("'{x}'")).join(", ") - )); - query + format_pull_query(ready_peek(tags, Some("$2"), "1")) } pub async fn store_pull_query(wc: &WorkerConfig) { diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 7b552fb02c..8f65638790 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -78,9 +78,9 @@ use windmill_common::{ users::{SUPERADMIN_NOTIFICATION_EMAIL, SUPERADMIN_SECRET_EMAIL}, utils::{not_found_if_none, report_critical_error, StripPath, WarnAfterExt}, worker::{ - to_raw_value, CLOUD_HOSTED, DISABLE_FLOW_SCRIPT, NO_LOGS, PREVIEW_TAGS_OVERRIDE, - WORKER_PULL_QUERIES, WORKER_PULL_QUERIES_FAIRNESS, WORKER_SUSPENDED_PULL_QUERY, - WORKSPACE_FAIRNESS_OVERLOADED, + make_batch_pull_query, to_raw_value, PullQueue, CLOUD_HOSTED, DISABLE_FLOW_SCRIPT, NO_LOGS, + PREVIEW_TAGS_OVERRIDE, WORKER_PULL_QUERIES, WORKER_PULL_QUERIES_FAIRNESS, + WORKER_SUSPENDED_PULL_QUERY, WORKSPACE_FAIRNESS_OVERLOADED, }, DB, METRICS_ENABLED, }; @@ -3570,7 +3570,7 @@ impl MiniPulledJob { } } -#[derive(sqlx::FromRow, Debug, Clone)] +#[derive(sqlx::FromRow, Debug, Clone, Serialize, Deserialize)] pub struct PulledJob { #[sqlx(flatten)] pub job: MiniPulledJob, @@ -3837,7 +3837,7 @@ pub async fn get_queued_job_v2<'c>( Ok(job) } -#[derive(Debug)] +#[derive(Debug, Serialize, Deserialize)] pub struct PulledJobResult { pub job: Option, pub suspended: bool, @@ -4228,6 +4228,13 @@ async fn clone_runnable(j: &mut PulledJob, db: &DB) -> error::Result<()> { Ok(()) } +// Cap to bound DB work per pull cycle. Each iteration on an over-limit job runs +// the full apply_concurrency_limit query stack (~7 queries), so without a tight +// cap a single pull() can issue thousands of queries when the queue is full of +// jobs sharing one over-limit concurrency key. Capping at 10 lets workers skip +// a few stale jobs in healthy conditions while preventing storm amplification. +const PULL_LOOP_LIMIT: i32 = 10; + // TODO: Factorize /// Pull the job from queue pub async fn pull( @@ -4242,14 +4249,9 @@ pub async fn pull( let mut pull_loop_count = 0; loop { pull_loop_count += 1; - // Cap to bound DB work per pull cycle. Each iteration on an over-limit job runs - // the full apply_concurrency_limit query stack (~7 queries), so without a tight - // cap a single pull() can issue thousands of queries when the queue is full of - // jobs sharing one over-limit concurrency key. Capping at 10 lets workers skip - // a few stale jobs in healthy conditions while preventing storm amplification. - if pull_loop_count > 10 { + if pull_loop_count > PULL_LOOP_LIMIT { tracing::warn!( - "Pull job loop count exceeded 10, backing off (likely concurrency re-queue storm)" + "Pull job loop count exceeded {PULL_LOOP_LIMIT}, backing off (likely concurrency re-queue storm)" ); return Ok(PulledJobResult { job: None, @@ -4389,68 +4391,257 @@ pub async fn pull( }); }; - let concurrency_settings = windmill_common::runnable_settings::prefetch_cached_from_handle( - job.runnable_settings_handle, - db, - ) - .await? - .1 - .maybe_fallback(None, job.concurrent_limit, job.concurrency_time_window_s); - - let has_concurent_limit = - has_active_concurrency_limit(concurrency_settings.concurrent_limit); - - #[cfg(not(feature = "enterprise"))] - if has_concurent_limit && !job.is_dependency() { - tracing::error!("Concurrent limits are an EE feature only, ignoring constraints") - } - - #[cfg(not(feature = "enterprise"))] - let has_concurent_limit = job.is_dependency() - && has_active_concurrency_limit(job.concurrent_limit) - && cfg!(feature = "private") - && !*WMDEBUG_NO_DEBOUNCING; - // if we don't have private flag, we don't have concurrency limit - - // concurrency check. If more than X jobs for this path are already running, we re-queue and pull another job from the queue - let pulled_job = job; - if pulled_job.runnable_path.is_none() - || !has_concurent_limit - || pulled_job.canceled_by.is_some() - { - #[cfg(feature = "prometheus")] - if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { - QUEUE_PULL_COUNT.inc(); - } - otel_incr_queue_pull_count(); - return Ok(PulledJobResult { - job: Some(pulled_job), - suspended, - missing_concurrency_key: false, - error_while_preprocessing: None, - }); - } - - #[cfg(feature = "private")] - if cfg!(feature = "enterprise") || (pulled_job.is_dependency() && !*WMDEBUG_NO_DEBOUNCING) { - if let Some(pulled_job_res) = timeout( - Duration::from_secs(15), - crate::jobs_ee::apply_concurrency_limit( - db, - pull_loop_count, - suspended, - pulled_job, - &concurrency_settings, - ), - ) - .await?? - { - return Ok(pulled_job_res); - } + if let Some(pulled_job_res) = admit_pulled_job(db, job, suspended, pull_loop_count).await? { + return Ok(pulled_job_res); } } } +/// Applies a just-claimed job's concurrency limit, if it has one. `None` means the job was over +/// its limit and has been re-queued for later, so the claiming worker should pull again. +async fn admit_pulled_job( + db: &Pool, + job: PulledJob, + suspended: bool, + pull_loop_count: i32, +) -> windmill_common::error::Result> { + let concurrency_settings = windmill_common::runnable_settings::prefetch_cached_from_handle( + job.runnable_settings_handle, + db, + ) + .await? + .1 + .maybe_fallback(None, job.concurrent_limit, job.concurrency_time_window_s); + + let has_concurent_limit = has_active_concurrency_limit(concurrency_settings.concurrent_limit); + + #[cfg(not(feature = "enterprise"))] + if has_concurent_limit && !job.is_dependency() { + tracing::error!("Concurrent limits are an EE feature only, ignoring constraints") + } + + #[cfg(not(feature = "enterprise"))] + let has_concurent_limit = job.is_dependency() + && has_active_concurrency_limit(job.concurrent_limit) + && cfg!(feature = "private") + && !*WMDEBUG_NO_DEBOUNCING; + // if we don't have private flag, we don't have concurrency limit + + // concurrency check. If more than X jobs for this path are already running, we re-queue and pull another job from the queue + let pulled_job = job; + if pulled_job.runnable_path.is_none() + || !has_concurent_limit + || pulled_job.canceled_by.is_some() + { + #[cfg(feature = "prometheus")] + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { + QUEUE_PULL_COUNT.inc(); + } + otel_incr_queue_pull_count(); + return Ok(Some(PulledJobResult { + job: Some(pulled_job), + suspended, + missing_concurrency_key: false, + error_while_preprocessing: None, + })); + } + + #[cfg(feature = "private")] + if cfg!(feature = "enterprise") || (pulled_job.is_dependency() && !*WMDEBUG_NO_DEBOUNCING) { + return Ok(timeout( + Duration::from_secs(15), + crate::jobs_ee::apply_concurrency_limit( + db, + pull_loop_count, + suspended, + pulled_job, + &concurrency_settings, + ), + ) + .await??); + } + #[cfg(not(feature = "private"))] + let _ = pull_loop_count; + Ok(None) +} + +/// Claims jobs for several waiting workers in one statement per queue walked, instead of one +/// pull per worker. Each distinct name in `worker_names` gets at most one job, marked running +/// under that name; a worker missing from the result found nothing to run. +/// +/// Walks the queues the way [`pull`] does for one worker: when `suspend_first`, suspended flows +/// ready to resume on any of the tags; then `tag_groups` from highest priority down, with the +/// workspace-fairness share of the workers first restricted to uncapped workspaces. A job over +/// its concurrency limit is re-queued for later and its worker is pulled for again. +/// +/// No authorization happens here: jobs from every workspace matching the tags are claimed and +/// returned with their args and permissions. Callers must authenticate each worker named and +/// pass only tags that worker is allowed to serve. +pub async fn pull_batch( + db: &Pool, + tag_groups: &[Vec], + worker_names: &[String], + suspend_first: bool, +) -> windmill_common::error::Result> { + let mut waiting: Vec = worker_names.iter().unique().cloned().collect(); + let mut admitted = vec![]; + for pull_loop_count in 1..=PULL_LOOP_LIMIT { + if waiting.is_empty() { + break; + } + let claimed = match claim_batch(db, tag_groups, &waiting, suspend_first).await + { + Ok(claimed) => claimed, + Err(e) if admitted.is_empty() => return Err(e), + // Jobs admitted by an earlier pass are running under their workers: return them. + Err(e) => { + tracing::error!( + "batch pull re-pull failed, returning the jobs admitted so far: {e:#}" + ); + break; + } + }; + let mut requeued = false; + for (worker_name, job, suspended) in claimed { + waiting.retain(|w| w != &worker_name); + let job_id = job.id; + match admit_pulled_job(db, job, suspended, pull_loop_count).await { + Ok(Some(res)) => admitted.push((worker_name, res)), + Ok(None) => { + waiting.push(worker_name); + requeued = true; + } + // The other claimed jobs are already running under their workers, so an error + // here must not abort the batch: hand this one back to the queue instead. + Err(e) => { + tracing::error!( + "error admitting batch-pulled job {job_id}, re-queuing it: {e:#}" + ); + if let Err(e) = sqlx::query( + "WITH ping AS ( + UPDATE v2_job_runtime SET ping = null WHERE id = $1 + ) + UPDATE v2_job_queue SET running = false, started_at = null + WHERE id = $1 AND worker = $2 AND running = true", + ) + .bind(job_id) + .bind(&worker_name) + .execute(db) + .await + { + tracing::error!("could not re-queue batch-pulled job {job_id}: {e:#}"); + } + } + } + } + if !requeued { + break; + } + if pull_loop_count == PULL_LOOP_LIMIT { + tracing::warn!( + "Batch pull loop count reached {PULL_LOOP_LIMIT}, backing off (likely concurrency re-queue storm)" + ); + } + } + Ok(admitted) +} + +/// One pass over the queues for `worker_names`, returning `(worker, job, suspended)` for each +/// job claimed. A failing statement ends the pass but keeps what was already claimed, since +/// those jobs are marked running and would otherwise be stranded until the zombie monitor. +async fn claim_batch( + db: &Pool, + tag_groups: &[Vec], + worker_names: &[String], + suspend_first: bool, +) -> windmill_common::error::Result> { + let mut claimed = vec![]; + let mut waiting = worker_names.to_vec(); + let pass: windmill_common::error::Result<()> = async { + let tags: Vec = tag_groups.iter().flatten().unique().cloned().collect(); + if suspend_first && !tags.is_empty() { + for (worker, job) in + claim_from_queue(db, &tags, PullQueue::Suspended, &[], &waiting).await? + { + waiting.retain(|w| w != &worker); + claimed.push((worker, job, true)); + } + } + + crate::workspace_fairness::maybe_refresh_overloaded(db); + let overloaded = WORKSPACE_FAIRNESS_OVERLOADED.load_full(); + if !overloaded.is_empty() { + // Same admission draw `pull` makes once per worker, made here once per waiting worker. + let mut fair: Vec = waiting + .iter() + .filter(|_| !crate::workspace_fairness::should_admit_capped()) + .cloned() + .collect(); + for tags in tag_groups.iter().filter(|t| !t.is_empty()) { + if fair.is_empty() { + break; + } + for (worker, job) in claim_from_queue( + db, + tags, + PullQueue::ReadyFairness, + overloaded.as_slice(), + &fair, + ) + .await? + { + fair.retain(|w| w != &worker); + waiting.retain(|w| w != &worker); + claimed.push((worker, job, false)); + } + } + } + + for tags in tag_groups.iter().filter(|t| !t.is_empty()) { + if waiting.is_empty() { + break; + } + for (worker, job) in claim_from_queue(db, tags, PullQueue::Ready, &[], &waiting).await? + { + waiting.retain(|w| w != &worker); + claimed.push((worker, job, false)); + } + } + Ok(()) + } + .await; + match pass { + Err(e) if claimed.is_empty() => Err(e), + Err(e) => { + tracing::error!( + "batch pull pass failed after claiming {} jobs: {e:#}", + claimed.len() + ); + Ok(claimed) + } + Ok(()) => Ok(claimed), + } +} + +async fn claim_from_queue( + db: &Pool, + tags: &[String], + queue: PullQueue, + overloaded: &[String], + worker_names: &[String], +) -> windmill_common::error::Result> { + use sqlx::{FromRow, Row}; + let query = make_batch_pull_query(tags, queue); + let mut query = sqlx::query(&query).bind(worker_names); + if queue == PullQueue::ReadyFairness { + query = query.bind(overloaded); + } + let rows = timeout(Duration::from_secs(15), query.fetch_all(db)).await??; + rows.iter() + .map(|row| Ok((row.try_get("assigned_worker")?, PulledJob::from_row(row)?))) + .collect() +} + async fn pull_single_job_and_mark_as_running_no_concurrency_limit<'c>( db: &Pool, suspend_first: bool,