From 4b09558e13fadca6fe604c231a5cdfb896d4d0f8 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Fri, 25 Sep 2026 12:43:46 +0200 Subject: [PATCH] feat: start a deferred queued job now without changing its id (#11347) * feat: start a deferred queued job now without changing its id Co-Authored-By: Claude Opus 5.5 (1M context) * fix: refuse starting a schedule's upcoming tick early Co-Authored-By: Claude Opus 5.5 (1M context) * fix: hide run now on upcoming schedule ticks and register its audit op Co-Authored-By: Claude Opus 5.5 (1M context) --------- Co-authored-by: Claude Opus 5.5 (1M context) --- ...751e1317e7398dfaf1dfde5c57169cc146c9d.json | 23 ++++ ...b30d71c74fffbba3f5ac6197c985fd73e5c11.json | 23 ++++ ...61af624ab6de67a5c82967f4df8117fe50bc7.json | 47 ++++++++ backend/tests/jobs_read_auth.rs | 93 ++++++++++++++- backend/windmill-api/openapi.yaml | 24 ++++ backend/windmill-api/src/jobs.rs | 112 +++++++++++++++++- .../auditLogs/AuditLogsFilters.svelte | 1 + .../(root)/(logged)/run/[...run]/+page.svelte | 36 +++++- 8 files changed, 354 insertions(+), 5 deletions(-) create mode 100644 backend/.sqlx/query-1740f8bd6f9e084a470a775151e751e1317e7398dfaf1dfde5c57169cc146c9d.json create mode 100644 backend/.sqlx/query-372db51a00e5cf37f00af4b7962b30d71c74fffbba3f5ac6197c985fd73e5c11.json create mode 100644 backend/.sqlx/query-c703da3c3e133536dda9721644561af624ab6de67a5c82967f4df8117fe50bc7.json diff --git a/backend/.sqlx/query-1740f8bd6f9e084a470a775151e751e1317e7398dfaf1dfde5c57169cc146c9d.json b/backend/.sqlx/query-1740f8bd6f9e084a470a775151e751e1317e7398dfaf1dfde5c57169cc146c9d.json new file mode 100644 index 0000000000..46b7d2509c --- /dev/null +++ b/backend/.sqlx/query-1740f8bd6f9e084a470a775151e751e1317e7398dfaf1dfde5c57169cc146c9d.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE v2_job_queue q SET scheduled_for = now()\n FROM (SELECT id, scheduled_for FROM v2_job_queue WHERE id = $1 AND workspace_id = $2 FOR UPDATE) prev\n WHERE q.id = prev.id AND NOT q.running AND q.suspend = 0 AND q.canceled_by IS NULL\n AND q.scheduled_for > now()\n RETURNING prev.scheduled_for", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "scheduled_for", + "type_info": "Timestamptz" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "1740f8bd6f9e084a470a775151e751e1317e7398dfaf1dfde5c57169cc146c9d" +} diff --git a/backend/.sqlx/query-372db51a00e5cf37f00af4b7962b30d71c74fffbba3f5ac6197c985fd73e5c11.json b/backend/.sqlx/query-372db51a00e5cf37f00af4b7962b30d71c74fffbba3f5ac6197c985fd73e5c11.json new file mode 100644 index 0000000000..f5a2ac20ed --- /dev/null +++ b/backend/.sqlx/query-372db51a00e5cf37f00af4b7962b30d71c74fffbba3f5ac6197c985fd73e5c11.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT EXISTS(SELECT 1 FROM v2_job_queue WHERE id = $1 AND workspace_id = $2)", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "exists", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "372db51a00e5cf37f00af4b7962b30d71c74fffbba3f5ac6197c985fd73e5c11" +} diff --git a/backend/.sqlx/query-c703da3c3e133536dda9721644561af624ab6de67a5c82967f4df8117fe50bc7.json b/backend/.sqlx/query-c703da3c3e133536dda9721644561af624ab6de67a5c82967f4df8117fe50bc7.json new file mode 100644 index 0000000000..4ed2afc15b --- /dev/null +++ b/backend/.sqlx/query-c703da3c3e133536dda9721644561af624ab6de67a5c82967f4df8117fe50bc7.json @@ -0,0 +1,47 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT q.scheduled_for, s.path, s.schedule, s.cron_version, s.timezone\n FROM v2_job j\n JOIN v2_job_queue q USING (id)\n JOIN schedule s ON s.workspace_id = j.workspace_id AND s.path = j.trigger\n WHERE j.id = $1 AND j.workspace_id = $2 AND j.trigger_kind = 'schedule'\n AND j.parent_job IS NULL AND q.scheduled_for > now()", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "scheduled_for", + "type_info": "Timestamptz" + }, + { + "ordinal": 1, + "name": "path", + "type_info": "Varchar" + }, + { + "ordinal": 2, + "name": "schedule", + "type_info": "Varchar" + }, + { + "ordinal": 3, + "name": "cron_version", + "type_info": "Text" + }, + { + "ordinal": 4, + "name": "timezone", + "type_info": "Varchar" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + false, + false, + false, + true, + false + ] + }, + "hash": "c703da3c3e133536dda9721644561af624ab6de67a5c82967f4df8117fe50bc7" +} diff --git a/backend/tests/jobs_read_auth.rs b/backend/tests/jobs_read_auth.rs index c8c154c250..8c5bfee0eb 100644 --- a/backend/tests/jobs_read_auth.rs +++ b/backend/tests/jobs_read_auth.rs @@ -516,10 +516,14 @@ async fn test_single_job_read_authorization(db: Pool) -> anyhow::Resul Some("SECRET_TOKEN_2"), ) .await; - assert!(status.is_success(), "owner must mint a resume secret: {secret}"); + assert!( + status.is_success(), + "owner must mint a resume secret: {secret}" + ); let secret = secret.trim().trim_matches('"').to_string(); - let approval_result = - format!("completed/get_result/{STEP_JOB}?suspended_job={STEP_JOB}&resume_id=0&secret={secret}"); + let approval_result = format!( + "completed/get_result/{STEP_JOB}?suspended_job={STEP_JOB}&resume_id=0&secret={secret}" + ); let (status, body) = get(&base, &approval_result, None).await; assert!( status.is_success(), @@ -991,6 +995,89 @@ async fn test_single_job_read_authorization(db: Pool) -> anyhow::Resul !status.is_success(), "a public token must not let an anonymous caller cancel the run (got {status}): {body}" ); + // ---- RUN_NOW is gated like cancel, and only moves a job still waiting for a + // future start. ---- + let run_now = format!("queue/run_now/{RUNNING_JOB}"); + let (status, body) = post(&authed_base, &run_now, Some("SECRET_TOKEN_3")).await; + assert_eq!( + status, + reqwest::StatusCode::FORBIDDEN, + "viewer must not start another user's job early (got {status}): {body}" + ); + let (status, body) = post(&authed_base, &run_now, Some("SECRET_TOKEN_2")).await; + assert_eq!( + status, + reqwest::StatusCode::BAD_REQUEST, + "a running job has no future start to move (got {status}): {body}" + ); + sqlx::query( + "UPDATE v2_job_queue SET running = false, scheduled_for = now() + interval '1 hour' + WHERE id = $1::uuid", + ) + .bind(RUNNING_JOB) + .execute(&db) + .await?; + let (status, body) = post(&authed_base, &run_now, Some("SECRET_TOKEN_2")).await; + assert!( + status.is_success(), + "owner must start their deferred job now (got {status}): {body}" + ); + let due: bool = + sqlx::query_scalar("SELECT scheduled_for <= now() FROM v2_job_queue WHERE id = $1::uuid") + .bind(RUNNING_JOB) + .fetch_one(&db) + .await?; + assert!(due, "run_now must make the job due immediately"); + + // A schedule tick that is not due yet is refused: the schedule would queue that same + // tick again on completion and run twice. + sqlx::query( + "INSERT INTO schedule (workspace_id, path, edited_by, schedule, script_path, permissioned_as) + VALUES ('test-workspace', 'u/test-user-2/daily', 'test-user-2', '0 0 3 * * *', + 'u/test-user-2/running_secret', 'u/test-user-2')", + ) + .execute(&db) + .await?; + sqlx::query( + "UPDATE v2_job SET trigger_kind = 'schedule', trigger = 'u/test-user-2/daily' + WHERE id = $1::uuid", + ) + .bind(RUNNING_JOB) + .execute(&db) + .await?; + sqlx::query( + "UPDATE v2_job_queue SET scheduled_for = + (date_trunc('day', now() AT TIME ZONE 'UTC') + interval '1 day 3 hours') AT TIME ZONE 'UTC' + WHERE id = $1::uuid", + ) + .bind(RUNNING_JOB) + .execute(&db) + .await?; + let (status, body) = post(&authed_base, &run_now, Some("SECRET_TOKEN_2")).await; + assert_eq!( + status, + reqwest::StatusCode::BAD_REQUEST, + "an upcoming schedule tick must not be started early (got {status}): {body}" + ); + // A tick a concurrency limit pushed past its time is off the cron occurrences and + // may still be started early. + sqlx::query( + "UPDATE v2_job_queue SET scheduled_for = now() + interval '1 hour 37 seconds 123 milliseconds' + WHERE id = $1::uuid", + ) + .bind(RUNNING_JOB) + .execute(&db) + .await?; + let (status, body) = post(&authed_base, &run_now, Some("SECRET_TOKEN_2")).await; + assert!( + status.is_success(), + "a deferred schedule tick must still start now (got {status}): {body}" + ); + sqlx::query("UPDATE v2_job SET trigger_kind = NULL, trigger = NULL WHERE id = $1::uuid") + .bind(RUNNING_JOB) + .execute(&db) + .await?; + // The owner still cancels their own job (no over-blocking). Keep this last: it // takes RUNNING_JOB out of the queue. let (status, body) = post( diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 4346dfcd0a..21ca421af3 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -16310,6 +16310,29 @@ paths: items: type: string + /w/{workspace}/jobs/queue/run_now/{id}: + post: + summary: start a queued job scheduled for later now + description: | + Moves the start of a queued job that is waiting for a future time (scheduled, debounced, + re-scheduled by a concurrency limit, or a flow step waiting on a sleep or a retry delay) + to now, keeping the same job id. A job with a + concurrency limit is still checked against its limit when it is pulled. A schedule's + upcoming tick is refused, since the schedule would queue that tick again. + operationId: runQueuedJobNow + tags: + - job + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/JobId" + responses: + "200": + description: job id + content: + text/plain: + schema: + type: string + /w/{workspace}/jobs/queue/cancel_selection: post: summary: cancel jobs based on the given uuids @@ -31177,6 +31200,7 @@ components: - "jobs" - "jobs.cancel" - "jobs.force_cancel" + - "jobs.run_now" - "jobs.disapproval" - "jobs.delete" - "account.delete" diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 052c9df4bb..b9e71271d4 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -116,7 +116,9 @@ use windmill_common::{ query_builders, scripts::{ScriptHash, ScriptLang}, users::username_to_permissioned_as, - utils::{not_found_if_none, now_from_db, paginate, require_admin, Pagination, StripPath}, + utils::{ + not_found_if_none, now_from_db, paginate, require_admin, Pagination, ScheduleType, StripPath, + }, }; use windmill_common::{ @@ -304,6 +306,7 @@ pub fn workspaced_service() -> Router { .route("/queue/position/{timestamp}", get(get_queue_position)) .route("/queue/scheduled_for/{id}", get(get_scheduled_for)) .route("/queue/cancel_selection", post(cancel_selection)) + .route("/queue/run_now/{id}", post(run_queued_job_now)) .route("/completed/count", get(count_completed_jobs)) .route("/completed/count_jobs", get(count_completed_jobs_detail)) .route( @@ -835,6 +838,113 @@ async fn force_cancel( } } +/// Moves a queued job's `scheduled_for` to now so the next pull picks it up, keeping its id. +/// A job with a concurrency limit still goes through the usual check at pull time and is +/// re-scheduled again if its key is still saturated. +async fn run_queued_job_now( + authed: ApiAuthed, + Extension(db): Extension, + Extension(user_db): Extension, + Path((w_id, id)): Path<(String, Uuid)>, +) -> error::Result { + require_job_update_read_access(&db, &user_db, &authed, &w_id, &id, None).await?; + refuse_upcoming_schedule_tick(&db, &w_id, id).await?; + + let mut tx = db.begin().await?; + let previous = sqlx::query_scalar!( + "UPDATE v2_job_queue q SET scheduled_for = now() + FROM (SELECT id, scheduled_for FROM v2_job_queue WHERE id = $1 AND workspace_id = $2 FOR UPDATE) prev + WHERE q.id = prev.id AND NOT q.running AND q.suspend = 0 AND q.canceled_by IS NULL + AND q.scheduled_for > now() + RETURNING prev.scheduled_for", + id, + w_id, + ) + .fetch_optional(&mut *tx) + .await?; + + let Some(previous) = previous else { + tx.commit().await?; + return Err( + if sqlx::query_scalar!( + "SELECT EXISTS(SELECT 1 FROM v2_job_queue WHERE id = $1 AND workspace_id = $2)", + id, + w_id + ) + .fetch_one(&db) + .await? + .unwrap_or(false) + { + Error::BadRequest(format!( + "job {id} is not waiting for a future start: it is running, suspended, \ + canceled or already due" + )) + } else { + Error::NotFound(format!("queued job id {id} does not exist")) + }, + ); + }; + + audit_log( + &mut *tx, + &authed, + "jobs.run_now", + ActionKind::Update, + &w_id, + Some(&id.to_string()), + Some([("previous_scheduled_for", previous.to_rfc3339().as_str())].into()), + ) + .await?; + tx.commit().await?; + + windmill_queue::append_logs( + &id, + &w_id, + format!( + "\nStart moved from {previous} to now by {}\n", + authed.display_username() + ), + &db.clone().into(), + ) + .await; + + Ok(id.to_string()) +} + +/// A schedule queues its next tick when the current one completes, computed from the +/// completion time. Starting a tick that is not yet due would therefore queue that same +/// tick again and run the schedule twice. A tick pushed past its time by a concurrency +/// limit no longer sits on a cron occurrence, so it may still be started early. +async fn refuse_upcoming_schedule_tick(db: &DB, w_id: &str, id: Uuid) -> error::Result<()> { + let Some(tick) = sqlx::query!( + "SELECT q.scheduled_for, s.path, s.schedule, s.cron_version, s.timezone + FROM v2_job j + JOIN v2_job_queue q USING (id) + JOIN schedule s ON s.workspace_id = j.workspace_id AND s.path = j.trigger + WHERE j.id = $1 AND j.workspace_id = $2 AND j.trigger_kind = 'schedule' + AND j.parent_job IS NULL AND q.scheduled_for > now()", + id, + w_id, + ) + .fetch_optional(db) + .await? + else { + return Ok(()); + }; + let sched = ScheduleType::from_str(&tick.schedule, tick.cron_version.as_deref(), false)?; + let tz = + chrono_tz::Tz::from_str(&tick.timezone).map_err(|e| Error::BadRequest(e.to_string()))?; + let just_before = (tick.scheduled_for - chrono::Duration::milliseconds(1)).with_timezone(&tz); + if sched.find_next(&just_before)? == tick.scheduled_for { + return Err(Error::BadRequest(format!( + "job {id} is the upcoming tick of schedule {}; run the schedule's script or flow \ + directly instead", + tick.path + ))); + } + Ok(()) +} + #[derive(Serialize)] struct QueuePosition { position: Option, diff --git a/frontend/src/lib/components/auditLogs/AuditLogsFilters.svelte b/frontend/src/lib/components/auditLogs/AuditLogsFilters.svelte index 96d341dc72..4090a39dbb 100644 --- a/frontend/src/lib/components/auditLogs/AuditLogsFilters.svelte +++ b/frontend/src/lib/components/auditLogs/AuditLogsFilters.svelte @@ -149,6 +149,7 @@ JOBS: 'jobs', JOBS_CANCEL: 'jobs.cancel', JOBS_FORCE_CANCEL: 'jobs.force_cancel', + JOBS_RUN_NOW: 'jobs.run_now', JOBS_DISAPPROVAL: 'jobs.disapproval', JOBS_DELETE: 'jobs.delete', JOBS_SHARE_PUBLICLY: 'jobs.share_publicly', diff --git a/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte b/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte index 932a76ae94..5fc7c1209d 100644 --- a/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte @@ -40,8 +40,10 @@ EllipsisVertical, Share2, Globe, - Users + Users, + Play } from 'lucide-svelte' + import { forLater } from '$lib/forLater' import { isJobResolvable } from '$lib/utils' import { @@ -280,6 +282,27 @@ } } + // A schedule's upcoming tick sits on a whole second and the backend refuses it; a tick + // deferred by a concurrency limit lands on a sub-second instant and can still start now. + let canRunNow = $derived( + job?.type === 'QueuedJob' && + !job.running && + !job.suspend && + !!job.scheduled_for && + forLater(job.scheduled_for) && + !(job.schedule_path && new Date(job.scheduled_for).getMilliseconds() === 0) + ) + + async function runJobNow(id: string) { + try { + await JobService.runQueuedJobNow({ workspace: $workspaceStore!, id }) + sendUserToast(`job ${id} will start as soon as a worker is available`) + getJob() + } catch (err) { + sendUserToast(`could not start job now: ${err?.body ?? err}`, true) + } + } + // Initialize view tab to logs since result is now outside tabs function initView(): void { // Result is now displayed outside tabs, so always default to logs @@ -886,6 +909,17 @@ Current runs {/if} + {#if canRunNow} + + {/if} {#if job && job?.type != 'CompletedJob' && (!job?.schedule_path || job?.['running'] == true)} {#if !forceCancel}