mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-03 16:02:12 +00:00
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) <noreply@anthropic.com> * fix: refuse starting a schedule's upcoming tick early Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * fix: hide run now on upcoming schedule ticks and register its audit op 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:
co-authored by
Claude Opus 5.5
parent
245628210f
commit
4b09558e13
+23
@@ -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"
|
||||
}
|
||||
+23
@@ -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"
|
||||
}
|
||||
+47
@@ -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"
|
||||
}
|
||||
@@ -516,10 +516,14 @@ async fn test_single_job_read_authorization(db: Pool<Postgres>) -> 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<Postgres>) -> 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(
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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<DB>,
|
||||
Extension(user_db): Extension<UserDB>,
|
||||
Path((w_id, id)): Path<(String, Uuid)>,
|
||||
) -> error::Result<String> {
|
||||
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<i64>,
|
||||
|
||||
@@ -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',
|
||||
|
||||
@@ -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
|
||||
</Button>
|
||||
{/if}
|
||||
{#if canRunNow}
|
||||
<Button
|
||||
unifiedSize="md"
|
||||
variant="default"
|
||||
startIcon={{ icon: Play }}
|
||||
on:click={() => job?.id && runJobNow(job.id)}
|
||||
title="Start this job now instead of at its scheduled time, skipping any remaining delay or sleep. It keeps the same id, and a concurrency limit is still enforced."
|
||||
>
|
||||
Run now
|
||||
</Button>
|
||||
{/if}
|
||||
{#if job && job?.type != 'CompletedJob' && (!job?.schedule_path || job?.['running'] == true)}
|
||||
{#if !forceCancel}
|
||||
<Button
|
||||
|
||||
Reference in New Issue
Block a user