mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-04 00:02:17 +00:00
* docs: plan for detecting skipped schedule occurrences Design plan only, no implementation. Records the scheduler's re-anchoring behaviour, the measurements behind it, and the three-piece design that came out of reviewing the alternatives. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * docs: state the user-facing outcome in the schedule plan The plan described the mechanism but never what a user would see, which made it hard to judge what the work is worth. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * docs: state which cause the schedule plan catches, and correct its scope Records which of the two causes each piece covers, and corrects the overrun scope: a script schedule carrying retry or dynamic_skip is pushed as a SingleStepFlow, so it re-arms at step 0 entry and its occurrences overlap like a flow's. Resolves the no_flow_overlap question, splits the read-time work into bounded detection and editor-only counting behind measured croner costs, and fixes the delivery order. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: count the occurrences a schedule skipped A schedule that overruns its interval, or waits for a worker, silently loses the occurrences in between: the scheduler keeps one queued occurrence and re-anchors on the clock, so nothing records that a run was due and never happened. Recovers the sequence from rows that already exist rather than writing per occurrence. `push_scheduled_job` anchors on `now_from_db` inside the transaction that inserts the job, and `v2_job.created_at` defaults to that same transaction timestamp, so `scheduled_for = find_next(created_at)` holds exactly and the whole occurrence history is derivable. The schedules list reports how many of the recent runs were followed by a lost occurrence, and a new occurrences endpoint carries the per-run wait and duration behind it. Detection is one `find_next` per gap, which stays bounded on a full page; counting walks the gap and runs only for a single schedule. The one write is `occurrence_baseline_at`, advanced at create, edit, re-enable and re-arm. Gaps older than it span a pause, a cron change, a re-enable or a reconciler re-arm, none of which mean runs were lost. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: show the wait and run time behind a schedule's skipped occurrences The list badge says a schedule is losing runs; this says which of the two causes did it. A large wait means not enough workers, a long run means the job outgrew its interval, and the pair is what tells them apart. Sits under the existing upcoming-events panel, so due and overdue read together. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: flag a schedule that is running late right now Reconstruction is retrospective: a gap only appears once the next occurrence has a row, which needs the current one to finish. A schedule wedged mid-run shows nothing until it moves, which is the case an operator most wants to see. An occurrence still in flight past the time its own successor was due will cost that successor, so `now > find_next(scheduled_for)` is the signal, needing no threshold and self-calibrating across a daily and a per-minute schedule. It applies only where occurrences serialize; an overlapping schedule starts its successor on time and would flag constantly while healthy. The queue is read in one aggregating pass keyed on (trigger, runnable_path) rather than a subquery per schedule, and an overlapping schedule holds more than one root row, hence the aggregate. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: run the schedule overrun alert from the monitor pass Wires `schedule_overrun_alerts` in next to `jobs_waiting_alerts`, every 30 iterations (~5 min). Its Enterprise implementation lives in windmill-labs/windmill-ee-private#772; only the wiring and the OSS stub are here. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * chore: refresh the sqlx offline cache Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: record and alert when a schedule skips occurrences push_scheduled_job compares each chained occurrence with the slot after the previous one. A gap is written to schedule.skipped_occurrences off the push transaction, alerts once when a clean schedule starts skipping, and recovers on the next clean chain. The schedules list shows a badge. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: alert only on a streak of skipping runs, keep a recent skip visible The skip state now describes the current streak and is written in the push transaction, so it commits or rolls back with the push. The alert fires once when 3 runs in a row skipped, and the list keeps a muted badge for 7 days after the latest skip. Editing or toggling a schedule resets it. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * fix: name the missed-occurrence state after what it counts, alert only once committed Renames the columns to late_run_streak, missed_occurrences and last_missed_at, keeps the missed count after a streak ends so the muted badge can show it, and rewords both badges. The alert task now reads the streak FOR SHARE, which waits for the push transaction, so a push that rolls back and retries alerts once. A failed slot count leaves the streak untouched, and a schedule deleted mid-push no longer fails it. Adds an integration test for the streak and its reset. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * fix: recover the late run alert, store the missed slot, name it missed throughout The alert now recovers (and so acknowledges itself) when a streak that alerted ends on a run on time, under the schedule:{path} resource used by the other trigger alerts. last_missed_at records the last missed cron slot rather than when the late run chained, and the counting helpers say missed, since skipped already names occurrences queued and not run. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * fix: scope the late run alert to its workspace, acknowledge it on edit, toggle and delete Recovery acknowledges alerts by resource alone, so the resource now carries the workspace. Editing, toggling or deleting a schedule clears its streak and a disabled or deleted one never chains a run on time, so those handlers acknowledge its open alert after committing. Past the 1000-slot cap, last_missed_at falls back to the detection time. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * refactor: raise the late run alert like the other critical alerts Drops the recovery, the workspace-scoped resource and the acknowledgement on edit, toggle and delete: the alert now fires once per streak with no resource and is acknowledged from the alerts feed, as the trigger and job failure alerts are. The FOR SHARE read stays, so a push that rolls back across the flow path's retries still alerts once. Notes in openapi that past 1000 misses in one late run the count is a floor and the time approximate. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
942 lines
34 KiB
Rust
942 lines
34 KiB
Rust
/*
|
|
* Author: Ruben Fiszel
|
|
* Copyright: Windmill Labs, Inc 2022
|
|
* This file and its contents are licensed under the AGPLv3 License.
|
|
* Please see the included NOTICE for copyright information and
|
|
* LICENSE-AGPL for a copy of the license.
|
|
*/
|
|
|
|
use crate::jobs::HTTP_CLIENT;
|
|
use crate::push;
|
|
use crate::PushIsolationLevel;
|
|
use anyhow::Context;
|
|
use chrono::DateTime;
|
|
use chrono::Utc;
|
|
use sqlx::{PgExecutor, Postgres, Transaction};
|
|
use std::collections::HashMap;
|
|
use std::str::FromStr;
|
|
use windmill_common::db::Authed;
|
|
use windmill_common::ee_oss::LICENSE_KEY_VALID;
|
|
use windmill_common::flows::Retry;
|
|
use windmill_common::get_flow_version_info_from_version;
|
|
use windmill_common::get_latest_flow_version_id_for_path;
|
|
use windmill_common::jobs::JobPayload;
|
|
use windmill_common::jobs::JobTriggerKind;
|
|
use windmill_common::jobs::OnBehalfOf;
|
|
use windmill_common::runnable_settings::ConcurrencySettings;
|
|
use windmill_common::runnable_settings::DebouncingSettings;
|
|
use windmill_common::schedule::schedule_to_user;
|
|
use windmill_common::scripts::get_full_hub_script_by_path;
|
|
use windmill_common::scripts::ScriptHash;
|
|
use windmill_common::triggers::TriggerMetadata;
|
|
use windmill_common::utils::WarnAfterExt;
|
|
use windmill_common::worker::to_raw_value;
|
|
use windmill_common::FlowVersionInfo;
|
|
use windmill_common::DB;
|
|
use windmill_common::{
|
|
error::{self, Result},
|
|
schedule::Schedule,
|
|
utils::{now_from_db, report_critical_error, ScheduleType, StripPath},
|
|
};
|
|
|
|
/// Helper to fetch metadata for a schedule's script or flow
|
|
async fn get_schedule_metadata<'c>(
|
|
tx: &mut sqlx::Transaction<'c, sqlx::Postgres>,
|
|
db: &DB,
|
|
schedule: &Schedule,
|
|
) -> Result<(
|
|
Option<String>, // tag
|
|
Option<i32>, // timeout
|
|
Option<OnBehalfOf>, // identity the runnable is deployed to run as
|
|
Option<ScriptHash>, // hash (for scripts)
|
|
Option<i64>, // flow_version (for flows)
|
|
Option<Retry>, // retry
|
|
)> {
|
|
let parsed_retry = schedule
|
|
.retry
|
|
.clone()
|
|
.and_then(|r| serde_json::from_value::<Retry>(r).ok());
|
|
|
|
if schedule.is_flow {
|
|
let version = get_latest_flow_version_id_for_path(
|
|
None,
|
|
&mut **tx,
|
|
&schedule.workspace_id,
|
|
&schedule.script_path,
|
|
false,
|
|
)
|
|
.await?;
|
|
|
|
let flow_info = get_flow_version_info_from_version(
|
|
&mut **tx,
|
|
version,
|
|
&schedule.workspace_id,
|
|
&schedule.script_path,
|
|
)
|
|
.await?;
|
|
|
|
Ok((
|
|
flow_info.tag.clone(),
|
|
None,
|
|
flow_info.on_behalf_of(&schedule.workspace_id, db).await?,
|
|
None,
|
|
Some(version),
|
|
parsed_retry,
|
|
))
|
|
} else if schedule.script_path.starts_with("hub/") {
|
|
Ok((None, None, None, None, None, parsed_retry))
|
|
} else {
|
|
let (
|
|
hash,
|
|
tag,
|
|
_custom_concurrency_key,
|
|
_concurrent_limit,
|
|
_concurrency_time_window_s,
|
|
_debounce_key,
|
|
_debounce_delay_s,
|
|
_cache_ttl,
|
|
_cache_ignore_s3_path,
|
|
_language,
|
|
_dedicated_worker,
|
|
_priority,
|
|
timeout,
|
|
on_behalf_of,
|
|
_runnable_settings_handle,
|
|
_labels,
|
|
) = windmill_common::get_latest_hash_for_path(
|
|
&mut **tx,
|
|
db,
|
|
&schedule.workspace_id,
|
|
&schedule.script_path,
|
|
false,
|
|
)
|
|
.await?;
|
|
|
|
Ok((tag, timeout, on_behalf_of, Some(hash), None, parsed_retry))
|
|
}
|
|
}
|
|
|
|
pub async fn push_scheduled_job<'c>(
|
|
db: &DB,
|
|
mut tx: Transaction<'c, Postgres>,
|
|
schedule: &Schedule,
|
|
authed: Option<&Authed>,
|
|
now_cutoff: Option<DateTime<Utc>>,
|
|
) -> Result<Transaction<'c, Postgres>> {
|
|
if !LICENSE_KEY_VALID.load(std::sync::atomic::Ordering::Relaxed) {
|
|
return Err(error::Error::BadRequest(
|
|
"License key is not valid. Go to your superadmin settings to update your license key."
|
|
.to_string(),
|
|
));
|
|
}
|
|
|
|
let sched =
|
|
ScheduleType::from_str(&schedule.schedule, schedule.cron_version.as_deref(), false)?;
|
|
|
|
let tz = chrono_tz::Tz::from_str(&schedule.timezone)
|
|
.map_err(|e| error::Error::BadRequest(e.to_string()))?;
|
|
|
|
let now = now_from_db(&mut *tx).await?;
|
|
|
|
let now = match now_cutoff {
|
|
Some(now_cutoff) if now_cutoff >= now => {
|
|
tracing::error!(
|
|
"now_cutoff ({:?}) is after now ({:?}) for schedule {}. Using now_cutoff + 1s. This likely means the pg clock was shifted backwards.",
|
|
now_cutoff,
|
|
now,
|
|
&schedule.path
|
|
);
|
|
now_cutoff + chrono::Duration::seconds(1)
|
|
}
|
|
_ => now,
|
|
};
|
|
|
|
let starting_from = match schedule.paused_until {
|
|
Some(paused_until) if paused_until > now => paused_until.with_timezone(&tz),
|
|
paused_until_o => {
|
|
if paused_until_o.is_some() {
|
|
sqlx::query!(
|
|
"UPDATE schedule SET paused_until = NULL WHERE workspace_id = $1 AND path = $2",
|
|
&schedule.workspace_id,
|
|
&schedule.path
|
|
)
|
|
.execute(&mut *tx)
|
|
.warn_after_seconds_with_sql(1, "update_schedule_paused_until".to_string())
|
|
.await
|
|
.context("Failed to clear paused_until for schedule")?;
|
|
}
|
|
now.with_timezone(&tz)
|
|
}
|
|
};
|
|
|
|
let next = sched.find_next(&starting_from)?;
|
|
|
|
// Scheduled events must be stored in the database in UTC
|
|
let next = next.with_timezone(&chrono::Utc);
|
|
let already_exists: bool = sqlx::query_scalar!(
|
|
// Query plan:
|
|
// - use of the `ix_v2_job_root_by_path` index; hence the `parent_job IS NULL` clause.
|
|
// - select from `v2_job` first, then join with `v2_job_queue` to avoid a full table scan
|
|
// on `scheduled_for = $3`.
|
|
"SELECT EXISTS (
|
|
SELECT 1 FROM v2_job j JOIN v2_job_queue USING (id)
|
|
WHERE j.workspace_id = $1 AND trigger_kind = 'schedule' AND trigger = $2 AND runnable_path = $4
|
|
AND parent_job IS NULL
|
|
AND scheduled_for = $3
|
|
)",
|
|
&schedule.workspace_id,
|
|
&schedule.path,
|
|
next,
|
|
&schedule.script_path
|
|
)
|
|
.fetch_one(&mut *tx)
|
|
.warn_after_seconds_with_sql(1, "already_exists_job".to_string())
|
|
.await?
|
|
.unwrap_or(false);
|
|
|
|
if already_exists {
|
|
tracing::warn!(
|
|
"Job for schedule {} at {} already exists",
|
|
&schedule.path,
|
|
next
|
|
);
|
|
return Ok(tx);
|
|
}
|
|
|
|
// Only a chained push (`now_cutoff` is the previous occurrence) can miss one: create,
|
|
// edit, enable and re-arm start fresh, and a pause, even one already over, is deliberate.
|
|
let missed = now_cutoff
|
|
.filter(|_| schedule.paused_until.is_none())
|
|
.and_then(|prev| {
|
|
// A failed count leaves the streak as it was rather than reading as a run on time.
|
|
count_missed_occurrences(&sched, &tz, prev, next)
|
|
.inspect_err(|e| {
|
|
tracing::warn!(
|
|
"failed to count the occurrences schedule {} missed: {e}",
|
|
&schedule.path
|
|
)
|
|
})
|
|
.ok()
|
|
});
|
|
if let Some((missed, last_missed)) = missed {
|
|
// Past the cap the walk stops short of the last miss; the time it was detected is at
|
|
// most one period after it.
|
|
let last_missed = last_missed.map(|l| if missed == MAX_COUNTED_MISSES { now } else { l });
|
|
if let Some(last_missed) = last_missed {
|
|
let streak = sqlx::query!(
|
|
"UPDATE schedule SET late_run_streak = late_run_streak + 1,
|
|
missed_occurrences = CASE WHEN late_run_streak = 0 THEN $3
|
|
ELSE missed_occurrences + $3 END,
|
|
last_missed_at = $4
|
|
WHERE workspace_id = $1 AND path = $2
|
|
RETURNING late_run_streak, missed_occurrences",
|
|
&schedule.workspace_id,
|
|
&schedule.path,
|
|
missed as i32,
|
|
last_missed,
|
|
)
|
|
.fetch_optional(&mut *tx)
|
|
.warn_after_seconds_with_sql(1, "update_schedule_late_run_streak".to_string())
|
|
.await?;
|
|
if let Some(streak) = streak.filter(|s| s.late_run_streak == LATE_RUNS_BEFORE_ALERT) {
|
|
tokio::spawn(alert_late_run_streak(
|
|
db.clone(),
|
|
schedule.workspace_id.clone(),
|
|
schedule.path.clone(),
|
|
streak.missed_occurrences,
|
|
));
|
|
}
|
|
} else if schedule.late_run_streak > 0 {
|
|
sqlx::query!(
|
|
"UPDATE schedule SET late_run_streak = 0
|
|
WHERE workspace_id = $1 AND path = $2",
|
|
&schedule.workspace_id,
|
|
&schedule.path,
|
|
)
|
|
.execute(&mut *tx)
|
|
.warn_after_seconds_with_sql(1, "reset_schedule_late_run_streak".to_string())
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
let mut args: HashMap<String, Box<serde_json::value::RawValue>> = HashMap::new();
|
|
|
|
if let Some(args_v) = &schedule.args {
|
|
if let Ok(args_m) =
|
|
serde_json::from_str::<HashMap<String, Box<serde_json::value::RawValue>>>(args_v.get())
|
|
{
|
|
args = args_m.clone()
|
|
} else {
|
|
return Err(error::Error::ExecutionErr(
|
|
"args of scripts needs to be dict".to_string(),
|
|
));
|
|
}
|
|
}
|
|
|
|
// Managed ducklake maintenance schedule (enterprise): the runnable is a
|
|
// generated DuckDB script, not a deployed one — built in the EE module.
|
|
// None (CE build, or no enabled maintenance config for the path's lake)
|
|
// falls through to normal script resolution, so a user schedule that
|
|
// pre-dates the reserved prefix keeps running its script and a stale
|
|
// managed row fails resolution with NotFound (auto-disabling it with
|
|
// schedule.error recorded).
|
|
let maintenance_payload =
|
|
if windmill_common::workspaces::lake_from_ducklake_maintenance_path(&schedule.path)
|
|
.is_some()
|
|
{
|
|
crate::ducklake_maintenance::build_maintenance_schedule_payload(&mut tx, schedule)
|
|
.await?
|
|
} else {
|
|
None
|
|
};
|
|
|
|
// If schedule handler is defined, wrap the scheduled job in a synthetic flow
|
|
// with the handler as the first step (with stop_after_if to skip if handler returns false)
|
|
let (payload, tag, timeout, on_behalf_of) = if let Some(maintenance_payload) =
|
|
maintenance_payload
|
|
{
|
|
maintenance_payload
|
|
} else if let Some(handler_path) = &schedule.dynamic_skip {
|
|
// Build skip handler args
|
|
let mut skip_handler_args = HashMap::<String, Box<serde_json::value::RawValue>>::new();
|
|
skip_handler_args.insert(
|
|
"scheduled_for".to_string(),
|
|
to_raw_value(&next.to_rfc3339()),
|
|
);
|
|
|
|
let stop_condition = "result !== true".to_string();
|
|
let stop_message = format!(
|
|
"Schedule handler {} did not return true for datetime {}. Handler must return boolean true to execute scheduled job.",
|
|
handler_path,
|
|
next.to_rfc3339()
|
|
);
|
|
|
|
// Get metadata from the scheduled script/flow for tag, timeout, etc.
|
|
let (tag, timeout, on_behalf_of, hash, flow_version, retry) =
|
|
get_schedule_metadata(&mut tx, db, schedule).await?;
|
|
|
|
(
|
|
JobPayload::SingleStepFlow {
|
|
path: schedule.script_path.clone(),
|
|
hash,
|
|
flow_version,
|
|
language: None,
|
|
args: args.clone(),
|
|
retry,
|
|
error_handler_path: None,
|
|
error_handler_args: None,
|
|
skip_handler: Some(windmill_common::jobs::SkipHandler {
|
|
path: handler_path.clone(),
|
|
args: skip_handler_args,
|
|
stop_condition,
|
|
stop_message,
|
|
}),
|
|
cache_ttl: None,
|
|
cache_ignore_s3_path: None,
|
|
priority: None,
|
|
tag_override: schedule.tag.clone(),
|
|
trigger_path: None,
|
|
apply_preprocessor: false,
|
|
concurrency_settings: ConcurrencySettings::default(),
|
|
debouncing_settings: DebouncingSettings::default(),
|
|
},
|
|
if schedule.tag.as_ref().is_some_and(|x| x != "") {
|
|
schedule.tag.clone()
|
|
} else {
|
|
tag
|
|
},
|
|
timeout,
|
|
on_behalf_of,
|
|
)
|
|
} else if schedule.is_flow {
|
|
let version = get_latest_flow_version_id_for_path(
|
|
None,
|
|
&mut *tx,
|
|
&schedule.workspace_id,
|
|
&schedule.script_path,
|
|
false,
|
|
)
|
|
.warn_after_seconds_with_sql(1, "get_latest_flow_version_id_for_path".to_string())
|
|
.await?;
|
|
|
|
let flow_info = get_flow_version_info_from_version(
|
|
&mut *tx,
|
|
version,
|
|
&schedule.workspace_id,
|
|
&schedule.script_path,
|
|
)
|
|
.warn_after_seconds_with_sql(1, "get_flow_version_info_from_version".to_string())
|
|
.await?;
|
|
let on_behalf_of = flow_info.on_behalf_of(&schedule.workspace_id, db).await?;
|
|
let FlowVersionInfo { version, tag, dedicated_worker, labels, .. } = flow_info;
|
|
|
|
(
|
|
JobPayload::Flow {
|
|
path: schedule.script_path.clone(),
|
|
dedicated_worker,
|
|
apply_preprocessor: false,
|
|
version,
|
|
labels,
|
|
},
|
|
tag,
|
|
None,
|
|
on_behalf_of,
|
|
)
|
|
} else if schedule.script_path.starts_with("hub/") {
|
|
let tag = schedule.tag.clone().filter(|t| !t.is_empty());
|
|
let payload = match &schedule.retry {
|
|
// The language is what lets `push` materialize this as a native retry
|
|
// instead of a flow wrapper, which would queue the next tick at start.
|
|
Some(retry) => JobPayload::SingleStepFlow {
|
|
path: schedule.script_path.clone(),
|
|
hash: None,
|
|
flow_version: None,
|
|
language: Some(
|
|
get_full_hub_script_by_path(
|
|
StripPath(schedule.script_path.clone()),
|
|
&HTTP_CLIENT,
|
|
Some(db),
|
|
)
|
|
.await?
|
|
.language,
|
|
),
|
|
retry: Some(serde_json::from_value::<Retry>(retry.clone()).map_err(|e| {
|
|
error::Error::internal_err(format!(
|
|
"Unable to parse retry information from schedule: {e}"
|
|
))
|
|
})?),
|
|
error_handler_path: None,
|
|
error_handler_args: None,
|
|
skip_handler: None,
|
|
args: args.clone(),
|
|
cache_ttl: None,
|
|
cache_ignore_s3_path: None,
|
|
priority: None,
|
|
tag_override: tag.clone(),
|
|
trigger_path: None,
|
|
apply_preprocessor: false,
|
|
concurrency_settings: ConcurrencySettings::default(),
|
|
debouncing_settings: DebouncingSettings::default(),
|
|
},
|
|
None => JobPayload::ScriptHub {
|
|
path: schedule.script_path.clone(),
|
|
apply_preprocessor: false,
|
|
},
|
|
};
|
|
(payload, tag, None, None)
|
|
} else {
|
|
let (
|
|
hash,
|
|
tag,
|
|
concurrency_key,
|
|
concurrent_limit,
|
|
concurrency_time_window_s,
|
|
debounce_key,
|
|
debounce_delay_s,
|
|
cache_ttl,
|
|
cache_ignore_s3_path,
|
|
language,
|
|
dedicated_worker,
|
|
priority,
|
|
timeout,
|
|
on_behalf_of,
|
|
runnable_settings_handle,
|
|
labels,
|
|
) = windmill_common::get_latest_hash_for_path(
|
|
&mut *tx,
|
|
db,
|
|
&schedule.workspace_id,
|
|
&schedule.script_path,
|
|
false,
|
|
)
|
|
.warn_after_seconds_with_sql(1, "get_latest_hash_for_path".to_string())
|
|
.await?;
|
|
|
|
// NB: read on the non-RLS pool (`db`), not `tx`. push_scheduled_job is
|
|
// also invoked with an RLS user_db transaction (api-schedule/api-flows),
|
|
// under which these lookups would resolve against the caller's row
|
|
// visibility rather than the full table. The dual-connection here is
|
|
// intentional and required for correctness.
|
|
let (debouncing_settings, concurrency_settings) =
|
|
windmill_common::runnable_settings::prefetch_cached_from_handle(
|
|
runnable_settings_handle,
|
|
db,
|
|
)
|
|
.await?;
|
|
|
|
if schedule.retry.is_some() {
|
|
let parsed_retry = serde_json::from_value::<Retry>(schedule.retry.clone().unwrap())
|
|
.map_err(|err| {
|
|
error::Error::internal_err(format!(
|
|
"Unable to parse retry information from schedule: {}",
|
|
err.to_string(),
|
|
))
|
|
})?;
|
|
let mut static_args = HashMap::<String, Box<serde_json::value::RawValue>>::new();
|
|
for (arg_name, arg_value) in args.clone() {
|
|
static_args.insert(arg_name, arg_value);
|
|
}
|
|
// A retry on a scheduled script is materialized into a native retry
|
|
// (see `push`): `Some(language)` opts in. Completion handlers are
|
|
// driven from the terminal attempt, and the per-occurrence
|
|
// failure/recovery counting queries (apply_schedule_handlers) resolve
|
|
// terminal status across the retry chain — so on_failure/on_recovery
|
|
// (incl. multi-count/exact) are all handled. A `retry_if` gate is
|
|
// evaluated at failure time; on a worker built without quickjs it
|
|
// cannot be evaluated and fails closed (no retry).
|
|
(
|
|
JobPayload::SingleStepFlow {
|
|
path: schedule.script_path.clone(),
|
|
hash: Some(hash),
|
|
flow_version: None,
|
|
language: Some(language),
|
|
retry: Some(parsed_retry),
|
|
error_handler_path: None,
|
|
error_handler_args: None,
|
|
skip_handler: None,
|
|
args: static_args,
|
|
cache_ttl,
|
|
cache_ignore_s3_path,
|
|
priority,
|
|
tag_override: schedule.tag.clone(),
|
|
trigger_path: None,
|
|
apply_preprocessor: false,
|
|
// Carry the script's concurrency/debounce settings (fetched
|
|
// above) into the native retry materialization, so a retrying
|
|
// concurrency-limited scheduled script still inserts its
|
|
// concurrency_key instead of running unbounded.
|
|
concurrency_settings,
|
|
debouncing_settings,
|
|
},
|
|
if schedule.tag.as_ref().is_some_and(|x| x != "") {
|
|
schedule.tag.clone()
|
|
} else {
|
|
tag
|
|
},
|
|
timeout,
|
|
on_behalf_of.clone(),
|
|
)
|
|
} else {
|
|
(
|
|
JobPayload::ScriptHash {
|
|
hash,
|
|
path: schedule.script_path.clone(),
|
|
cache_ttl,
|
|
cache_ignore_s3_path,
|
|
dedicated_worker,
|
|
language,
|
|
priority,
|
|
apply_preprocessor: false,
|
|
debouncing_settings: debouncing_settings
|
|
.maybe_fallback(debounce_key, debounce_delay_s),
|
|
concurrency_settings: concurrency_settings.maybe_fallback(
|
|
concurrency_key,
|
|
concurrent_limit,
|
|
concurrency_time_window_s,
|
|
),
|
|
labels,
|
|
},
|
|
if schedule.tag.as_ref().is_some_and(|x| x != "") {
|
|
schedule.tag.clone()
|
|
} else {
|
|
tag
|
|
},
|
|
timeout,
|
|
on_behalf_of,
|
|
)
|
|
}
|
|
};
|
|
|
|
if let Err(e) = sqlx::query!(
|
|
"UPDATE schedule SET error = NULL WHERE workspace_id = $1 AND path = $2",
|
|
&schedule.workspace_id,
|
|
&schedule.path
|
|
)
|
|
.execute(&mut *tx)
|
|
.warn_after_seconds_with_sql(1, "clear_schedule_error".to_string())
|
|
.await
|
|
{
|
|
tracing::error!(
|
|
"Failed to clear error for schedule {}: {}",
|
|
&schedule.path,
|
|
e
|
|
);
|
|
};
|
|
|
|
let (email, permissioned_as, push_authed, revert_to_windmill_user) = if let Some(obo) =
|
|
on_behalf_of.as_ref()
|
|
{
|
|
let is_windmill_user =
|
|
sqlx::query_scalar!("SELECT CURRENT_USER = 'windmill_user' as \"is_windmill_user!\"")
|
|
.fetch_one(&mut *tx)
|
|
.warn_after_seconds_with_sql(1, "is_windmill_user".to_string())
|
|
.await?;
|
|
if is_windmill_user {
|
|
sqlx::query!("SET LOCAL ROLE NONE")
|
|
.execute(&mut *tx)
|
|
.warn_after_seconds_with_sql(1, "set_local_role_none".to_string())
|
|
.await?;
|
|
}
|
|
(
|
|
obo.email.clone(),
|
|
obo.permissioned_as.clone(),
|
|
None,
|
|
is_windmill_user,
|
|
)
|
|
} else {
|
|
let permissioned_as = schedule.permissioned_as.clone();
|
|
let resolved_email = windmill_common::users::get_email_from_permissioned_as(
|
|
&permissioned_as,
|
|
&schedule.workspace_id,
|
|
db,
|
|
)
|
|
.await?;
|
|
(resolved_email, permissioned_as, authed, false)
|
|
};
|
|
|
|
let obo_authed;
|
|
let push_authed = match push_authed {
|
|
Some(a) => Some(a),
|
|
None => {
|
|
obo_authed = windmill_common::auth::fetch_authed_from_permissioned_as(
|
|
&permissioned_as,
|
|
&email,
|
|
&schedule.workspace_id,
|
|
&mut *tx,
|
|
)
|
|
.await
|
|
.ok();
|
|
obo_authed.as_ref()
|
|
}
|
|
};
|
|
|
|
if let Some(tag) = tag.as_deref().filter(|t| !t.is_empty()) {
|
|
let is_super_admin = windmill_common::auth::is_super_admin_email(db, &email).await?;
|
|
crate::check_tag_available_for_push(
|
|
db,
|
|
&schedule.workspace_id,
|
|
&tag,
|
|
&crate::PushArgs::from(&args),
|
|
is_super_admin,
|
|
None, // no token for schedules so no scopes so no scope_tags
|
|
)
|
|
.warn_after_seconds_with_sql(1, "check_tag_available_for_push".to_string())
|
|
.await?;
|
|
}
|
|
|
|
tracing::info!(
|
|
"Pushing next scheduled job for schedule {} at {} (schedule: {})",
|
|
&schedule.path,
|
|
next,
|
|
&schedule.schedule
|
|
);
|
|
let tx = PushIsolationLevel::Transaction(tx);
|
|
let (_, mut tx) = push(
|
|
&db,
|
|
tx,
|
|
&schedule.workspace_id,
|
|
payload,
|
|
crate::PushArgs { args: &args, extra: None },
|
|
&schedule_to_user(&schedule.path),
|
|
&email,
|
|
permissioned_as,
|
|
Some(&schedule.path),
|
|
None,
|
|
Some(next),
|
|
Some(schedule.path.clone()),
|
|
None,
|
|
None,
|
|
None,
|
|
None,
|
|
false,
|
|
false,
|
|
None,
|
|
true,
|
|
tag,
|
|
timeout,
|
|
None,
|
|
None,
|
|
push_authed,
|
|
false,
|
|
None,
|
|
Some(TriggerMetadata::new(
|
|
Some(schedule.path.clone()),
|
|
JobTriggerKind::Schedule,
|
|
)),
|
|
None,
|
|
)
|
|
.warn_after_seconds_with_sql(1, "push in push_scheduled_job".to_string())
|
|
.await?;
|
|
|
|
if revert_to_windmill_user {
|
|
sqlx::query!("SET LOCAL ROLE windmill_user")
|
|
.execute(&mut *tx)
|
|
.warn_after_seconds_with_sql(1, "set_local_role_windmill_user".to_string())
|
|
.await?;
|
|
}
|
|
|
|
Ok(tx) // TODO: Bubble up pushed UUID from here
|
|
}
|
|
|
|
const MAX_COUNTED_MISSES: u32 = 1000;
|
|
|
|
/// Due slots of the cron strictly between the previous occurrence and the next one, and
|
|
/// the last of them. A chain on time costs one `find_next`: its first slot is `next` itself.
|
|
fn count_missed_occurrences(
|
|
sched: &ScheduleType,
|
|
tz: &chrono_tz::Tz,
|
|
prev: DateTime<Utc>,
|
|
next: DateTime<Utc>,
|
|
) -> Result<(u32, Option<DateTime<Utc>>)> {
|
|
let mut count = 0;
|
|
let mut last = None;
|
|
let mut slot = prev.with_timezone(tz);
|
|
while count < MAX_COUNTED_MISSES {
|
|
match sched.find_next(&slot) {
|
|
Ok(s) if s.with_timezone(&Utc) < next => {
|
|
count += 1;
|
|
last = Some(s.with_timezone(&Utc));
|
|
slot = s;
|
|
}
|
|
Ok(_) => break,
|
|
Err(e) => return Err(e),
|
|
}
|
|
}
|
|
Ok((count, last))
|
|
}
|
|
|
|
/// A single late run is a blip (a slow run, a worker restart, an edit mid-run); only a
|
|
/// streak alerts, once, when it reaches this length. The schedules list shows every one.
|
|
const LATE_RUNS_BEFORE_ALERT: i32 = 3;
|
|
|
|
/// Spawned: it reaches the instance alert channels, which must not hold up the push.
|
|
async fn alert_late_run_streak(db: DB, w_id: String, path: String, missed: i32) {
|
|
// The push transaction holds this row until it ends, so FOR SHARE waits for it: a push
|
|
// that rolled back leaves the streak short of the threshold, and only its committed
|
|
// retry alerts.
|
|
let committed = sqlx::query_scalar!(
|
|
"SELECT late_run_streak FROM schedule WHERE workspace_id = $1 AND path = $2 FOR SHARE",
|
|
&w_id,
|
|
&path,
|
|
)
|
|
.fetch_optional(&db)
|
|
.await;
|
|
match committed {
|
|
Ok(Some(streak)) if streak >= LATE_RUNS_BEFORE_ALERT => {}
|
|
Ok(_) => return,
|
|
Err(e) => {
|
|
tracing::warn!("failed to confirm the late run streak of schedule {w_id}/{path}: {e}");
|
|
return;
|
|
}
|
|
}
|
|
report_critical_error(
|
|
format!(
|
|
"Schedule {path} missed {missed} occurrences: its last {LATE_RUNS_BEFORE_ALERT} runs in a row started or finished too late"
|
|
),
|
|
db,
|
|
Some(&w_id),
|
|
None,
|
|
)
|
|
.await;
|
|
}
|
|
|
|
/// Enabled schedules with no occurrence in the queue, as `(workspace_id, path)`.
|
|
///
|
|
/// Every path that completes a scheduled job pushes the next occurrence in the
|
|
/// same transaction (for flows, on entry to step 0), so an enabled schedule
|
|
/// always has a queued occurrence — a run in progress is itself one. A run that
|
|
/// dies through an abnormal path can skip that push though, leaving the schedule
|
|
/// enabled yet dead until it is manually disabled and re-enabled. This is how the
|
|
/// monitor spots that state; see `rearm_schedule` for the recovery.
|
|
///
|
|
/// Not an authorization boundary: it reports schedules across every workspace, so
|
|
/// this is for system callers (the monitor's reconciliation pass) only and its
|
|
/// result must never be returned to a user unfiltered.
|
|
pub async fn find_unarmed_schedules(db: &DB) -> Result<Vec<(String, String)>> {
|
|
let rows = sqlx::query!(
|
|
// Query plan: the anti-join builds from `v2_job_queue` (only pending and
|
|
// running jobs) rather than probing `v2_job` once per schedule.
|
|
"SELECT s.workspace_id, s.path
|
|
FROM schedule s JOIN workspace w ON w.id = s.workspace_id AND NOT w.deleted
|
|
WHERE s.enabled IS TRUE
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM v2_job_queue q JOIN v2_job j USING (id)
|
|
WHERE j.workspace_id = s.workspace_id
|
|
AND j.trigger_kind = 'schedule'
|
|
AND j.trigger = s.path
|
|
AND j.runnable_path = s.script_path
|
|
AND j.parent_job IS NULL
|
|
)"
|
|
)
|
|
.fetch_all(db)
|
|
.await?;
|
|
Ok(rows.into_iter().map(|r| (r.workspace_id, r.path)).collect())
|
|
}
|
|
|
|
#[derive(Debug, PartialEq, Eq)]
|
|
pub enum RearmOutcome {
|
|
/// The next occurrence was pushed.
|
|
Rearmed,
|
|
/// Nothing to do: the schedule was deleted or disabled since it was found.
|
|
NoOp,
|
|
}
|
|
|
|
/// Push the next occurrence of a schedule that has none queued.
|
|
///
|
|
/// Only ever starts a schedule, never stops one: re-arming something that did not
|
|
/// need it costs one extra run, whereas wrongly disabling one is the silent
|
|
/// permanent stoppage this whole mechanism exists to prevent. So an occurrence
|
|
/// that cannot be pushed is logged and left alone — the schedule is already not
|
|
/// running, and `try_schedule_next_job` still disables on the completion path,
|
|
/// where the population is limited to actively-cycling schedules. Keep it that
|
|
/// way: this sweeps *every* enabled schedule, including ones broken long before
|
|
/// this code existed and never swept before.
|
|
///
|
|
/// Not an authorization boundary: it pushes under the schedule's own
|
|
/// `permissioned_as` identity for any `(w_id, path)`, so this is for system
|
|
/// callers (the monitor's reconciliation pass) only. A caller acting for a user
|
|
/// MUST already have enforced their permissions on `w_id` and `path`.
|
|
pub async fn rearm_schedule(db: &DB, w_id: &str, path: &str) -> Result<RearmOutcome> {
|
|
let mut tx = db.begin().await?;
|
|
// Lock the row for the whole push: an edit or a disable committing between the
|
|
// read and the push would otherwise leave a queued occurrence for a schedule
|
|
// that is disabled, or one built from superseded settings.
|
|
let schedule = sqlx::query_as::<_, Schedule>(
|
|
"SELECT workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, args, extra_perms, email, permissioned_as, error, on_failure, on_failure_times, on_failure_exact, on_failure_extra_args, on_recovery, on_recovery_times, on_recovery_extra_args, on_success, on_success_extra_args, ws_error_handler_muted, retry, no_flow_overlap, summary, description, tag, paused_until, cron_version, dynamic_skip, labels FROM schedule WHERE path = $1 AND workspace_id = $2 FOR UPDATE",
|
|
)
|
|
.bind(path)
|
|
.bind(w_id)
|
|
.fetch_optional(&mut *tx)
|
|
.await?;
|
|
let Some(schedule) = schedule else {
|
|
return Ok(RearmOutcome::NoOp);
|
|
};
|
|
if !schedule.enabled {
|
|
return Ok(RearmOutcome::NoOp);
|
|
}
|
|
// Re-check for a queued occurrence now that the row is locked: a normal
|
|
// completion, an edit, or a re-enable could have pushed one between the unarmed
|
|
// scan and this lock. push_scheduled_job only dedups the exact computed
|
|
// scheduled_for, so re-arming a schedule that has since become armed and crossed a
|
|
// cron boundary would queue a second root occurrence. Mirrors the anti-join in
|
|
// find_unarmed_schedules.
|
|
let already_armed: bool = sqlx::query_scalar(
|
|
"SELECT EXISTS (
|
|
SELECT 1 FROM v2_job_queue q JOIN v2_job j USING (id)
|
|
WHERE j.workspace_id = $1
|
|
AND j.trigger_kind = 'schedule'
|
|
AND j.trigger = $2
|
|
AND j.runnable_path = $3
|
|
AND j.parent_job IS NULL
|
|
)",
|
|
)
|
|
.bind(w_id)
|
|
.bind(path)
|
|
.bind(&schedule.script_path)
|
|
.fetch_one(&mut *tx)
|
|
.await?;
|
|
if already_armed {
|
|
return Ok(RearmOutcome::NoOp);
|
|
}
|
|
match push_scheduled_job(db, tx, &schedule, None, None).await {
|
|
Ok(tx) => {
|
|
tx.commit().await?;
|
|
Ok(RearmOutcome::Rearmed)
|
|
}
|
|
// An occurrence that can never be pushed (runnable gone, quota blown) is
|
|
// reported, not acted on — see the note above on why this never disables.
|
|
Err(err @ (error::Error::NotFound(_) | error::Error::QuotaExceeded(_))) => {
|
|
tracing::error!(
|
|
"Could not re-arm schedule {path} in {w_id}: {err}. Leaving it enabled; it will not run until the cause is fixed."
|
|
);
|
|
Ok(RearmOutcome::NoOp)
|
|
}
|
|
Err(err) => Err(err),
|
|
}
|
|
}
|
|
|
|
pub async fn get_schedule_opt<'c>(
|
|
e: impl PgExecutor<'c>,
|
|
w_id: &str,
|
|
path: &str,
|
|
) -> Result<Option<Schedule>> {
|
|
let schedule_opt = sqlx::query_as::<_, Schedule>(
|
|
"SELECT workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, args, extra_perms, email, permissioned_as, error, on_failure, on_failure_times, on_failure_exact, on_failure_extra_args, on_recovery, on_recovery_times, on_recovery_extra_args, on_success, on_success_extra_args, ws_error_handler_muted, retry, no_flow_overlap, summary, description, tag, paused_until, cron_version, dynamic_skip, labels, late_run_streak FROM schedule WHERE path = $1 AND workspace_id = $2",
|
|
)
|
|
.bind(path)
|
|
.bind(w_id)
|
|
.fetch_optional(e)
|
|
.await?;
|
|
Ok(schedule_opt)
|
|
}
|
|
|
|
pub async fn exists_schedule(
|
|
tx: &mut Transaction<'_, Postgres>,
|
|
w_id: String,
|
|
path: StripPath,
|
|
) -> Result<bool> {
|
|
let path = path.to_path();
|
|
|
|
let exists = sqlx::query_scalar!(
|
|
"SELECT EXISTS(SELECT 1 FROM schedule WHERE path = $1 AND workspace_id = $2)",
|
|
path,
|
|
w_id
|
|
)
|
|
.fetch_one(&mut **tx)
|
|
.await?
|
|
.unwrap_or(false);
|
|
|
|
Ok(exists)
|
|
}
|
|
|
|
pub async fn clear_schedule<'c>(
|
|
tx: &mut Transaction<'c, Postgres>,
|
|
path: &str,
|
|
w_id: &str,
|
|
) -> Result<()> {
|
|
tracing::info!("Clearing schedule {}", path);
|
|
// Delete the queued jobs (cascading their v2_job_queue-keyed side tables), then route the
|
|
// freed ids through delete_jobs so v2_job and its no-longer-cascading side tables go too.
|
|
let deleted_ids: Vec<uuid::Uuid> = sqlx::query_scalar!(
|
|
"WITH to_delete AS (
|
|
SELECT id FROM v2_job_queue
|
|
JOIN v2_job j USING (id)
|
|
WHERE trigger_kind = 'schedule'
|
|
AND trigger = $1
|
|
AND j.workspace_id = $2
|
|
AND flow_step_id IS NULL
|
|
AND running = false
|
|
FOR UPDATE
|
|
)
|
|
DELETE FROM v2_job_queue
|
|
WHERE id IN (SELECT id FROM to_delete)
|
|
RETURNING id",
|
|
path,
|
|
w_id
|
|
)
|
|
.fetch_all(&mut **tx)
|
|
.await?;
|
|
|
|
windmill_common::jobs::delete_jobs(&mut **tx, &deleted_ids).await?;
|
|
Ok(())
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn counts_the_slots_missed_between_two_occurrences() {
|
|
let every_30s = ScheduleType::from_str("*/30 * * * * *", None, false).unwrap();
|
|
let at = |s: &str| DateTime::parse_from_rfc3339(s).unwrap().with_timezone(&Utc);
|
|
let tz = chrono_tz::UTC;
|
|
let prev = at("2026-09-02T08:20:00Z");
|
|
let count = |next| count_missed_occurrences(&every_30s, &tz, prev, at(next)).unwrap();
|
|
assert_eq!(count("2026-09-02T08:20:30Z"), (0, None));
|
|
assert_eq!(
|
|
count("2026-09-02T08:21:30Z"),
|
|
(2, Some(at("2026-09-02T08:21:00Z")))
|
|
);
|
|
}
|
|
}
|