feat: add pull_batch to claim jobs for many waiting workers at once (#11350)

* feat: add pull_batch to claim jobs for many waiting workers at once

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: keep jobs admitted by earlier batch passes when a re-pull fails

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: keep suspended flows first on every batch re-pull pass

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-09-29 15:24:39 +02:00
committed by GitHub
co-authored by Claude Opus 5.5
parent ede4103b00
commit 9f40cdca62
3 changed files with 652 additions and 116 deletions
+278
View File
@@ -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<Postgres>, tag: &str, suspend: Option<i32>) -> anyhow::Result<Uuid> {
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<String> {
(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<Postgres>,
) -> 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<Uuid> = 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<String>) =
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<Postgres>) -> 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<Uuid> = 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<Uuid> = 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<Postgres>,
) -> 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<Postgres>) -> 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<Value>) {
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<Postgres>) -> 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(())
}
+113 -46
View File
@@ -774,6 +774,60 @@ impl From<&Pool<Postgres>> 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) {
+261 -70
View File
@@ -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<PulledJob>,
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<Postgres>,
job: PulledJob,
suspended: bool,
pull_loop_count: i32,
) -> windmill_common::error::Result<Option<PulledJobResult>> {
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<Postgres>,
tag_groups: &[Vec<String>],
worker_names: &[String],
suspend_first: bool,
) -> windmill_common::error::Result<Vec<(String, PulledJobResult)>> {
let mut waiting: Vec<String> = 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<Postgres>,
tag_groups: &[Vec<String>],
worker_names: &[String],
suspend_first: bool,
) -> windmill_common::error::Result<Vec<(String, PulledJob, bool)>> {
let mut claimed = vec![];
let mut waiting = worker_names.to_vec();
let pass: windmill_common::error::Result<()> = async {
let tags: Vec<String> = 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<String> = 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<Postgres>,
tags: &[String],
queue: PullQueue,
overloaded: &[String],
worker_names: &[String],
) -> windmill_common::error::Result<Vec<(String, PulledJob)>> {
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<Postgres>,
suspend_first: bool,