Files
windmill/backend/windmill-queue/tests/schedule_push.rs
Ruben Fiszel dda59767c2 feat: stamp webhook trigger_kind on token-driven job runs (#10431)
* feat: stamp ui vs webhook trigger_kind on direct job runs

* fix: gate ui trigger kind on min worker version and dedupe display names

* docs: state that the ui trigger kind attributes rather than proves

* refactor: fold the trigger fallback into one trigger_or_fallback helper

* feat: hold trigger_kind as a tolerant label on the worker paths

* chore: refresh the sqlx offline cache for the trigger_kind label queries

* chore: update ee-repo-ref to 7de7daff5eed410e0c815ad6b292d2b4303f02f2

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

Previous ee-repo-ref: 974ab910d9a30c5565e1198ee312acc6d11239f3

New ee-repo-ref: 7de7daff5eed410e0c815ad6b292d2b4303f02f2

Automated by sync-ee-ref workflow.

* fix: keep the API job structs tolerant of unknown trigger kinds too

* chore: point ee-repo-ref at the merged EE main

---------

Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
2026-07-31 14:46:08 +00:00

1719 lines
76 KiB
Rust

mod schedule_push {
use chrono::Utc;
use sqlx::{Pool, Postgres};
use windmill_common::db::Authed;
use windmill_common::jobs::{JobKind, JobTriggerKind};
use windmill_common::runnable_settings::{
from_handle, insert_rs, ConcurrencySettings, RunnableSettings, RunnableSettingsTrait,
};
use windmill_common::schedule::Schedule;
use windmill_common::scripts::ScriptHash;
use windmill_common::users::username_to_permissioned_as;
use windmill_queue::jobs::{try_schedule_next_job, MiniCompletedJob};
use windmill_queue::schedule::{
find_unarmed_schedules, push_scheduled_job, rearm_schedule, RearmOutcome,
};
fn make_schedule(overrides: impl FnOnce(&mut Schedule)) -> Schedule {
let mut s = Schedule {
workspace_id: "test-workspace".to_string(),
path: "f/system/test_schedule".to_string(),
edited_by: "test-user".to_string(),
edited_at: Utc::now(),
schedule: "0 0 */5 * * *".to_string(),
timezone: "UTC".to_string(),
enabled: true,
script_path: "f/system/test_script".to_string(),
is_flow: false,
args: None,
extra_perms: serde_json::json!({}),
email: "test@windmill.dev".to_string(),
permissioned_as: "u/test-user".to_string(),
error: None,
on_failure: None,
on_failure_times: None,
on_failure_exact: None,
on_failure_extra_args: None,
on_recovery: None,
on_recovery_times: None,
on_recovery_extra_args: None,
on_success: None,
on_success_extra_args: None,
ws_error_handler_muted: false,
retry: None,
no_flow_overlap: false,
summary: None,
description: None,
tag: None,
paused_until: None,
cron_version: None,
dynamic_skip: None,
labels: None,
};
overrides(&mut s);
s
}
fn make_authed() -> Authed {
Authed {
email: "test@windmill.dev".to_string(),
username: "test-user".to_string(),
is_admin: true,
is_operator: false,
groups: vec![],
folders: vec![],
scopes: None,
token_prefix: None,
}
}
fn make_completed_job(schedule: &Schedule) -> MiniCompletedJob {
MiniCompletedJob {
id: uuid::Uuid::new_v4(),
workspace_id: schedule.workspace_id.clone(),
runnable_id: Some(ScriptHash(100001)),
scheduled_for: Utc::now() - chrono::Duration::minutes(5),
parent_job: None,
flow_innermost_root_job: None,
runnable_path: Some(schedule.script_path.clone()),
kind: JobKind::Script,
started_at: Some(Utc::now() - chrono::Duration::minutes(4)),
permissioned_as: username_to_permissioned_as(&schedule.edited_by),
created_by: schedule.edited_by.clone(),
script_lang: None,
permissioned_as_email: schedule.email.clone(),
flow_step_id: None,
trigger_kind: Some(JobTriggerKind::Schedule.into()),
trigger: Some(schedule.path.clone()),
priority: None,
concurrent_limit: None,
tag: "deno".to_string(),
cache_ttl: None,
cache_ignore_s3_path: None,
runnable_settings_handle: None,
}
}
async fn count_queued_jobs(db: &Pool<Postgres>) -> i64 {
sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM v2_job_queue")
.fetch_one(db)
.await
.unwrap()
}
async fn get_queued_job(
db: &Pool<Postgres>,
) -> Option<(
String, // workspace_id
Option<String>, // runnable_path
Option<String>, // trigger
Option<String>, // trigger_kind as text
)> {
sqlx::query_as::<_, (String, Option<String>, Option<String>, Option<String>)>(
"SELECT j.workspace_id, j.runnable_path, j.trigger, j.trigger_kind::text
FROM v2_job j JOIN v2_job_queue q ON j.id = q.id
LIMIT 1",
)
.fetch_optional(db)
.await
.unwrap()
}
// -----------------------------------------------------------------------
// push_scheduled_job: basic script schedule
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_script_schedule(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|_| {});
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
let (ws, path, trigger, trigger_kind) = get_queued_job(&db).await.unwrap();
assert_eq!(ws, "test-workspace");
assert_eq!(path.as_deref(), Some("f/system/test_script"));
assert_eq!(trigger.as_deref(), Some("f/system/test_schedule"));
assert_eq!(trigger_kind.as_deref(), Some("schedule"));
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: flow schedule
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_flow_schedule(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.is_flow = true;
s.script_path = "f/system/test_flow".to_string();
s.path = "f/system/flow_schedule".to_string();
});
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
let (_, path, trigger, _) = get_queued_job(&db).await.unwrap();
assert_eq!(path.as_deref(), Some("f/system/test_flow"));
assert_eq!(trigger.as_deref(), Some("f/system/flow_schedule"));
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: on_behalf_of_email (script)
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_script_on_behalf_of_email(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.script_path = "f/system/obo_script".to_string();
s.path = "f/system/obo_schedule".to_string();
});
// No pre-computed authed: forces the obo path inside push_scheduled_job
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, None, None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
let email = sqlx::query_scalar::<_, String>(
"SELECT permissioned_as_email FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(email, "obo@windmill.dev");
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: on_behalf_of_email (flow)
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_flow_on_behalf_of_email(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.is_flow = true;
s.script_path = "f/system/obo_flow".to_string();
s.path = "f/system/obo_flow_schedule".to_string();
});
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, None, None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
let email = sqlx::query_scalar::<_, String>(
"SELECT permissioned_as_email FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(email, "obo@windmill.dev");
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: with retry wraps in SingleStepFlow
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_script_with_retry(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.retry = Some(serde_json::json!({
"constant": { "attempts": 3, "seconds": 10 }
}));
});
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
// Native retry: a scheduled script with a retry policy is pushed as a plain
// Script carrying the policy via runnable_settings_handle — no SingleStepFlow.
let (kind, handle) = sqlx::query_as::<_, (String, Option<i64>)>(
"SELECT kind::text, q.runnable_settings_handle FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(kind, "script");
assert!(
handle.is_some(),
"retry policy carried via runnable_settings_handle"
);
Ok(())
}
// A scheduled, concurrency-limited script with a retry policy: the materialized
// root attempt's handle must resolve to BOTH the retry policy and the script's
// concurrency settings — otherwise the retry chain runs unbounded. Regression: P1.
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_script_with_retry_keeps_concurrency(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
let concurrency = ConcurrencySettings {
concurrency_key: Some("f/system/test_script".to_string()),
concurrent_limit: Some(1),
concurrency_time_window_s: Some(60),
};
let script_handle = insert_rs(
RunnableSettings {
debouncing_settings: None,
concurrency_settings: concurrency.insert_cached(&db).await?,
retry_settings: None,
},
&db,
)
.await?;
sqlx::query(
"UPDATE script SET runnable_settings_handle = $1 WHERE workspace_id = 'test-workspace' AND path = 'f/system/test_script'",
)
.bind(script_handle)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.retry = Some(serde_json::json!({ "constant": { "attempts": 3, "seconds": 10 } }));
});
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
let handle = sqlx::query_scalar::<_, Option<i64>>(
"SELECT q.runnable_settings_handle FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
let rs = from_handle(handle, &db).await?;
assert!(
rs.retry_settings.is_some(),
"root attempt carries the retry policy"
);
let resolved = ConcurrencySettings::get(
rs.concurrency_settings
.expect("root attempt must carry concurrency settings, not just retry"),
&db,
)
.await?;
assert_eq!(resolved.concurrent_limit, Some(1));
assert_eq!(
resolved.concurrency_key.as_deref(),
Some("f/system/test_script")
);
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: duplicate detection (same schedule + time = skip)
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_duplicate_skipped(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|_| {});
let authed = make_authed();
// First push
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
// Second push with same schedule — should be idempotent
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1); // Still 1
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: invalid timezone
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_invalid_timezone(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.timezone = "Invalid/Timezone".to_string();
});
let authed = make_authed();
let tx = db.begin().await?;
let result = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await;
assert!(result.is_err());
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: invalid cron expression
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_invalid_cron(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.schedule = "not a cron".to_string();
});
let authed = make_authed();
let tx = db.begin().await?;
let result = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await;
assert!(result.is_err());
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: invalid args (not a dict)
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_invalid_args(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
let raw = serde_json::value::RawValue::from_string("[1,2,3]".to_string()).unwrap();
s.args = Some(sqlx::types::Json(raw));
});
let authed = make_authed();
let tx = db.begin().await?;
let result = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await;
assert!(result.is_err());
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: with schedule args passed to job
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_with_args(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
let raw =
serde_json::value::RawValue::from_string(r#"{"key":"value"}"#.to_string()).unwrap();
s.args = Some(sqlx::types::Json(raw));
});
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
let args = sqlx::query_scalar::<_, serde_json::Value>(
"SELECT args FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(args, serde_json::json!({"key": "value"}));
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: script not found
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_script_not_found(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.script_path = "f/system/nonexistent".to_string();
});
let authed = make_authed();
let tx = db.begin().await?;
let result = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await;
assert!(result.is_err());
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: flow not found
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_flow_not_found(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.is_flow = true;
s.script_path = "f/system/nonexistent_flow".to_string();
});
let authed = make_authed();
let tx = db.begin().await?;
let result = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await;
assert!(result.is_err());
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: paused schedule (paused_until in future)
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_paused_schedule(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.paused_until = Some(Utc::now() + chrono::Duration::hours(1));
});
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
// Job is still pushed, but scheduled_for will be after paused_until
assert_eq!(count_queued_jobs(&db).await, 1);
let scheduled_for = sqlx::query_scalar::<_, chrono::DateTime<Utc>>(
"SELECT scheduled_for FROM v2_job_queue LIMIT 1",
)
.fetch_one(&db)
.await?;
assert!(scheduled_for > Utc::now());
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: clock shift detection (now_cutoff >= now)
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_clock_shift(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|_| {});
let authed = make_authed();
// Pass a now_cutoff far in the future — simulates clock shift
let future_cutoff = Utc::now() + chrono::Duration::hours(24);
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), Some(future_cutoff)).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
// The scheduled_for should be after the cutoff
let scheduled_for = sqlx::query_scalar::<_, chrono::DateTime<Utc>>(
"SELECT scheduled_for FROM v2_job_queue LIMIT 1",
)
.fetch_one(&db)
.await?;
assert!(scheduled_for > future_cutoff);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: disabled schedule does not push
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_handle_disabled_schedule(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.enabled = false;
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
tx.commit().await?;
assert!(err.is_none());
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: script path mismatch does not push
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_handle_path_mismatch(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, "f/system/different_script").await;
tx.commit().await?;
assert!(err.is_none());
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: enabled + matching path pushes next job
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_handle_enabled_pushes_next_job(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
tx.commit().await?;
assert!(err.is_none());
assert_eq!(count_queued_jobs(&db).await, 1);
let (_, path, trigger, trigger_kind) = get_queued_job(&db).await.unwrap();
assert_eq!(path.as_deref(), Some("f/system/test_script"));
assert_eq!(trigger.as_deref(), Some("f/system/test_schedule"));
assert_eq!(trigger_kind.as_deref(), Some("schedule"));
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: on_behalf_of_email via handle path
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_handle_on_behalf_of_email(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.script_path = "f/system/obo_script".to_string();
s.path = "f/system/obo_schedule".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
tx.commit().await?;
assert!(err.is_none());
assert_eq!(count_queued_jobs(&db).await, 1);
let email = sqlx::query_scalar::<_, String>(
"SELECT permissioned_as_email FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(email, "obo@windmill.dev");
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: push failure returns error to caller
// (caller is responsible for retry + eventual disable)
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_handle_push_failure_disables_schedule(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/bad_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
// NotFound: schedule disabled internally, no error returned (caller commits)
assert!(err.is_none());
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 0);
// Schedule should be disabled (NotFound is non-retryable)
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/bad_schedule'",
)
.fetch_one(&db)
.await?;
assert!(!enabled);
assert!(error.is_some());
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: successful push is atomic with tx commit
// If the caller commits, both the next tick and any prior writes persist.
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_success_atomic_with_commit(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
assert!(err.is_none());
tx.commit().await?;
// Next tick was pushed
assert_eq!(count_queued_jobs(&db).await, 1);
// Schedule stayed enabled
let enabled: bool = sqlx::query_scalar(
"SELECT enabled FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/test_schedule'",
)
.fetch_one(&db)
.await?;
assert!(enabled);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: successful push rolls back if tx is dropped
// Ensures no next tick leaks when the outer tx is not committed.
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_success_rolls_back_on_drop(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
assert!(err.is_none());
// Intentionally drop tx without committing (simulates caller failure)
drop(tx);
// Nothing should be visible — the next tick must NOT leak
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: schedule disable rolls back if tx is dropped
// The schedule must stay enabled when the caller doesn't commit.
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_failure_disable_rolls_back_on_drop(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/bad_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
// NotFound: schedule disabled in tx, no error returned
assert!(err.is_none());
// Drop without commit — simulates zombie retry path
drop(tx);
// Schedule should STILL be enabled (disable was rolled back with tx)
let enabled: bool = sqlx::query_scalar(
"SELECT enabled FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/bad_schedule'",
)
.fetch_one(&db)
.await?;
assert!(enabled, "schedule must stay enabled when tx is rolled back");
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: tx remains usable after successful push
// The caller can perform additional writes on the returned tx.
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_tx_usable_after_success(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (mut tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
assert!(err.is_none());
// Write something else on the same tx
sqlx::query(
"INSERT INTO global_settings (name, value) VALUES ('_test_after_push', '42'::jsonb)",
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
// Both the pushed job and the extra write should be visible
assert_eq!(count_queued_jobs(&db).await, 1);
let val: serde_json::Value =
sqlx::query_scalar("SELECT value FROM global_settings WHERE name = '_test_after_push'")
.fetch_one(&db)
.await?;
assert_eq!(val, serde_json::json!(42));
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: tx remains usable after push failure + disable
// The caller can still write on the returned tx after a schedule disable.
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_tx_usable_after_failure(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/bad_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (mut tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
// NotFound: schedule disabled in tx, no error returned
assert!(err.is_none());
// Write something else on the returned tx — tx is still usable
sqlx::query(
"INSERT INTO global_settings (name, value) VALUES ('_test_after_fail', '99'::jsonb)",
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
// Schedule should be disabled (NotFound is non-retryable)
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/bad_schedule'",
)
.fetch_one(&db)
.await?;
assert!(!enabled);
assert!(error.is_some());
// Extra write should still be committed (tx is usable after push failure)
let val: serde_json::Value =
sqlx::query_scalar("SELECT value FROM global_settings WHERE name = '_test_after_fail'")
.fetch_one(&db)
.await?;
assert_eq!(val, serde_json::json!(99));
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: flow schedule pushes next job
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_try_schedule_flow(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.is_flow = true;
s.script_path = "f/system/test_flow".to_string();
s.path = "f/system/flow_schedule".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
tx.commit().await?;
assert!(err.is_none());
assert_eq!(count_queued_jobs(&db).await, 1);
let (_, path, trigger, _) = get_queued_job(&db).await.unwrap();
assert_eq!(path.as_deref(), Some("f/system/test_flow"));
assert_eq!(trigger.as_deref(), Some("f/system/flow_schedule"));
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: script with retry is a native Script (no wrapping)
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_try_schedule_with_retry(db: Pool<Postgres>) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.retry = Some(serde_json::json!({
"constant": { "attempts": 3, "seconds": 10 }
}));
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
tx.commit().await?;
assert!(err.is_none());
assert_eq!(count_queued_jobs(&db).await, 1);
let (kind, handle) = sqlx::query_as::<_, (String, Option<i64>)>(
"SELECT kind::text, q.runnable_settings_handle FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(kind, "script");
assert!(
handle.is_some(),
"retry policy carried via runnable_settings_handle"
);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: push failure error message is stored on schedule
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_failure_stores_error_message(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/bad_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
// NotFound: schedule disabled in tx, no error returned (caller commits)
assert!(err.is_none());
tx.commit().await?;
let error: String = sqlx::query_scalar(
"SELECT error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/bad_schedule'",
)
.fetch_one(&db)
.await?;
// Error should mention the script that couldn't be found
assert!(
error.contains("nonexistent"),
"error message should describe the failure, got: {error}"
);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: disabled schedule leaves no side effects
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_disabled_schedule_no_side_effects(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/disabled_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', false, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/disabled_schedule".to_string();
s.enabled = false;
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
tx.commit().await?;
assert!(err.is_none());
// Schedule should remain disabled (not re-enabled) and no error set
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/disabled_schedule'",
)
.fetch_one(&db)
.await?;
assert!(!enabled);
assert!(
error.is_none(),
"disabled schedule should not get an error set"
);
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: path mismatch leaves schedule unchanged
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_path_mismatch_no_side_effects(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, "f/system/different_script").await;
tx.commit().await?;
assert!(err.is_none());
// Schedule should remain enabled, no error, no jobs pushed
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/test_schedule'",
)
.fetch_one(&db)
.await?;
assert!(enabled, "schedule must stay enabled on path mismatch");
assert!(error.is_none(), "no error should be set on path mismatch");
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: push failure with schedule not in DB
// When the schedule row doesn't exist, the UPDATE affects 0 rows but
// doesn't error. The function should return (tx, None).
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_failure_schedule_not_in_db(db: Pool<Postgres>) -> anyhow::Result<()> {
// Do NOT insert a schedule row — the disable UPDATE will match 0 rows
let schedule = make_schedule(|s| {
s.path = "f/system/ghost_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
drop(tx);
// NotFound: disable succeeds (UPDATE 0 rows is not an error), no error returned
assert!(err.is_none());
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// -----------------------------------------------------------------------
// Critical invariant: after commit, it's impossible to have schedule
// enabled + no next tick + function returned success.
// We verify: if push succeeds, both the job AND the schedule's
// unmodified enabled state are committed together.
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_invariant_success_means_tick_committed(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
assert!(err.is_none());
tx.commit().await?;
// After commit: schedule enabled AND next tick exists — invariant holds
let enabled: bool = sqlx::query_scalar(
"SELECT enabled FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/test_schedule'",
)
.fetch_one(&db)
.await?;
assert!(enabled, "schedule must be enabled after successful push");
assert_eq!(
count_queued_jobs(&db).await,
1,
"next tick must exist after successful push + commit"
);
Ok(())
}
// -----------------------------------------------------------------------
// Critical invariant: after commit with push failure, the schedule is
// disabled. It's never the case that we commit with the schedule still
// enabled and no next tick.
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_invariant_failure_means_disabled_after_commit(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/bad_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
// NotFound: schedule disabled in tx, no error returned (caller commits)
assert!(err.is_none());
tx.commit().await?;
// After commit: no next tick, but schedule is disabled — invariant holds
assert_eq!(count_queued_jobs(&db).await, 0);
let enabled: bool = sqlx::query_scalar(
"SELECT enabled FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/bad_schedule'",
)
.fetch_one(&db)
.await?;
assert!(
!enabled,
"schedule must be disabled when push fails with NotFound and tx commits"
);
Ok(())
}
// -----------------------------------------------------------------------
// Critical invariant: if tx is NOT committed (zombie path), neither
// the next tick nor the schedule disable persists. The schedule stays
// enabled so that zombie retry can re-attempt.
// -----------------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_invariant_rollback_preserves_schedule_for_retry(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/bad_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
let tx = db.begin().await?;
let (tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path).await;
// NotFound: schedule disabled in tx, no error returned
assert!(err.is_none());
// Simulate zombie path: drop tx without commit
drop(tx);
// Schedule still enabled (disable was rolled back with tx) — ready for retry
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/bad_schedule'",
)
.fetch_one(&db)
.await?;
assert!(enabled, "schedule must remain enabled after rollback");
assert!(error.is_none(), "error must not persist after rollback");
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
// ===================================================================
// Failpoint tests — feature-gated, only compiled under `failpoints`
// ===================================================================
#[cfg(feature = "failpoints")]
mod failpoint_tests {
use super::*;
use windmill_queue::jobs::schedule_failpoints::{ScheduleFailPoint, ACTIVE};
// ---------------------------------------------------------------
// SavepointCreate failpoint → schedule disabled, 0 jobs
// ---------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_failpoint_savepoint_create_disables(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
ACTIVE.scope(ScheduleFailPoint::SavepointCreate, async {
let tx = db.begin().await.unwrap();
let (tx, err) = try_schedule_next_job(
&db, tx, &job, &schedule, &schedule.script_path,
).await;
// Transient error returned to caller for retry
assert!(err.is_some());
drop(tx);
assert_eq!(count_queued_jobs(&db).await, 0);
let enabled: bool = sqlx::query_scalar(
"SELECT enabled FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/test_schedule'",
)
.fetch_one(&db)
.await
.unwrap();
assert!(enabled, "schedule must stay enabled for caller retry");
}).await;
Ok(())
}
// ---------------------------------------------------------------
// Push failpoint → schedule disabled, 0 jobs
// ---------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_failpoint_push_disables(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
ACTIVE.scope(ScheduleFailPoint::Push, async {
let tx = db.begin().await.unwrap();
let (tx, err) = try_schedule_next_job(
&db, tx, &job, &schedule, &schedule.script_path,
).await;
// Transient error returned to caller for retry
assert!(err.is_some());
drop(tx);
assert_eq!(count_queued_jobs(&db).await, 0);
let enabled: bool = sqlx::query_scalar(
"SELECT enabled FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/test_schedule'",
)
.fetch_one(&db)
.await
.unwrap();
assert!(enabled, "schedule must stay enabled for caller retry");
}).await;
Ok(())
}
// ---------------------------------------------------------------
// SavepointCommit failpoint → schedule disabled, 0 jobs
// ---------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_failpoint_savepoint_commit_disables(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
ACTIVE.scope(ScheduleFailPoint::SavepointCommit, async {
let tx = db.begin().await.unwrap();
let (tx, err) = try_schedule_next_job(
&db, tx, &job, &schedule, &schedule.script_path,
).await;
// Transient error returned to caller for retry
assert!(err.is_some());
drop(tx);
assert_eq!(count_queued_jobs(&db).await, 0, "pushed job must be rolled back");
let enabled: bool = sqlx::query_scalar(
"SELECT enabled FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/test_schedule'",
)
.fetch_one(&db)
.await
.unwrap();
assert!(enabled, "schedule must stay enabled for caller retry");
}).await;
Ok(())
}
// ---------------------------------------------------------------
// ScheduleDisable failpoint → returns Some(err), caller doesn't commit
// ---------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_failpoint_schedule_disable_returns_err(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/bad_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
ACTIVE
.scope(ScheduleFailPoint::ScheduleDisable, async {
let tx = db.begin().await.unwrap();
let (_tx, err) =
try_schedule_next_job(&db, tx, &job, &schedule, &schedule.script_path)
.await;
assert!(err.is_some(), "must return Some(err) when disable fails");
})
.await;
Ok(())
}
// ---------------------------------------------------------------
// ScheduleDisable failpoint + tx drop → schedule stays enabled
// ---------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_failpoint_disable_failure_rollback(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/system/bad_schedule".to_string();
s.script_path = "f/system/nonexistent".to_string();
});
let job = make_completed_job(&schedule);
ACTIVE.scope(ScheduleFailPoint::ScheduleDisable, async {
let tx = db.begin().await.unwrap();
let (tx, err) = try_schedule_next_job(
&db, tx, &job, &schedule, &schedule.script_path,
).await;
assert!(err.is_some(), "must return Some(err) when disable fails");
drop(tx);
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/bad_schedule'",
)
.fetch_one(&db)
.await
.unwrap();
assert!(enabled, "schedule must stay enabled when tx is dropped after disable failure");
assert!(error.is_none(), "error must not persist after rollback");
}).await;
Ok(())
}
// ---------------------------------------------------------------
// PushQuotaExceeded failpoint (script) → schedule disabled, 0 jobs,
// no error handler notification (QuotaExceeded is silenced)
// ---------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_failpoint_push_quota_exceeded_script(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|_| {});
let job = make_completed_job(&schedule);
ACTIVE.scope(ScheduleFailPoint::PushQuotaExceeded, async {
let tx = db.begin().await.unwrap();
let (tx, err) = try_schedule_next_job(
&db, tx, &job, &schedule, &schedule.script_path,
).await;
// QuotaExceeded: schedule disabled internally, no error returned
assert!(err.is_none(), "QuotaExceeded should be handled internally (returns None)");
tx.commit().await.unwrap();
assert_eq!(count_queued_jobs(&db).await, 0);
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/test_schedule'",
)
.fetch_one(&db)
.await
.unwrap();
assert!(!enabled, "schedule must be disabled after QuotaExceeded");
assert!(error.is_some(), "error must be set on schedule");
assert!(error.unwrap().contains("quota"), "error message should mention quota");
}).await;
Ok(())
}
// ---------------------------------------------------------------
// PushQuotaExceeded failpoint (flow) → schedule disabled, 0 jobs
// ---------------------------------------------------------------
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_failpoint_push_quota_exceeded_flow(db: Pool<Postgres>) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', 'f/system/flow_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_flow', true, 'test@windmill.dev', '{}', false, false, 'u/test-user')"
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.is_flow = true;
s.script_path = "f/system/test_flow".to_string();
s.path = "f/system/flow_schedule".to_string();
});
let job = make_completed_job(&schedule);
ACTIVE.scope(ScheduleFailPoint::PushQuotaExceeded, async {
let tx = db.begin().await.unwrap();
let (tx, err) = try_schedule_next_job(
&db, tx, &job, &schedule, &schedule.script_path,
).await;
// QuotaExceeded: schedule disabled internally, no error returned
assert!(err.is_none(), "QuotaExceeded should be handled internally (returns None)");
tx.commit().await.unwrap();
assert_eq!(count_queued_jobs(&db).await, 0);
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/flow_schedule'",
)
.fetch_one(&db)
.await
.unwrap();
assert!(!enabled, "flow schedule must be disabled after QuotaExceeded");
assert!(error.is_some(), "error must be set on flow schedule");
assert!(error.unwrap().contains("quota"), "error message should mention quota");
}).await;
Ok(())
}
}
// -----------------------------------------------------------------------
// push_scheduled_job: reserved ducklake-maintenance prefix
// -----------------------------------------------------------------------
// A schedule that pre-dates the reserved prefix (a user schedule under a
// real `ducklake_maintenance` folder) must fall through to normal script
// resolution when its path's lake has no enabled maintenance config —
// never be hijacked into the maintenance payload builder and auto-disabled.
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_reserved_prefix_no_config_falls_through_to_script(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
let schedule = make_schedule(|s| {
s.path = "f/ducklake_maintenance/legacy".to_string();
});
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
let (ws, path, trigger, _) = get_queued_job(&db).await.unwrap();
assert_eq!(ws, "test-workspace");
assert_eq!(
path.as_deref(),
Some("f/system/test_script"),
"must resolve the schedule's script_path, not the maintenance builder"
);
assert_eq!(trigger.as_deref(), Some("f/ducklake_maintenance/legacy"));
Ok(())
}
// With maintenance enabled for the path's lake, the occurrence is a
// raw-code duckdb job (kind preview, runnable_path = schedule path,
// duckdb tag pinned) — enterprise builds only.
#[cfg(feature = "private")]
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_push_reserved_prefix_with_config_builds_maintenance_job(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
sqlx::query(
r#"UPDATE workspace_settings SET ducklake = '{"ducklakes": {"legacy": {
"catalog": {"resource_type": "postgresql", "resource_path": "u/test/pg"},
"storage": {"path": "legacy"},
"maintenance": {"enabled": true, "retention_days": 3}
}}}'::jsonb WHERE workspace_id = 'test-workspace'"#,
)
.execute(&db)
.await?;
let schedule = make_schedule(|s| {
s.path = "f/ducklake_maintenance/legacy".to_string();
s.script_path = "f/ducklake_maintenance/legacy".to_string();
s.tag = Some("duckdb".to_string());
});
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
let (kind, path, tag, raw_code) =
sqlx::query_as::<_, (String, Option<String>, String, Option<String>)>(
"SELECT j.kind::text, j.runnable_path, j.tag, j.raw_code
FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(kind, "preview");
assert_eq!(path.as_deref(), Some("f/ducklake_maintenance/legacy"));
assert_eq!(tag, "duckdb");
let raw_code = raw_code.expect("maintenance job must carry generated SQL");
assert!(raw_code.contains("ducklake_expire_snapshots"));
assert!(raw_code.contains("INTERVAL '3 days'"));
Ok(())
}
// Saving maintenance off must remove the managed row AND its queued
// occurrence in BOTH builds: the enterprise sync reconciles, and the
// public stub must not leave a job pushed under the enterprise edition
// to run after the admin disabled maintenance (EE-to-CE downgrade).
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_sync_disable_clears_managed_row_and_queued_occurrence(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
use std::collections::HashMap;
use windmill_common::workspaces::{
Ducklake, DucklakeCatalog, DucklakeCatalogResourceType, DucklakeMaintenance,
DucklakeStorage,
};
use windmill_queue::ducklake_maintenance::sync_ducklake_maintenance_schedules;
fn lake(maintenance_enabled: bool) -> Ducklake {
Ducklake {
catalog: DucklakeCatalog {
resource_type: DucklakeCatalogResourceType::Postgresql,
resource_path: "u/test/pg".to_string(),
},
storage: DucklakeStorage { storage: None, path: "legacy".to_string() },
extra_args: None,
fork_behavior: None,
maintenance: Some(DucklakeMaintenance {
enabled: maintenance_enabled,
schedule: None,
retention_days: None,
compaction: None,
orphan_cleanup: None,
}),
}
}
// a managed row with a queued occurrence (queued via fall-through: no
// lake config exists yet, so the push resolves the script path)
let schedule = make_schedule(|s| {
s.path = "f/ducklake_maintenance/legacy".to_string();
});
sqlx::query(
"INSERT INTO schedule (workspace_id, path, schedule, timezone, edited_by, script_path,
is_flow, enabled, email, permissioned_as, cron_version)
VALUES ($1, $2, $3, 'UTC', $4, $5, false, true, $6, $7, 'v2')",
)
.bind(&schedule.workspace_id)
.bind(&schedule.path)
.bind(&schedule.schedule)
.bind(&schedule.edited_by)
.bind(&schedule.script_path)
.bind(&schedule.email)
.bind(&schedule.permissioned_as)
.execute(&db)
.await?;
let authed = make_authed();
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &schedule, Some(&authed), None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
// maintenance saved off
let previous = HashMap::from([("legacy".to_string(), lake(true))]);
let current = HashMap::from([("legacy".to_string(), lake(false))]);
let tx = db.begin().await?;
let tx = sync_ducklake_maintenance_schedules(
&db,
tx,
&schedule.workspace_id,
&current,
&previous,
"test-user",
"test@windmill.dev",
)
.await?;
tx.commit().await?;
assert_eq!(
count_queued_jobs(&db).await,
0,
"queued occurrence must be cleared when maintenance is saved off"
);
let row_exists: bool = sqlx::query_scalar(
"SELECT EXISTS(SELECT 1 FROM schedule WHERE workspace_id = $1 AND path = $2)",
)
.bind(&schedule.workspace_id)
.bind(&schedule.path)
.fetch_one(&db)
.await?;
assert!(!row_exists, "managed schedule row must be deleted");
Ok(())
}
// -----------------------------------------------------------------------
// find_unarmed_schedules / rearm_schedule: recovery for a schedule left
// enabled with no queued occurrence (a run that died on an abnormal path
// skipped its next-occurrence push). Without this the chain stays dead
// until the schedule is manually disabled and re-enabled.
// -----------------------------------------------------------------------
async fn insert_schedule(db: &Pool<Postgres>, path: &str, script_path: &str, enabled: bool) {
sqlx::query(
"INSERT INTO schedule (workspace_id, path, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, email, extra_perms, ws_error_handler_muted, no_flow_overlap, permissioned_as)
VALUES ('test-workspace', $1, 'test-user', now(), '0 0 */5 * * *', 'UTC', $3, $2, false, 'test@windmill.dev', '{}', false, true, 'u/test-user')",
)
.bind(path)
.bind(script_path)
.bind(enabled)
.execute(db)
.await
.unwrap();
}
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_find_unarmed_schedules(db: Pool<Postgres>) -> anyhow::Result<()> {
insert_schedule(&db, "f/system/test_schedule", "f/system/test_script", true).await;
insert_schedule(&db, "f/system/disabled", "f/system/test_script", false).await;
// No occurrence queued yet: the enabled schedule is unarmed, the disabled one is ignored.
assert_eq!(
find_unarmed_schedules(&db).await?,
vec![(
"test-workspace".to_string(),
"f/system/test_schedule".to_string()
)]
);
// Once an occurrence is queued it is armed and must not be reported —
// re-arming it would double-push the occurrence.
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &make_schedule(|_| {}), None, None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
assert!(find_unarmed_schedules(&db).await?.is_empty());
Ok(())
}
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_rearm_schedule_pushes_next_occurrence(db: Pool<Postgres>) -> anyhow::Result<()> {
insert_schedule(&db, "f/system/test_schedule", "f/system/test_script", true).await;
assert_eq!(
rearm_schedule(&db, "test-workspace", "f/system/test_schedule").await?,
RearmOutcome::Rearmed
);
assert_eq!(count_queued_jobs(&db).await, 1);
assert!(find_unarmed_schedules(&db).await?.is_empty());
Ok(())
}
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_rearm_schedule_skips_disabled(db: Pool<Postgres>) -> anyhow::Result<()> {
// A disable that lands between the scan and the re-arm must win: pushing an
// occurrence for a disabled schedule would resurrect a schedule the user
// just turned off.
insert_schedule(&db, "f/system/test_schedule", "f/system/test_script", false).await;
assert_eq!(
rearm_schedule(&db, "test-workspace", "f/system/test_schedule").await?,
RearmOutcome::NoOp
);
assert_eq!(count_queued_jobs(&db).await, 0);
Ok(())
}
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_rearm_schedule_never_disables(db: Pool<Postgres>) -> anyhow::Result<()> {
insert_schedule(&db, "f/system/bad_schedule", "f/system/nonexistent", true).await;
// Reconciliation only ever starts a schedule. An unpushable occurrence is
// reported and left alone: this sweeps every enabled schedule in the
// instance, so disabling here would turn a wrong invariant into the exact
// silent stoppage the reconciler exists to undo.
assert_eq!(
rearm_schedule(&db, "test-workspace", "f/system/bad_schedule").await?,
RearmOutcome::NoOp
);
assert_eq!(count_queued_jobs(&db).await, 0);
let (enabled, error): (bool, Option<String>) = sqlx::query_as(
"SELECT enabled, error FROM schedule WHERE workspace_id = 'test-workspace' AND path = 'f/system/bad_schedule'",
)
.fetch_one(&db)
.await?;
assert!(enabled, "reconciliation must never disable a schedule");
assert!(error.is_none());
Ok(())
}
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
async fn test_rearm_schedule_skips_already_armed(db: Pool<Postgres>) -> anyhow::Result<()> {
// An occurrence can be queued (a normal completion, an edit, a re-enable)
// between the unarmed scan and rearm_schedule acquiring the row lock. Re-arming
// then would double-push, since push_scheduled_job only dedups the exact
// computed scheduled_for.
insert_schedule(&db, "f/system/test_schedule", "f/system/test_script", true).await;
let tx = db.begin().await?;
let tx = push_scheduled_job(&db, tx, &make_schedule(|_| {}), None, None).await?;
tx.commit().await?;
assert_eq!(count_queued_jobs(&db).await, 1);
assert_eq!(
rearm_schedule(&db, "test-workspace", "f/system/test_schedule").await?,
RearmOutcome::NoOp
);
assert_eq!(count_queued_jobs(&db).await, 1);
Ok(())
}
}