diff --git a/backend/tests/batch_pull.rs b/backend/tests/batch_pull.rs deleted file mode 100644 index 184b85dd06..0000000000 --- a/backend/tests/batch_pull.rs +++ /dev/null @@ -1,278 +0,0 @@ -//! 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 449fbd7606..d480a89851 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -774,60 +774,6 @@ 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 ( @@ -857,7 +803,16 @@ fn format_pull_query(peek: String) -> String { raw_flow, script_entrypoint_override, preprocessed FROM v2_job WHERE id = (SELECT id FROM peek) - ) SELECT {PULLED_JOB_COLUMNS} + ) 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 FROM q, j LEFT JOIN v2_job_status f USING (id) LEFT JOIN job_perms p ON p.job_id = j.id @@ -869,63 +824,22 @@ fn format_pull_query(peek: String) -> String { r } -/// 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), - }) -} - +// 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. pub fn make_suspended_pull_query(tags: &[String]) -> String { - format_pull_query(suspended_peek(tags, "1")) + 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(", ") + )) } // pub async fn make_suspended pub async fn store_suspended_pull_query(wc: &WorkerConfig) { @@ -938,7 +852,16 @@ pub async fn store_suspended_pull_query(wc: &WorkerConfig) { } pub fn make_pull_query(tags: &[String]) -> String { - format_pull_query(ready_peek(tags, None, "1")) + 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 } // Variant of `make_pull_query` that additionally excludes jobs whose workspace_id is in the @@ -949,7 +872,17 @@ 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 { - format_pull_query(ready_peek(tags, Some("$2"), "1")) + 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 } 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 8f65638790..7b552fb02c 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::{ - 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, + 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, }, DB, METRICS_ENABLED, }; @@ -3570,7 +3570,7 @@ impl MiniPulledJob { } } -#[derive(sqlx::FromRow, Debug, Clone, Serialize, Deserialize)] +#[derive(sqlx::FromRow, Debug, Clone)] pub struct PulledJob { #[sqlx(flatten)] pub job: MiniPulledJob, @@ -3837,7 +3837,7 @@ pub async fn get_queued_job_v2<'c>( Ok(job) } -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug)] pub struct PulledJobResult { pub job: Option, pub suspended: bool, @@ -4228,13 +4228,6 @@ 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( @@ -4249,9 +4242,14 @@ pub async fn pull( let mut pull_loop_count = 0; loop { pull_loop_count += 1; - if pull_loop_count > PULL_LOOP_LIMIT { + // 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 { tracing::warn!( - "Pull job loop count exceeded {PULL_LOOP_LIMIT}, backing off (likely concurrency re-queue storm)" + "Pull job loop count exceeded 10, backing off (likely concurrency re-queue storm)" ); return Ok(PulledJobResult { job: None, @@ -4391,255 +4389,66 @@ pub async fn pull( }); }; - 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, - ), + let concurrency_settings = windmill_common::runnable_settings::prefetch_cached_from_handle( + job.runnable_settings_handle, + db, ) - .await??); - } - #[cfg(not(feature = "private"))] - let _ = pull_loop_count; - Ok(None) -} + .await? + .1 + .maybe_fallback(None, job.concurrent_limit, job.concurrency_time_window_s); -/// 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 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") } - let claimed = match claim_batch(db, tag_groups, &waiting, suspend_first).await + + #[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() { - 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)); + #[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, + }); } - 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( + #[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, - 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? + pull_loop_count, + suspended, + pulled_job, + &concurrency_settings, + ), + ) + .await?? { - waiting.retain(|w| w != &worker); - claimed.push((worker, job, false)); + return Ok(pulled_job_res); } } - 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>(