Files
windmill/backend/windmill-api-schedule/src/lib.rs
T
d60dd745e4 feat(forks): handle triggers and schedules in workspace forks (#8976)
* feat(forks): strip operational state from triggers/schedules on git-sync export

When the source workspace is a fork (`wm-fork-*`), the tarball export now
omits `mode` from triggers and `enabled` from schedules. The trigger update
handler also preserves the existing DB `mode` when both fields are absent
from the request, instead of falling back to the BaseTriggerData default.

This prevents a fork's git-sync round-trip from flipping the parent
workspace's enabled/disabled state when a merge applies the fork's YAML
back to main.

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

* feat(forks): opt-in fork_triggers flag clones triggers/schedules disabled

Adds `workspace.fork_triggers` (default false) and a matching field on
CreateWorkspaceFork. When the user opts in, fork creation also runs
clone_triggers_and_schedules: every row in schedule and the ten
*_trigger tables is copied to the fork with mode='disabled' /
enabled=false. Listener identifiers (group_id, replication_slot_name,
subscription_name, …) are copied verbatim — the runtime suffix that
prevents the fork from competing with the parent ships in a follow-up
PR.

native_trigger is intentionally skipped: those triggers manage external
webhook state we don't want duplicated.

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

* feat(forks): warn before enabling triggers/schedules that conflict with parent

set_trigger_mode and schedule's set_enabled now check whether the parent
workspace has the same path actively enabled. If so, the call is rejected
with a `fork-conflict:<kind>:<parent_id>` error unless the request includes
`force=true`. The frontend interprets the prefix to surface a confirm-to-
proceed dialog.

This is the placeholder safety net until the Phase 3 listener-suffix work
removes the conflict for the namespaceable kinds (Kafka/MQTT/NATS/Postgres/
Azure/GCP-CreateNew). For SQS, GCP-Existing, and schedules — where there's
no namespacing fix — the warning is the durable solution.

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

* feat(forks): UI: opt-in clone-triggers checkbox + confirm-on-fork-conflict

Adds the user-facing surface for the fork-trigger work:

- CreateWorkspaceInner: new "Clone triggers and schedules" toggle in the
  fork-creation dialog (default off). Sends fork_triggers in the request.

- forkConflict utility: detects the `fork-conflict:<kind>:<parent_id>`
  error string from the backend, shows a confirm() dialog explaining
  why the action is blocked, retries with `force: true` if accepted.

- Wires withForkConflictRetry into every trigger setMode and the
  schedule setEnabled call, both in the per-kind editor components and
  the +page.svelte list views (HTTP, websocket, kafka, NATS, SQS, MQTT,
  GCP, Azure, Postgres, email, schedule).

OpenAPI spec gains the `force` field on each setmode/setenabled body.

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

* feat(forks): CLI --fork-triggers flag, fork-trigger docs, skill update

- Adds --fork-triggers boolean to wmill workspace fork; passes
  fork_triggers through to the create_fork API call.
- New docs/fork-triggers.md describing the model end-to-end (default,
  opt-in clone, merge-direction filter, conflict warning, future
  runtime-suffix work).
- Updates the adding-a-trigger SKILL.md to mention the fork-export
  ignore-keys participation and the clone_triggers_and_schedules
  block that new trigger kinds must extend.

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

* chore: regenerate sqlx offline query cache for fork-trigger SQL

* fix(forks): replace browser confirm() with ConfirmationModal for fork conflict

The fork-conflict warning previously used the browser's native confirm()
which doesn't match Windmill's design system. Switches to a singleton
ConfirmationModal mounted at the (logged) layout root, driven by a new
forkConflictModal store. The withForkConflictRetry helper now sets the
store and awaits the user's choice via a Promise, instead of blocking
on window.confirm.

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

* fix(forks): filter unchanged triggers in merge UI, add diff view, surface parent-only ones

The fork merge UI listed every trigger from the fork as a deployable item
regardless of whether it differed from the parent — so a fork created with
fork_triggers=true (which clones triggers in disabled state, otherwise
identical) showed every trigger as a "Fork-only" change. The 'Update
current' tab also missed triggers newly created in the parent that the
fork hadn't pulled yet.

This refactor:

- fetchAllTriggers now lists both fork and parent in parallel for each
  trigger kind, then merges by path.
- Computes a per-trigger `changeKind` (new / modified / deleted-in-source)
  using a JSON comparison that strips runtime + fork-local fields
  (mode/enabled/server_id/last_server_ping/edited_at/edited_by/etc.) so
  the disabled-on-clone difference doesn't show up as a change.
- Filters the trigger items in deployableItems by the current direction:
  Deploy mode shows fork-side new/modified, Update mode shows parent-side
  new/modified.
- Replaces the always-on "Fork-only" badge with proper New/Modified
  badges and surfaces a Diff button (modal Drawer + Monaco DiffEditor)
  for modified triggers — the diff strips the same ignored fields so
  users see only the meaningful config differences.

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

* fix(forks): always clone triggers/schedules disabled, drop opt-in flag

Disabled triggers and schedules are inert — no listener attaches, no cron
fires — so cloning them by default is safe by construction. Drops the
fork_triggers opt-in flag introduced earlier in this PR:

- Drops workspace.fork_triggers column (migration removed)
- Removes fork_triggers from CreateWorkspaceFork (API + OpenAPI)
- Removes the conditional in create_workspace_fork — clone always runs
- Removes the toggle from the fork-creation dialog
- Removes --fork-triggers from `wmill workspace fork`
- Updates docs/fork-triggers.md and adding-a-trigger SKILL.md

The merge UI continues to exclude triggers from the deploy/update default
selection, so a routine merge from a fork doesn't accidentally push
trigger config the user hasn't intentionally changed.

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

* fix(http-triggers): scope route exists check by workspace, skip non-workspaced clones in forks

The non-CLOUD branch of `route_path_key_exists` self-excluded by trigger
path alone, which silently masked cross-workspace collisions once forks
started cloning trigger rows verbatim. Tighten it to exclude only the
exact `(workspace_id, path)` row.

Fork creation also now skips non-workspaced HTTP triggers — their URL
has no workspace prefix, so a clone collides with the parent at the
matchit router (which silently drops one of two duplicates) and there is
no namespacing escape hatch. The clone copies all rows when CLOUD_HOSTED
or HTTP_ROUTE_WORKSPACED_ROUTE forces every route workspaced regardless.

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

* fix(forks-ui): silent cancel on enable conflict, clean up trigger rows in compare view

forkConflict helper now returns undefined when the user dismisses the
modal instead of throwing, so the redundant 'Cannot enable: undefined'
toast no longer appears.

CompareWorkspaces trigger rows now mirror the script row layout: drop
the redundant Disabled badge and the Trash/Details buttons (both belong
on the dedicated trigger pages, not in the deploy/compare view); pass
triggerKind through so RowIcon picks the right kind-specific icon; move
extraLabel into the summary line; replace the yellow Modified badge
with the same green ↗ ahead / blue ↘ behind treatment scripts use.

Trigger diff drawer: switch JSON → YAML for parity with DiffDrawer, fix
zero-height monaco render with className=!h-full, drop the redundant
Original/Modified label banner above the diff.

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

* fix(email-trigger): scope local_part exists check, skip non-workspaced clones in forks

Mirrors the HTTP route fix for the email-trigger non-CLOUD `email_exists`
check (in EE) which had the same path-only self-exclusion bug, and the
fork clone of `email_trigger` rows which copied non-workspaced
`local_part` verbatim. Skip non-workspaced rows in the clone unless the
instance is CLOUD_HOSTED (where lookup is workspace-scoped natively).

EE companion change in windmill-trigger-email/src/handler_ee.rs.

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

* chore: update ee-repo-ref to 78512dd73b4a1c9f70574cff863374179e3a621b

This commit updates the EE repository reference after PR #554 was merged in windmill-ee-private.

Previous ee-repo-ref: 1ac77f50747b58e720a11162dfd309bc252a24ab

New ee-repo-ref: 78512dd73b4a1c9f70574cff863374179e3a621b

Automated by sync-ee-ref workflow.

* fix(forks): always-warn on parent row, kind-specific modal copy, cancel-aware toggles

- Conflict check now fires whenever the parent has the path (regardless of
  parent's mode), since the cloned upstream identifier is shared by
  construction; closes the Postgres slot-takeover gap when the parent is
  disabled. Schedule's set_schedule_enabled gets the same treatment.
- Skip the warning entirely for HTTP and Email via a new
  TriggerCrud::FORK_CONFLICT_ON_ENABLE const — both kinds are workspace-
  scoped at runtime so cloned rows can't collide with the parent.
- Modal copy branches by failure family: split-events (Kafka/NATS/MQTT/SQS/
  GCP/Azure), duplicate-firing (Websocket/Schedule), slot-takeover
  (Postgres). Generic fallback for unknown kinds.
- withForkConflictRetry now returns boolean (true=committed, false=
  cancelled). TriggerModeToggle reuses its existing innerTriggerMode local
  state via a function binding for the regular Toggle, snapping back to
  the prop when onToggleMode signals a cancel — needed because the native
  bind:checked diverges from the parent's prop after a click and Svelte's
  reactivity won't re-push a same-valued prop down. Schedule list page
  uses {#key} on a reset version since it renders Toggle directly.
- Editor inners revert mode = previousMode on cancel; list pages skip the
  re-fetch (loadTriggers/loadSchedules) on cancel to avoid pointless
  network traffic and the schedule "Job stats loading..." flash.
- Drop withForkConflictRetry from HTTP and Email editors + list pages
  since the backend never emits the conflict for those kinds.

* fix(forks-ui): widen onToggleMode types, scope schedule toggle reset by path

- TriggerEditorToolbar and TriggerSuspendedJobsModal forwarded
  onToggleMode as `(mode) => void`, dropping the new boolean return so
  any caller wired through them would silently no-op the cancel-revert.
  Match the wider TriggerModeToggle signature.
- Schedule list page used a single resetVersion counter for every row's
  {#key}, so cancelling on any one schedule remounted every <Toggle> on
  the page. Switch to a per-path Record<string, number> bumped only for
  the affected row.

* chore: bump ee-repo-ref to c3a4553 (email FORK_CONFLICT_ON_ENABLE override)

* fix(forks): include Suspended in conflict gate, use parent_workspace_id for fork detection

Three fixes from the Claude review on PR #8976:

- Suspended mode still attaches the listener (it just pauses auto-run of
  queued jobs); two suspended fork+parent listeners would still split
  Kafka events / share a PG slot. Gate set_trigger_mode on
  `mode != Disabled` instead of `mode == Enabled` so Suspended also
  surfaces the warning.
- workspaces_export.rs::fork_*_ignore_keys keyed off the wm-fork-* prefix
  while set_trigger_mode and set_schedule_enabled key off
  parent_workspace_id. Switch the export filter to query
  parent_workspace_id once at the top of tarball_workspace and pass
  is_fork through. The column is the contract; the prefix is a
  creation-time naming convention that could in principle drift.
- TriggerModeToggle's suspend-dropdown action reassigned the non-bindable
  `triggerMode` prop instead of the local `innerTriggerMode` mirror,
  leaking inconsistent state if the dispatch was cancelled. Now writes
  to innerTriggerMode like the Toggle's on:change handler does.

* fix(cli): skip setScheduleEnabled when local YAML lacks `enabled`

Tarball export from a fork strips `enabled` from schedules so the
fork→parent git-sync round-trip can't flip the parent's operational
state. The CLI's pushSchedule called setScheduleEnabled whenever
`localSchedule.enabled != schedule.enabled`, which evaluates truthy
when local is undefined (fork-pulled YAML) and remote is true/false —
sending `{ enabled: undefined }` that serializes to `{}` and gets
rejected by the backend (`SetEnabled.enabled` is required).

Skip the call when `localSchedule.enabled === undefined` so a sync push
of fork-pulled YAMLs preserves the target's existing enabled state
instead of erroring out. Trigger updates were already safe — the
backend's update_trigger preserves `mode` when the request omits it.

* Revert "fix(cli): skip setScheduleEnabled when local YAML lacks `enabled`"

This reverts commit 23ba7e72fc.

* feat(cli): --force flag and friendlier error on fork-conflict for schedule enable

`wmill schedule enable foo/bar` against a fork whose parent has the same
path used to surface the raw `fork-conflict:schedule:<parent>` error
body. The CLI now:

- accepts `--force` to bypass the warning (mirrors the API field and the
  UI's "Enable anyway" confirmation),
- detects the `fork-conflict:` prefix on errors and prints a one-screen
  explanation pointing at --force instead of the raw body.

Disable doesn't trigger the warning (the gate fires only on transitions
to listener-attaching modes), so no flag there. Trigger enable/disable
isn't exposed as a standalone CLI command — sync push goes through
updateTrigger which has its own backend mode-preservation, so no
fork-conflict surfaces from the CLI for those.

* chore: regenerate cli-commands docs after adding --force to schedule enable

* chore: update ee-repo-ref to 967f961f0a88b027d894aebd03977181129477a8

This commit updates the EE repository reference after PR #555 was merged in windmill-ee-private.

Previous ee-repo-ref: c3a4553296473932e15392a06415dd7fb9aa6591

New ee-repo-ref: 967f961f0a88b027d894aebd03977181129477a8

Automated by sync-ee-ref workflow.

* fix(forks): address CI dead-code, claude/cubic review feedback

- backend: cfg-gate `fork_trigger_ignore_keys` to match its already-gated
  callsite. CI compiles with `-D warnings`, so the unused-fn under feature
  combos that disable all trigger crates was breaking check_oss/check_ee/
  cargo_test/test-linux/test-windows.
- cli: re-apply the `pushSchedule` undefined-skip (originally 23ba7e7,
  reverted in 4d172a1). Tarball export from forks strips `enabled`, so
  fork-pulled YAMLs that get sync-pushed back via `wmill schedule push`
  would otherwise serialize `{ enabled: undefined }` → `{}` and the
  backend's required `SetEnabled.enabled` rejects the body. Skipping
  preserves the target's existing flag, which is the round-trip-safe
  behavior. (`wmill workspace merge` extension to triggers/schedules is
  tracked in #9001 — until then sync push is the only CLI path.)
- TriggerModeToggle suspend-dropdown action awaits onToggleMode and
  resets `innerTriggerMode = triggerMode` on cancel, matching the Toggle
  on:change handler. Without this, dismissing the fork-conflict modal on
  a Suspend transition leaves the toggle stuck in 'suspended'.
- forkConflict: when a new modal opens with a previous resolver still
  pending, resolve the older promise to false. Avoids a dangling promise
  if the user clicks toggles on two rows in quick succession.
- schedules list: bump `toggleResetVersions[path]` on the
  permission-denied branch so the Toggle re-mounts back to the prop's
  `enabled` value. Without this, a user without write permission could
  click the toggle and have it stick visually flipped.
- docs/fork-triggers.md: switch the merge-direction filter description
  from `wm-fork-*` prefix to `parent_workspace_id IS NOT NULL` (matches
  the code after 4dd38fe). Drop the misleading "merge-direction filter
  strips identifier columns too" line in Future Work — the runtime
  suffix is applied at listener attach, the stored column never carries
  it, so no export filtering is needed there.

---------

Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
2026-05-01 20:55:24 +00:00

1328 lines
42 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 axum::{
extract::{Extension, Path, Query},
routing::{delete, get, post},
Json, Router,
};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use sql_builder::{prelude::Bind, SqlBuilder};
use sqlx::{Postgres, Transaction};
use std::str::FromStr;
use windmill_api_auth::{check_scopes, maybe_refresh_folders, require_super_admin, ApiAuthed};
use windmill_audit::audit_oss::audit_log;
use windmill_audit::ActionKind;
use windmill_common::DB;
use windmill_common::{
can_preserve_on_behalf_of,
db::UserDB,
error::{Error, JsonResult, Result},
schedule::Schedule,
utils::{
escape_ilike_pattern, not_found_if_none, paginate, Pagination, ScheduleType, StripPath,
},
worker::to_raw_value,
};
use windmill_git_sync::{handle_deployment_metadata, DeployedObject};
use windmill_queue::schedule::push_scheduled_job;
/// Resolves the permissioned_as value for a schedule.
/// When preserving, uses the provided permissioned_as value directly.
fn resolve_permissioned_as(
permissioned_as: Option<&String>,
preserve_permissioned_as: Option<bool>,
authed: &ApiAuthed,
) -> String {
if let Some(permissioned_as) = permissioned_as {
if preserve_permissioned_as.unwrap_or(false) && can_preserve_on_behalf_of(authed) {
return permissioned_as.clone();
}
}
windmill_common::users::username_to_permissioned_as(&authed.username)
}
/// Create-time variant: applies the folder's `default_permissioned_as` rule when no
/// explicit preserved value is provided and the caller can preserve (admin / wm_deployers).
async fn resolve_permissioned_as_for_create(
permissioned_as: Option<&String>,
preserve_permissioned_as: Option<bool>,
path: &str,
authed: &ApiAuthed,
db: &DB,
w_id: &str,
) -> Result<String> {
if let Some(pa) = permissioned_as {
if preserve_permissioned_as.unwrap_or(false) && can_preserve_on_behalf_of(authed) {
return Ok(pa.clone());
}
}
if can_preserve_on_behalf_of(authed) {
if let Some(default) =
windmill_common::folders::resolve_folder_default_permissioned_as(db, w_id, path).await?
{
return Ok(default);
}
}
Ok(windmill_common::users::username_to_permissioned_as(
&authed.username,
))
}
fn resolve_edited_by(authed: &ApiAuthed) -> String {
authed.username.clone()
}
pub fn workspaced_service() -> Router {
Router::new()
.route("/list", get(list_schedule))
.route("/list_with_jobs", get(list_schedule_with_jobs))
.route("/get/{*path}", get(get_schedule))
.route("/exists/{*path}", get(exists_schedule))
.route("/create", post(create_schedule))
.route("/update/{*path}", post(edit_schedule))
.route("/delete/{*path}", delete(delete_schedule))
.route("/setenabled/{*path}", post(set_enabled))
.route("/setdefaulthandler", post(set_default_error_handler))
// .route("/catchup/*path", post(do_catchup).get(list_catchup))
}
pub fn global_service() -> Router {
Router::new().route("/preview", post(preview_schedule))
}
#[derive(Deserialize)]
pub struct NewSchedule {
pub path: String,
pub schedule: String,
pub timezone: String,
pub summary: Option<String>,
pub description: Option<String>,
pub no_flow_overlap: Option<bool>,
pub script_path: String,
pub is_flow: bool,
pub args: Option<serde_json::Value>,
pub enabled: Option<bool>,
pub on_failure: Option<String>,
pub on_failure_times: Option<i32>,
pub on_failure_exact: Option<bool>,
pub on_failure_extra_args: Option<serde_json::Value>,
pub on_recovery: Option<String>,
pub on_recovery_times: Option<i32>,
pub on_recovery_extra_args: Option<serde_json::Value>,
pub on_success: Option<String>,
pub on_success_extra_args: Option<serde_json::Value>,
pub ws_error_handler_muted: Option<bool>,
pub retry: Option<serde_json::Value>,
pub tag: Option<String>,
pub paused_until: Option<DateTime<Utc>>,
pub cron_version: Option<String>,
pub dynamic_skip: Option<String>,
pub permissioned_as: Option<String>,
pub preserve_permissioned_as: Option<bool>,
#[serde(default)]
pub labels: Option<Vec<String>>,
}
#[derive(Serialize, Deserialize)]
pub struct ErrorOrRecoveryHandler {
pub handler_type: HandlerType,
pub override_existing: bool,
pub path: Option<String>,
pub extra_args: Option<serde_json::Value>,
pub number_of_occurence: Option<i32>,
pub number_of_occurence_exact: Option<bool>,
pub workspace_handler_muted: Option<bool>,
}
#[derive(Serialize, Deserialize)]
#[serde(rename_all(serialize = "lowercase", deserialize = "lowercase"))]
pub enum HandlerType {
Error,
Recovery,
Success,
}
async fn check_path_conflict<'c>(
tx: &mut Transaction<'c, Postgres>,
w_id: &str,
path: &str,
) -> Result<()> {
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);
if exists {
return Err(Error::BadRequest(format!(
"Schedule {} already exists",
path
)));
}
return Ok(());
}
fn to_json_raw_opt(
value: Option<&serde_json::Value>,
) -> Option<sqlx::types::Json<Box<serde_json::value::RawValue>>> {
value.map(|v| sqlx::types::Json(to_raw_value(&v)))
}
/// Validate that a dynamic skip handler (script or flow) exists
async fn validate_dynamic_skip<'c>(
tx: &mut Transaction<'c, Postgres>,
w_id: &str,
handler_path: &str,
) -> Result<()> {
// Check for script only (flows are not supported in the UI)
let exists = sqlx::query_scalar!(
"SELECT EXISTS(
SELECT 1 FROM script
WHERE workspace_id = $1 AND path = $2 AND archived = false AND deleted = false
)",
w_id,
handler_path
)
.fetch_one(&mut **tx)
.await?
.unwrap_or(false);
if exists {
Ok(())
} else {
Err(Error::BadRequest(format!(
"Dynamic skip handler '{}' not found. The handler must be an existing, non-archived script at schedule creation time.",
handler_path
)))
}
}
async fn create_schedule(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
Json(ns): Json<NewSchedule>,
) -> Result<String> {
check_scopes(&authed, || format!("schedules:write:{}", ns.path))?;
let authed = maybe_refresh_folders(&ns.path, &w_id, authed, &db).await;
#[cfg(not(feature = "enterprise"))]
if ns.on_recovery.is_some() {
return Err(Error::BadRequest(
"on_recovery is only available in enterprise version".to_string(),
));
}
#[cfg(not(feature = "enterprise"))]
if ns.on_success.is_some() {
return Err(Error::BadRequest(
"on_success is only available in enterprise version".to_string(),
));
}
#[cfg(not(feature = "enterprise"))]
if ns.on_failure_times.is_some() && ns.on_failure_times.unwrap() > 1 {
return Err(Error::BadRequest(
"on_failure with a number of times > 1 is only available in enterprise version"
.to_string(),
));
}
let mut tx: Transaction<'_, Postgres> = user_db.begin(&authed).await?;
// Check schedule for error
ScheduleType::from_str(&ns.schedule, ns.cron_version.as_deref(), true)?;
check_path_conflict(&mut tx, &w_id, &ns.path).await?;
check_flow_conflict(&mut tx, &w_id, &ns.path, ns.is_flow, &ns.script_path).await?;
// Validate dynamic_skip if provided
if let Some(handler_path) = &ns.dynamic_skip {
validate_dynamic_skip(&mut tx, &w_id, handler_path).await?;
}
let resolved_edited_by = resolve_edited_by(&authed);
let resolved_permissioned_as = resolve_permissioned_as_for_create(
ns.permissioned_as.as_ref(),
ns.preserve_permissioned_as,
&ns.path,
&authed,
&db,
&w_id,
)
.await?;
// email is still written for backwards compat with old workers that don't know about permissioned_as
let resolved_email = windmill_common::users::get_email_from_permissioned_as(
&resolved_permissioned_as,
&w_id,
&db,
)
.await?;
let schedule = sqlx::query_as!(
Schedule,
r#"
INSERT INTO schedule (
workspace_id, path, schedule, timezone, edited_by, script_path,
is_flow, args, enabled, email, permissioned_as,
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, summary, no_flow_overlap,
tag, paused_until, cron_version, description, dynamic_skip, labels
) VALUES (
$1, $2, $3, $4, $5, $6,
$7, $8, $9, $10, $11,
$12, $13, $14, $15,
$16, $17, $18,
$19, $20,
$21, $22, $23, $24,
$25, $26, $27, $28, $29, $30
)
RETURNING
workspace_id,
path,
edited_by,
edited_at,
schedule,
timezone,
enabled,
script_path,
is_flow,
args AS "args: _",
extra_perms,
email,
permissioned_as,
error,
on_failure,
on_failure_times,
on_failure_exact,
on_failure_extra_args AS "on_failure_extra_args: _",
on_recovery,
on_recovery_times,
on_recovery_extra_args AS "on_recovery_extra_args: _",
on_success,
on_success_extra_args AS "on_success_extra_args: _",
ws_error_handler_muted,
retry,
no_flow_overlap,
summary,
description,
tag,
paused_until,
cron_version,
dynamic_skip,
labels
"#,
w_id,
ns.path,
ns.schedule,
ns.timezone,
resolved_edited_by,
ns.script_path,
ns.is_flow,
to_json_raw_opt(ns.args.as_ref())
as Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
ns.enabled.unwrap_or(false),
resolved_email,
resolved_permissioned_as,
ns.on_failure,
ns.on_failure_times,
ns.on_failure_exact,
to_json_raw_opt(ns.on_failure_extra_args.as_ref())
as Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
ns.on_recovery,
ns.on_recovery_times,
to_json_raw_opt(ns.on_recovery_extra_args.as_ref())
as Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
ns.on_success,
to_json_raw_opt(ns.on_success_extra_args.as_ref())
as Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
ns.ws_error_handler_muted.unwrap_or(false),
ns.retry,
ns.summary,
ns.no_flow_overlap.unwrap_or(false),
ns.tag,
ns.paused_until,
ns.cron_version.clone().unwrap_or_else(|| "v2".to_string()),
ns.description,
ns.dynamic_skip,
ns.labels.as_deref() as Option<&[String]>
)
.fetch_one(&mut *tx)
.await
.map_err(|e| Error::internal_err(format!("inserting schedule in {w_id}: {e:#}")))?;
audit_log(
&mut *tx,
&authed,
"schedule.create",
ActionKind::Create,
&w_id,
Some(&ns.path.to_string()),
Some(
[
Some(("schedule", ns.schedule.as_str())),
Some(("script_path", ns.script_path.as_str())),
]
.into_iter()
.flatten()
.collect(),
),
)
.await?;
if let Some(on_behalf_of) = windmill_common::check_on_behalf_of_preservation(
ns.permissioned_as.as_deref(),
ns.preserve_permissioned_as.unwrap_or(false),
&authed,
&authed.username,
) {
audit_log(
&mut *tx,
&authed,
"schedule.on_behalf_of",
ActionKind::Create,
&w_id,
Some(&ns.path),
Some(
[
("on_behalf_of", on_behalf_of.as_str()),
("action", "create"),
]
.into(),
),
)
.await?;
}
if ns.enabled.unwrap_or(true) {
tx = push_scheduled_job(&db, tx, &schedule, Some(&authed.clone().into()), None).await?
}
tx.commit().await?;
handle_deployment_metadata(
&authed.email,
&authed.username,
&db,
&w_id,
DeployedObject::Schedule { path: ns.path.clone() },
Some(format!("Schedule '{}' created", ns.path.clone())),
true,
None,
)
.await?;
Ok(ns.path.to_string())
}
async fn edit_schedule(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
Json(es): Json<EditSchedule>,
) -> Result<String> {
let path = path.to_path();
check_scopes(&authed, || format!("schedules:write:{}", path))?;
let authed = maybe_refresh_folders(&path, &w_id, authed, &db).await;
let mut tx = user_db.begin(&authed).await?;
// Check schedule for error
ScheduleType::from_str(&es.schedule, es.cron_version.as_deref(), true)?;
// Validate dynamic_skip if provided
if let Some(handler_path) = &es.dynamic_skip {
validate_dynamic_skip(&mut tx, &w_id, handler_path).await?;
}
let resolved_edited_by = resolve_edited_by(&authed);
let resolved_permissioned_as = resolve_permissioned_as(
es.permissioned_as.as_ref(),
es.preserve_permissioned_as,
&authed,
);
// email is still written for backwards compat with old workers that don't know about permissioned_as.
// When permissioned_as is preserved to a different user, derive email from it.
let resolved_email = if resolved_permissioned_as
!= windmill_common::users::username_to_permissioned_as(&authed.username)
{
windmill_common::users::get_email_from_permissioned_as(
&resolved_permissioned_as,
&w_id,
&db,
)
.await?
} else {
authed.email.clone()
};
let schedule = sqlx::query_as!(
Schedule,
r#"
UPDATE schedule SET
schedule = $1,
timezone = $2,
args = $3,
on_failure = $4,
on_failure_times = $5,
on_failure_exact = $6,
on_failure_extra_args = $7,
on_recovery = $8,
on_recovery_times = $9,
on_recovery_extra_args = $10,
on_success = $11,
on_success_extra_args = $12,
ws_error_handler_muted = $13,
retry = $14,
summary = $15,
no_flow_overlap = $16,
tag = $17,
paused_until = $18,
path = $19,
workspace_id = $20,
cron_version = COALESCE($21, cron_version),
description = $22,
dynamic_skip = $23,
email = $24,
edited_by = $25,
permissioned_as = $26,
labels = COALESCE($27, labels)
WHERE path = $19 AND workspace_id = $20
RETURNING
workspace_id,
path,
edited_by,
edited_at,
schedule,
timezone,
enabled,
script_path,
is_flow,
args AS "args: _",
extra_perms,
email,
permissioned_as,
error,
on_failure,
on_failure_times,
on_failure_exact,
on_failure_extra_args AS "on_failure_extra_args: _",
on_recovery,
on_recovery_times,
on_recovery_extra_args AS "on_recovery_extra_args: _",
on_success,
on_success_extra_args AS "on_success_extra_args: _",
ws_error_handler_muted,
retry,
no_flow_overlap,
summary,
description,
tag,
paused_until,
cron_version,
dynamic_skip,
labels
"#,
es.schedule,
es.timezone,
to_json_raw_opt(es.args.as_ref())
as Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
es.on_failure,
es.on_failure_times,
es.on_failure_exact,
to_json_raw_opt(es.on_failure_extra_args.as_ref())
as Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
es.on_recovery,
es.on_recovery_times,
to_json_raw_opt(es.on_recovery_extra_args.as_ref())
as Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
es.on_success,
to_json_raw_opt(es.on_success_extra_args.as_ref())
as Option<sqlx::types::Json<Box<serde_json::value::RawValue>>>,
es.ws_error_handler_muted.unwrap_or(false),
es.retry,
es.summary,
es.no_flow_overlap.unwrap_or(false),
es.tag,
es.paused_until,
path,
w_id,
es.cron_version,
es.description,
es.dynamic_skip,
resolved_email,
resolved_edited_by,
resolved_permissioned_as,
es.labels.as_deref() as Option<&[String]>
)
.fetch_one(&mut *tx)
.await
.map_err(|e| Error::internal_err(format!("updating schedule in {w_id}: {e:#}")))?;
// clear_schedule must come AFTER UPDATE schedule to maintain consistent lock ordering
// (schedule row first, then v2_job_queue) and avoid deadlocks with concurrent operations
// like set_enabled, flow updates, and worker job completions.
clear_schedule(&mut tx, path, &w_id).await?;
audit_log(
&mut *tx,
&authed,
"schedule.edit",
ActionKind::Update,
&w_id,
Some(&path.to_string()),
Some(
[Some(("schedule", es.schedule.as_str()))]
.into_iter()
.flatten()
.collect(),
),
)
.await?;
if let Some(on_behalf_of) = windmill_common::check_on_behalf_of_preservation(
es.permissioned_as.as_deref(),
es.preserve_permissioned_as.unwrap_or(false),
&authed,
&authed.username,
) {
audit_log(
&mut *tx,
&authed,
"schedule.on_behalf_of",
ActionKind::Update,
&w_id,
Some(&path.to_string()),
Some([("on_behalf_of", on_behalf_of.as_str()), ("action", "edit")].into()),
)
.await?;
}
if schedule.enabled {
tx = push_scheduled_job(&db, tx, &schedule, None, None).await?;
}
tx.commit().await?;
handle_deployment_metadata(
&authed.email,
&authed.username,
&db,
&w_id,
DeployedObject::Schedule { path: path.to_string() },
None,
true,
None,
)
.await?;
Ok(path.to_string())
}
#[derive(Deserialize)]
pub struct ListScheduleQuery {
pub page: Option<usize>,
pub per_page: Option<usize>,
pub path: Option<String>,
pub is_flow: Option<bool>,
// filter by matching a subset of the args using base64 encoded json subset
pub args: Option<String>,
pub path_start: Option<String>,
// exact match on schedule path
pub schedule_path: Option<String>,
// filter on description (pattern match)
pub description: Option<String>,
// filter on summary (pattern match)
pub summary: Option<String>,
pub broad_filter: Option<String>,
pub label: Option<String>,
}
#[derive(sqlx::FromRow, Serialize, Deserialize, Debug, Clone)]
pub struct ScheduleLight {
pub workspace_id: String,
pub path: String,
pub edited_by: String,
pub edited_at: DateTime<chrono::Utc>,
pub schedule: String,
pub timezone: String,
pub enabled: bool,
pub script_path: String,
pub is_flow: bool,
pub summary: Option<String>,
pub extra_perms: serde_json::Value,
#[serde(skip_serializing_if = "Option::is_none")]
pub labels: Option<Vec<String>>,
}
async fn list_schedule(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
Query(lsq): Query<ListScheduleQuery>,
) -> JsonResult<Vec<ScheduleLight>> {
let mut tx = user_db.begin(&authed).await?;
let (per_page, offset) = paginate(Pagination { per_page: lsq.per_page, page: lsq.page });
let mut sqlb = SqlBuilder::select_from("schedule")
.fields(&[
"workspace_id",
"path",
"edited_by",
"edited_at",
"schedule",
"timezone",
"enabled",
"script_path",
"is_flow",
"summary",
"extra_perms",
"labels",
])
.order_by("edited_at", true)
.and_where("workspace_id = ?".bind(&w_id))
.offset(offset)
.limit(per_page)
.clone();
if let Some(path) = lsq.path {
sqlb.and_where_eq("script_path", "?".bind(&path));
}
if let Some(is_flow) = lsq.is_flow {
sqlb.and_where_eq("is_flow", "?".bind(&is_flow));
}
if let Some(args) = &lsq.args {
if let Ok(v) = serde_json::from_str::<serde_json::Value>(args) {
sqlb.and_where("args @> ?".bind(&v.to_string()));
} else {
sqlb.and_where("FALSE");
}
}
if let Some(path_start) = &lsq.path_start {
sqlb.and_where_like_left("path", path_start);
}
if let Some(schedule_path) = &lsq.schedule_path {
sqlb.and_where_eq("path", "?".bind(schedule_path));
}
if let Some(description) = &lsq.description {
let pat = format!("%{}%", escape_ilike_pattern(description));
sqlb.and_where("description ILIKE ?".bind(&pat));
}
if let Some(summary) = &lsq.summary {
let pat = format!("%{}%", escape_ilike_pattern(summary));
sqlb.and_where("summary ILIKE ?".bind(&pat));
}
if let Some(broad_filter) = &lsq.broad_filter {
let pat = format!("%{}%", escape_ilike_pattern(broad_filter));
sqlb.and_where(
"(path ILIKE ? OR script_path ILIKE ? OR description ILIKE ? OR summary ILIKE ? OR schedule ILIKE ?)"
.bind(&pat).bind(&pat).bind(&pat).bind(&pat).bind(&pat)
);
}
if let Some(label) = &lsq.label {
for l in label.split(',') {
sqlb.and_where("labels @> ARRAY[?]".bind(&l.trim()));
}
}
let sql = sqlb.sql().map_err(|e| Error::internal_err(e.to_string()))?;
let rows = sqlx::query_as::<_, ScheduleLight>(&sql)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
Ok(Json(rows))
}
#[derive(Serialize, Deserialize, Debug)]
pub struct ScheduleWJobs {
pub path: String,
pub jobs: Option<Vec<serde_json::Value>>,
}
async fn list_schedule_with_jobs(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
Query(pagination): Query<Pagination>,
) -> JsonResult<Vec<ScheduleWJobs>> {
let mut tx = user_db.begin(&authed).await?;
let (per_page, offset) = paginate(pagination);
let rows = sqlx::query_as!(ScheduleWJobs,
// Query plan:
// - use of the `ix_completed_job_workspace_id_started_at_new_2` index first, then;
// - use of the `ix_v2_job_root_by_path` index; hence the `parent_job IS NULL` clause.
// - both `workspace_id = $1` checks are required to hit both indexes.
"SELECT
schedule.path, t.jobs FROM schedule,
LATERAL(SELECT ARRAY(
SELECT json_build_object('id', id, 'success', status = 'success', 'duration_ms', duration_ms)
FROM v2_job_completed c JOIN v2_job j USING (id)
WHERE trigger_kind = 'schedule'
AND trigger = schedule.path
AND c.workspace_id = $1
AND j.workspace_id = $1
AND parent_job IS NULL AND runnable_path = schedule.script_path
AND status <> 'skipped'
ORDER BY completed_at DESC
LIMIT 20
) AS jobs) t
WHERE workspace_id = $1
ORDER BY edited_at DESC
LIMIT $2 OFFSET $3",
w_id,
per_page as i64,
offset as i64
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
Ok(Json(rows))
}
// SELECT id, title AS item_title, t.tag_array
// FROM items i, LATERAL ( -- this is an implicit CROSS JOIN
// SELECT ARRAY (
// SELECT t.title
// FROM items_tags it
// JOIN tags t ON t.id = it.tag_id
// WHERE it.item_id = i.id
// ) AS tag_array
// ) t;
async fn get_schedule(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<Schedule> {
let path = path.to_path();
check_scopes(&authed, || format!("schedules:read:{}", path))?;
let mut tx = user_db.begin(&authed).await?;
let schedule_o = windmill_queue::schedule::get_schedule_opt(&mut *tx, &w_id, path).await?;
let schedule = not_found_if_none(schedule_o, "Schedule", path)?;
tx.commit().await?;
Ok(Json(schedule))
}
async fn exists_schedule(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<bool> {
let mut tx = db.begin().await?;
let res = windmill_queue::schedule::exists_schedule(&mut tx, w_id, path).await?;
tx.commit().await?;
Ok(Json(res))
}
#[derive(Deserialize)]
pub struct PreviewPayload {
pub schedule: String,
pub timezone: String,
pub cron_version: Option<String>,
}
pub async fn preview_schedule(
Json(payload): Json<PreviewPayload>,
) -> JsonResult<Vec<DateTime<Utc>>> {
let schedule =
ScheduleType::from_str(&payload.schedule, payload.cron_version.as_deref(), true)?;
let tz =
chrono_tz::Tz::from_str(&payload.timezone).map_err(|e| Error::BadRequest(e.to_string()))?;
let upcoming: Vec<DateTime<Utc>> = schedule.upcoming(tz, 5)?;
Ok(Json(upcoming))
}
pub async fn set_enabled(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
Json(payload): Json<SetEnabled>,
) -> Result<String> {
let mut tx = user_db.begin(&authed).await?;
let path = path.to_path();
check_scopes(&authed, || format!("schedules:write:{}", path))?;
// Block enabling a schedule in a fork when the parent has the same path
// (regardless of parent's enabled flag), unless force=true. Two enabled
// crons fire in lockstep; even when the parent is currently disabled the
// user is likely to re-enable it later, at which point both fire — better
// to surface that risk at every fork-side enable. There's no namespacing
// fix for schedules (Phase 3 doesn't help cron); the user has to confirm
// or point the script at fork-only side effects.
if payload.enabled && !payload.force {
let parent_id: Option<String> = sqlx::query_scalar!(
"SELECT parent_workspace_id FROM workspace WHERE id = $1",
&w_id
)
.fetch_optional(&mut *tx)
.await?
.flatten();
if let Some(parent_id) = parent_id {
let exists: Option<bool> = sqlx::query_scalar!(
"SELECT EXISTS(SELECT 1 FROM schedule WHERE workspace_id = $1 AND path = $2)",
&parent_id,
path,
)
.fetch_one(&mut *tx)
.await?;
if exists == Some(true) {
return Err(Error::BadRequest(format!(
"fork-conflict:schedule:{}",
parent_id
)));
}
}
}
// email is still written for backwards compat with old workers that don't know about permissioned_as
let schedule_o = sqlx::query_as!(
Schedule,
r#"
UPDATE schedule SET
enabled = $1,
email = $2
WHERE path = $3 AND workspace_id = $4
RETURNING
workspace_id,
path,
edited_by,
edited_at,
schedule,
timezone,
enabled,
script_path,
is_flow,
args AS "args: _",
extra_perms,
email,
permissioned_as,
error,
on_failure,
on_failure_times,
on_failure_exact,
on_failure_extra_args AS "on_failure_extra_args: _",
on_recovery,
on_recovery_times,
on_recovery_extra_args AS "on_recovery_extra_args: _",
on_success,
on_success_extra_args AS "on_success_extra_args: _",
ws_error_handler_muted,
retry,
no_flow_overlap,
summary,
description,
tag,
paused_until,
cron_version,
dynamic_skip,
labels
"#,
payload.enabled,
authed.email,
path,
w_id
)
.fetch_optional(&mut *tx)
.await?;
let schedule = not_found_if_none(schedule_o, "Schedule", path)?;
clear_schedule(&mut tx, path, &w_id).await?;
audit_log(
&mut *tx,
&authed,
"schedule.setenabled",
ActionKind::Update,
&w_id,
Some(path),
Some([("enabled", payload.enabled.to_string().as_ref())].into()),
)
.await?;
if payload.enabled {
tx = push_scheduled_job(&db, tx, &schedule, None, None).await?;
}
tx.commit().await?;
handle_deployment_metadata(
&authed.email,
&authed.username,
&db,
&w_id,
DeployedObject::Schedule { path: path.to_string() },
None,
true,
None,
)
.await?;
Ok(format!(
"succesfully updated schedule at path {} to status {}",
path, payload.enabled
))
}
// pub async fn do_catchup(
// authed: ApiAuthed,
// Extension(db): Extension<DB>,
// Extension(user_db): Extension<UserDB>,
// // Path((w_id, path)): Path<(String, StripPath)>,
// Json(payload): Json<SetEnabled>,
// ) -> Result<String> {
// let mut tx: QueueTransaction<'_, rsmq_async::MultiplexedRsmq> =
// (user_db.begin(&authed).await?).into();
// let path = path.to_path();
// let schedule_o = sqlx::query_as!(
// Schedule,
// "UPDATE schedule SET enabled = $1, email = $2 WHERE path = $3 AND workspace_id = $4 RETURNING *",
// &payload.enabled,
// authed.email,
// path,
// w_id
// )
// .fetch_optional(&mut *tx)
// .await?;
// let schedule = not_found_if_none(schedule_o, "Schedule", path)?;
// clear_schedule(&mut tx, path, &w_id).await?;
// audit_log(
// &mut *tx,
// &authed,
// "schedule.setenabled",
// ActionKind::Update,
// &w_id,
// Some(path),
// Some([("enabled", payload.enabled.to_string().as_ref())].into()),
// )
// .await?;
// if payload.enabled {
// tx = push_scheduled_job(&db, tx, &schedule, None).await?;
// }
// tx.commit().await?;
// Ok(format!(
// "succesfully updated schedule at path {} to status {}",
// path, payload.enabled
// ))
// }
async fn delete_schedule(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> Result<String> {
let path = path.to_path();
check_scopes(&authed, || format!("schedules:write:{}", path))?;
let mut tx = user_db.begin(&authed).await?;
clear_schedule(&mut tx, path, &w_id).await?;
let exists = sqlx::query_scalar!(
"SELECT 1 FROM schedule WHERE path = $1 AND workspace_id = $2",
path,
w_id
)
.fetch_optional(&mut *tx)
.await?
.flatten();
if exists.is_none() {
return Err(windmill_common::error::Error::NotFound(format!(
"Schedule {} not found",
path
)));
}
// Capture row for trashbin before deleting
let trash_data: Option<serde_json::Value> = sqlx::query_scalar(
"SELECT jsonb_build_object('row', to_jsonb(t)) FROM schedule t WHERE path = $1 AND workspace_id = $2",
)
.bind(path)
.bind(&w_id)
.fetch_optional(&mut *tx)
.await?;
let del = sqlx::query_scalar!(
"DELETE FROM schedule WHERE path = $1 AND workspace_id = $2 RETURNING 1",
path,
w_id
)
.fetch_optional(&mut *tx)
.await?
.flatten();
if del.is_none() {
return Err(windmill_common::error::Error::NotAuthorized(format!(
"Not authorized to delete schedule {}",
path
)));
}
if let Some(data) = trash_data {
windmill_common::trashbin::move_to_trash(
&mut *tx,
&w_id,
"schedule",
path,
data,
&authed.username,
)
.await?;
}
audit_log(
&mut *tx,
&authed,
"schedule.delete",
ActionKind::Delete,
&w_id,
Some(path),
None,
)
.await?;
tx.commit().await?;
handle_deployment_metadata(
&authed.email,
&authed.username,
&db,
&w_id,
DeployedObject::Schedule { path: path.to_string() },
Some(format!("Schedule '{}' deleted", path)),
true,
None,
)
.await?;
Ok(format!("schedule {} deleted", path))
}
async fn set_default_error_handler(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Path(w_id): Path<String>,
Json(payload): Json<ErrorOrRecoveryHandler>,
) -> Result<()> {
require_super_admin(&db, &authed.email).await?;
let (key, value) = match payload.handler_type {
HandlerType::Error => {
let key = format!("default_error_handler_{}", w_id);
if let Some(payload_path) = payload.path.as_ref() {
let value = serde_json::json!({
"wsErrorHandlerMuted": payload.workspace_handler_muted,
"errorHandlerPath": payload_path,
"errorHandlerExtraArgs": payload.extra_args,
"failedTimes": payload.number_of_occurence,
"failedExact": payload.number_of_occurence_exact,
});
(key, Some(value))
} else {
(key, None)
}
}
HandlerType::Recovery => {
let key = format!("default_recovery_handler_{}", w_id);
if let Some(payload_path) = payload.path.as_ref() {
let value = serde_json::json!({
"recoveryHandlerPath": payload_path,
"recoveryHandlerExtraArgs": payload.extra_args,
"recoveredTimes": payload.number_of_occurence,
});
(key, Some(value))
} else {
(key, None)
}
}
HandlerType::Success => {
let key = format!("default_success_handler_{}", w_id);
if let Some(payload_path) = payload.path.as_ref() {
let value = serde_json::json!({
"successHandlerPath": payload_path,
"successHandlerExtraArgs": payload.extra_args,
});
(key, Some(value))
} else {
(key, None)
}
}
};
if let Some(value_content) = value {
windmill_api_settings::set_global_setting_internal(&db, key, value_content).await?;
} else {
windmill_api_settings::delete_global_setting(&db, key.as_str()).await?;
}
if payload.override_existing {
let updated_schedules: Vec<String>;
match payload.handler_type {
HandlerType::Error => {
if payload.path.is_some() {
updated_schedules = sqlx::query_scalar!(
"UPDATE schedule SET ws_error_handler_muted = $1, on_failure = $2, on_failure_extra_args = $3, on_failure_times = $4, on_failure_exact = $5 WHERE workspace_id = $6 RETURNING path",
payload.workspace_handler_muted,
payload.path,
payload.extra_args,
payload.number_of_occurence,
payload.number_of_occurence_exact,
w_id,
)
.fetch_all(&db)
.await?;
} else {
updated_schedules = sqlx::query_scalar!(
"UPDATE schedule SET ws_error_handler_muted = false, on_failure = NULL, on_failure_extra_args = NULL, on_failure_times = NULL, on_failure_exact = NULL WHERE workspace_id = $1 RETURNING path",
w_id,
)
.fetch_all(&db)
.await?;
}
}
HandlerType::Recovery => {
if payload.path.is_some() {
updated_schedules = sqlx::query_scalar!(
"UPDATE schedule SET on_recovery = $1, on_recovery_extra_args = $2, on_recovery_times = $3 WHERE workspace_id = $4 RETURNING path",
payload.path,
payload.extra_args,
payload.number_of_occurence,
w_id,
)
.fetch_all(&db)
.await?;
} else {
updated_schedules = sqlx::query_scalar!(
"UPDATE schedule SET on_recovery = NULL, on_recovery_extra_args = NULL, on_recovery_times = NULL WHERE workspace_id = $1 RETURNING path",
w_id,
)
.fetch_all(&db)
.await?;
}
}
HandlerType::Success => {
if payload.path.is_some() {
updated_schedules = sqlx::query_scalar!(
"UPDATE schedule SET on_success = $1, on_success_extra_args = $2 WHERE workspace_id = $3 RETURNING path",
payload.path,
payload.extra_args,
w_id,
)
.fetch_all(&db)
.await?;
} else {
updated_schedules = sqlx::query_scalar!(
"UPDATE schedule SET on_success = NULL, on_success_extra_args = NULL WHERE workspace_id = $1 RETURNING path",
w_id,
)
.fetch_all(&db)
.await?;
}
}
}
for updated_schedule_path in updated_schedules {
handle_deployment_metadata(
&authed.email,
&authed.username,
&db,
&w_id,
DeployedObject::Schedule { path: updated_schedule_path },
None,
true,
None,
)
.await?;
}
}
Ok(())
}
async fn check_flow_conflict<'c>(
tx: &mut Transaction<'c, Postgres>,
w_id: &str,
path: &str,
is_flow: bool,
script_path: &str,
) -> Result<()> {
if path != script_path || !is_flow {
let exists_flow = sqlx::query_scalar!(
"SELECT EXISTS (SELECT 1 FROM flow WHERE path = $1 AND workspace_id = $2)",
path,
w_id
)
.fetch_one(&mut **tx)
.await?
.unwrap_or(false);
if exists_flow {
return Err(Error::BadRequest(format!(
"The path is the same as a flow, it can only trigger that flow.
However the provided path is: {script_path} and is_flow is {is_flow}"
)));
};
}
Ok(())
}
#[derive(Deserialize)]
pub struct EditSchedule {
pub schedule: String,
pub timezone: String,
pub args: Option<serde_json::Value>,
pub summary: Option<String>,
pub description: Option<String>,
pub on_failure: Option<String>,
pub on_failure_times: Option<i32>,
pub on_failure_exact: Option<bool>,
pub on_failure_extra_args: Option<serde_json::Value>,
pub on_recovery: Option<String>,
pub on_recovery_times: Option<i32>,
pub on_recovery_extra_args: Option<serde_json::Value>,
pub on_success: Option<String>,
pub on_success_extra_args: Option<serde_json::Value>,
pub ws_error_handler_muted: Option<bool>,
pub retry: Option<serde_json::Value>,
pub no_flow_overlap: Option<bool>,
pub tag: Option<String>,
pub paused_until: Option<DateTime<Utc>>,
pub cron_version: Option<String>,
pub dynamic_skip: Option<String>,
pub permissioned_as: Option<String>,
pub preserve_permissioned_as: Option<bool>,
#[serde(default)]
pub labels: Option<Vec<String>>,
}
pub use windmill_queue::schedule::clear_schedule;
#[derive(Deserialize)]
pub struct SetEnabled {
pub enabled: bool,
/// Bypass the parent-state warning when enabling a schedule in a fork
/// whose parent has the same path enabled. The frontend sets this after
/// the user confirms the duplicate-firing dialog.
#[serde(default)]
pub force: bool,
}
// #[derive(Deserialize)]
// pub struct Catchup {
// pub from: DateTime<Utc>,
// pub to: Option<DateTime<Utc>>,
// }