diff --git a/README.md b/README.md index cdbebee661..0afd99a4b5 100644 --- a/README.md +++ b/README.md @@ -312,6 +312,8 @@ upcoming CLI tool. | PIP_LOCAL_DEPENDENCIES | None | Specify dependencies that are installed locally and do not need to be solved nor installed again | | ADDITIONAL_PYTHON_PATHS | None | Specify python paths (separated by a :) to be appended to the PYTHONPATH of the python jobs. To be used with PIP_LOCAL_DEPENDENCIES to use python codebases within Windmill | Worker | | INCLUDE_HEADERS | None | Whitelist of headers that are passed to jobs as args (separated by a comma) | Server | +| WHITELIST_WORKSPACES | None | Whitelist of workspaces this worker takes job from | Worker | +| BLACKLIST_WORKSPACES | None | Blacklist of workspaces this worker takes job from | Worker | diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index c0f267aaa4..b0d115ffe8 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -82,7 +82,32 @@ pub async fn cancel_job<'c>( Ok((tx, job_option)) } -pub async fn pull(db: &Pool) -> windmill_common::error::Result> { +pub async fn pull( + db: &Pool, + whitelist_workspaces: Option>, + blacklist_workspaces: Option>, +) -> windmill_common::error::Result> { + let mut workspaces_filter = String::new(); + if let Some(whitelist) = whitelist_workspaces { + workspaces_filter.push_str(&format!( + " AND workspace_id IN ({})", + whitelist + .into_iter() + .map(|x| format!("'{x}'")) + .collect::>() + .join(",") + )); + } + if let Some(blacklist) = blacklist_workspaces { + workspaces_filter.push_str(&format!( + " AND workspace_id NOT IN ({})", + blacklist + .into_iter() + .map(|x| format!("'{x}'")) + .collect::>() + .join(",") + )); + } /* Jobs can be started if they: * - haven't been started before, * running = false @@ -90,7 +115,7 @@ pub async fn pull(db: &Pool) -> windmill_common::error::Result = sqlx::query_as::<_, QueuedJob>( + let job: Option = sqlx::query_as::<_, QueuedJob>(&format!( "UPDATE queue SET running = true , started_at = coalesce(started_at, now()) @@ -99,17 +124,17 @@ pub async fn pull(db: &Pool) -> windmill_common::error::Result { + (job, timer) = {let timer = worker_pull_duration.start_timer(); pull(&db, whitelist_workspaces.clone(), blacklist_workspaces.clone()).map(|x| (x, timer)) } => { drop(timer); (false, job) },