Files
windmill/backend/tests/schedule_push.rs
T
Ruben Fiszel bbb397b6ad fix: improve scheduling reliability in extreme pool contention conditions (#7825)
* fix: reuse outer tx for schedule push in commit_completed_job

Instead of calling handle_maybe_scheduled_job(db) which opens its own
connections (peak=3), inline the schedule push using a savepoint on the
outer transaction. Auth is fetched via the tx connection using
fetch_authed_from_permissioned_as_conn, and push_scheduled_job runs
on a savepoint so failures roll back only the push, not the completion.

On push failure: savepoint rolls back, schedule is disabled on the outer
tx, and the zombie return path is preserved if disabling also fails.

Peak connections drop from 3 to 1 (or 2 on cold RunnableSettings cache).

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* all

* fix: extract shared try_schedule_next_job to unify schedule push paths

Replace the two diverging schedule-push implementations (inlined in
commit_completed_job and standalone handle_maybe_scheduled_job) with a
single try_schedule_next_job that reuses the caller's transaction via
savepoints. This eliminates extra pool connection usage in the
worker_flow.rs path and ensures consistent retry/error semantics.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test: add failpoint markers to try_schedule_next_job

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* chore: remove plan.md

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: remove inner retry loop from try_schedule_next_job, add caller-level retries

The 10-retry x 5s-sleep loop inside try_schedule_next_job held locks on
v2_job_completed/v2_job_queue for up to ~45s when running inside the
outer commit_completed_job transaction.

Now try_schedule_next_job makes a single attempt and returns errors to
the caller. Non-retryable errors (QuotaExceeded, NotFound) disable the
schedule immediately inside the function. Transient errors are returned
for the caller to retry:

- commit_completed_job path: outer backon retry (10x3s) retries the
  entire transaction including the schedule push, so no locks are held
  during sleep.
- handle_flow path: new backon retry (10x3s) wraps begin/push/commit
  with a fresh transaction per attempt.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: clear push_err after successful schedule disable to prevent stuck schedules

When try_schedule_next_job disables the schedule for non-retryable errors
(NotFound, QuotaExceeded), clear the error so the caller commits the tx
(persisting the disable). Previously, the error propagated up, causing the
tx to be dropped and rolling back the disable — leaving the schedule
permanently enabled but broken.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: add 5s timeout on push_scheduled_job, clean up handle_flow error handling

- Add tokio::time::timeout(5s) around push_scheduled_job inside
  try_schedule_next_job to bound worst-case lock holding per attempt
- Remove unreachable QuotaExceeded/NotFound match arms in handle_flow
  (these errors are handled internally by try_schedule_next_job)
- Add report_error_to_workspace_handler_or_critical_side_channel in
  handle_flow when post-exhaustion schedule disable fails

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: return SchedulePushZombieError when both schedule push and disable fail

When handle_flow cannot push the next scheduled job AND cannot disable the
schedule, return a SchedulePushZombieError so the worker leaves the flow job
in the queue for zombie detection to restart. This prevents stuck schedules
where neither the next tick was pushed nor the schedule was disabled.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-02-07 13:16:32 +00:00

1616 lines
67 KiB
Rust

mod common;
mod schedule_push {
use chrono::Utc;
use sqlx::{Pool, Postgres};
use windmill_common::db::Authed;
use windmill_common::jobs::{JobKind, JobTriggerKind};
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::push_scheduled_job;
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(),
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,
};
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),
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(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(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(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(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(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);
// When retry is set, the job kind is singlescriptflow (SingleStepFlow wraps it)
let kind = sqlx::query_scalar::<_, String>(
"SELECT kind::text FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(kind, "singlestepflow");
Ok(())
}
// -----------------------------------------------------------------------
// push_scheduled_job: duplicate detection (same schedule + time = skip)
// -----------------------------------------------------------------------
#[sqlx::test(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(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(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(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(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(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(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(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(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(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(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(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(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(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)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false)"
)
.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(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)
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)"
)
.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(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(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)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false)"
)
.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(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(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)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false)"
)
.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(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 wraps in SingleStepFlow
// -----------------------------------------------------------------------
#[sqlx::test(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 = sqlx::query_scalar::<_, String>(
"SELECT kind::text FROM v2_job j JOIN v2_job_queue q ON j.id = q.id LIMIT 1",
)
.fetch_one(&db)
.await?;
assert_eq!(kind, "singlestepflow");
Ok(())
}
// -----------------------------------------------------------------------
// try_schedule_next_job: push failure error message is stored on schedule
// -----------------------------------------------------------------------
#[sqlx::test(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)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false)"
)
.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(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)
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)"
)
.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(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)
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)"
)
.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(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(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)
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)"
)
.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(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)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false)"
)
.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(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)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false)"
)
.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`
// ===================================================================
mod failpoint_tests {
use super::*;
use windmill_queue::jobs::schedule_failpoints::{ScheduleFailPoint, ACTIVE};
// ---------------------------------------------------------------
// SavepointCreate failpoint → schedule disabled, 0 jobs
// ---------------------------------------------------------------
#[sqlx::test(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)
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)"
)
.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(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)
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)"
)
.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(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)
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)"
)
.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(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)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false)"
)
.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(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)
VALUES ('test-workspace', 'f/system/bad_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/nonexistent', false, 'test@windmill.dev', '{}', false, false)"
)
.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(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)
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)"
)
.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(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)
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)"
)
.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(())
}
}
// ===================================================================
// Zombie detection tests — verify that a flow left in queue after
// SchedulePushZombieError meets the restart criteria in monitor.rs
// ===================================================================
// -----------------------------------------------------------------------
// When both schedule push AND post-retry disable fail, the flow job stays
// in v2_job_queue as a zombie. This test simulates that state and verifies:
//
// 1. The zombie detection query (from handle_zombie_flows) finds the flow
// 2. The flow meets restart criteria: first module is WaitingForPriorSteps
// and same_worker is false
// 3. After the restart UPDATE, the flow is re-queued (running=false)
// 4. The schedule remains enabled for retry on next flow execution
// -----------------------------------------------------------------------
#[sqlx::test(fixtures("base", "schedule_push"))]
async fn test_zombie_flow_after_schedule_push_failure_meets_restart_criteria(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
use windmill_common::flow_status::{FlowStatus, FlowStatusModule};
let flow_job_id = uuid::Uuid::new_v4();
let now = Utc::now();
let stale_ping = now - chrono::Duration::minutes(5);
// Schedule still enabled — simulates both push and disable failing
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)
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 */5 * * *', 'UTC', true, 'f/system/test_flow', true, 'test@windmill.dev', '{}', false, false)"
)
.execute(&db)
.await?;
// Flow job in v2_job (kind=flow, triggered by schedule, same_worker=false)
sqlx::query(
"INSERT INTO v2_job (id, workspace_id, created_at, created_by, permissioned_as, permissioned_as_email, kind, runnable_path, trigger, trigger_kind, same_worker, visible_to_owner, tag)
VALUES ($1, 'test-workspace', $2, 'test-user', 'u/test-user', 'test@windmill.dev', 'flow', 'f/system/test_flow', 'f/system/test_schedule', 'schedule', false, false, 'flow')"
)
.bind(flow_job_id)
.bind(now - chrono::Duration::minutes(5))
.execute(&db)
.await?;
// Queue entry: running=true (worker caught SchedulePushZombieError and returned Ok)
sqlx::query(
"INSERT INTO v2_job_queue (id, workspace_id, created_at, scheduled_for, running, started_at, tag, suspend)
VALUES ($1, 'test-workspace', $2, $2, true, $3, 'flow', 0)"
)
.bind(flow_job_id)
.bind(now - chrono::Duration::minutes(5))
.bind(now - chrono::Duration::minutes(4))
.execute(&db)
.await?;
// Stale ping — older than the 60s zombie transition timeout
sqlx::query("INSERT INTO v2_job_runtime (id, ping) VALUES ($1, $2)")
.bind(flow_job_id)
.bind(stale_ping)
.execute(&db)
.await?;
// Flow status at step 0, first module = WaitingForPriorSteps
// This is the initial state of a flow that hasn't started any steps yet
let flow_status = serde_json::json!({
"step": 0,
"modules": [{"type": "WaitingForPriorSteps", "id": "a"}],
"failure_module": {"type": "WaitingForPriorSteps", "id": "failure"},
"retry": {"fail_count": 0, "failed_jobs": []},
"cleanup_module": {"flow_jobs_to_clean": []}
});
sqlx::query("INSERT INTO v2_job_status (id, flow_status) VALUES ($1, $2::jsonb)")
.bind(flow_job_id)
.bind(&flow_status)
.execute(&db)
.await?;
// Run the same zombie detection query from handle_zombie_flows (60s timeout)
let zombie_flows = sqlx::query_as::<_, (uuid::Uuid, String, Option<bool>, Option<String>)>(
r#"
SELECT
j.id, j.workspace_id, j.same_worker,
COALESCE(s.flow_status, s.workflow_as_code_status)::text AS flow_status
FROM v2_job_queue q
JOIN v2_job j USING (id)
LEFT JOIN v2_job_runtime r USING (id)
LEFT JOIN v2_job_status s USING (id)
WHERE q.running = true AND q.suspend = 0 AND q.suspend_until IS null
AND q.scheduled_for <= now()
AND (j.kind = 'flow' OR j.kind = 'flowpreview' OR j.kind = 'flownode')
AND r.ping IS NOT NULL
AND r.ping < NOW() - ('60' || ' seconds')::interval
AND q.canceled_by IS NULL
"#,
)
.fetch_all(&db)
.await?;
assert_eq!(zombie_flows.len(), 1, "zombie flow must be detected");
let (id, _ws, same_worker, flow_status_json) = &zombie_flows[0];
assert_eq!(*id, flow_job_id);
// Replicate the exact branching logic from handle_zombie_flows (monitor.rs:2711-2754).
// Only flows matching the restart condition get restarted; others are cancelled.
let status = flow_status_json
.as_deref()
.and_then(|x| serde_json::from_str::<FlowStatus>(x).ok());
let should_restart = !same_worker.unwrap_or(false)
&& status.is_some_and(|s| {
s.modules
.get(0)
.is_some_and(|x| matches!(x, FlowStatusModule::WaitingForPriorSteps { .. }))
});
assert!(
should_restart,
"flow must match the restart branch (not the cancel branch) in handle_zombie_flows"
);
// Apply the restart action — same UPDATE as handle_zombie_flows
sqlx::query(
"UPDATE v2_job_queue SET running = false, started_at = null
WHERE id = $1 AND canceled_by IS NULL",
)
.bind(flow_job_id)
.execute(&db)
.await?;
// Flow is re-queued for processing
let (running, started_at): (bool, Option<chrono::DateTime<Utc>>) = sqlx::query_as(
"SELECT running, started_at FROM v2_job_queue WHERE id = $1",
)
.bind(flow_job_id)
.fetch_one(&db)
.await?;
assert!(!running, "flow must not be running after zombie restart");
assert!(started_at.is_none(), "started_at must be null after zombie restart");
// Schedule still enabled — will be retried when flow re-executes
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 remain enabled for retry after zombie restart");
// Flow job is NOT in v2_job_completed (it was never completed with error)
let completed_count: i64 = sqlx::query_scalar(
"SELECT COUNT(*) FROM v2_job_completed WHERE id = $1",
)
.bind(flow_job_id)
.fetch_one(&db)
.await?;
assert_eq!(completed_count, 0, "flow must not be in completed_job — it's a zombie, not an error");
Ok(())
}
}