Files
windmill/backend/windmill-queue/tests/schedule_push.rs
hugocasa efb4a27d51 fix: replace email with permissioned_as for triggers/schedules (#8439)
* refactor: replace email with permissioned_as for triggers/schedules

Add a new `permissioned_as` column (format: `u/{username}`, `g/{group}`,
or raw email) to all trigger tables and schedule. This value is used
directly for job permission checks, removing the need for email lookups
when creating/updating triggers.

- Migration: add permissioned_as to all 9 trigger tables + schedule,
  drop email from trigger tables (schedule keeps it for backwards compat)
- Backend: resolve_email() (async, DB) -> resolve_permissioned_as() (sync)
- Email cache: get_email_from_permissioned_as() with quick_cache for
  places that still need email (fetch_api_authed, schedule backwards compat)
- Frontend: rename email/preserve_email -> permissioned_as/preserve_permissioned_as
  in deploy data and OpenAPI schemas
- Tests updated for new field names and u/{username} format

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

* fix sqlx/build

* update ee ref

* refactor: simplify resolve_edited_by to always use authed username

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

* fix compile + migration

* update ref

* test: add trigger trait method tests for permissioned_as queries

Add tests that call TriggerCrud and Listener trait methods directly
to verify dynamic SQL correctly references the permissioned_as column.
Covers get_trigger_by_path, list_triggers, set_trigger_mode, and
fetch_enabled_unlistened_triggers for all trigger types.

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

* update sqlx

* fix: use permissioned_as directly for schedules and fix audit RLS for groups

- Schedule: permissioned_as only set on create, not on edit/set_enabled
- Schedule: stop reading email column, use get_email_from_permissioned_as
- Triggers: use fetch_api_authed_from_permissioned_as instead of edited_by
- Triggers: rename listener fields for clarity (username -> edited_by)
- Fix audit author username for group permissioned_as (g/test -> group-test)
  to match session.user, preventing RLS policy violations on audit_partitioned
- OpenAPI: remove permissioned_as/preserve_permissioned_as from EditSchedule
- Add backwards-compat comments for schedule email writes

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

* chore: regenerate system prompts for permissioned_as field

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

* fix build

* refactor: generalize onBehalfOf naming, add permissioned_as to EditSchedule

- Frontend: rename onBehalfOfPermissionedAs -> onBehalfOf with comments
  explaining it carries emails for flows/scripts and permissioned_as for
  triggers/schedules
- Frontend: rename getOnBehalfOfEmail -> getOnBehalfOf,
  getOnBehalfOfPermissionedAsForDeploy -> getOnBehalfOfForDeploy,
  customOnBehalfOfEmails -> customOnBehalfOf
- Backend: add optional permissioned_as/preserve_permissioned_as to
  EditSchedule with COALESCE (only updates when provided)
- Backend: add on_behalf_of audit log for schedule edit
- Backend: remove unused resolve_on_behalf_of_permissioned_as
- Tests: remove email assertions from schedule update test (email is
  just backwards compat, only permissioned_as matters)

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

* fix: preserve email column when permissioned_as is preserved on schedule edit

Derive email from the preserved permissioned_as via cache lookup instead
of always writing authed.email. This keeps the email column consistent
with the old behavior for backwards compat with old workers.

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

* fix: update deploy UI labels from "edited by" to "run as" for triggers

Triggers now use permissioned_as (not edited_by) for permissions, so
update the deploy UI wording to reflect this. Also update wm_deployers
group description to mention schedules and permissioned_as.

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

* fix: use u/username format for custom trigger/schedule deploy selection

When picking a custom user for trigger/schedule deployment, store
u/${username} (permissioned_as format) instead of the email. Flows/scripts
continue to use email format for on_behalf_of_email.

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

* fix: show u/username format for "me" option in trigger deploy selector

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

* refactor: simplify OnBehalfOfSelector to return the right format per kind

OnBehalfOfSelector now handles the email vs permissioned_as format
internally based on kind:
- triggers: returns u/username, displays u/username in all options
- flows/scripts/apps: returns email, displays username

The onSelect callback now takes (choice, value?) where value is already
in the correct format. Parent components just store it directly without
needing to know about the format difference.

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

* fix: always show u/username format in OnBehalfOfSelector for all kinds

Display is now consistent: all kinds show u/username in the selector.
The returned value still differs (email for flows/scripts, u/username
for triggers) since the backend APIs expect different formats.

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

* fix: replace email with permissioned_as in http_trigger test insert

The email column was dropped from trigger tables in the migration.

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

* fix: review fixes — migration, app policy, capture cleanup, naming

- Migration: remove DEFAULT '', use nullable → populate → SET NOT NULL
- App policy: set both on_behalf_of and on_behalf_of_email for all choices
- OnBehalfOfSelector: return OnBehalfOfDetails {email, permissionedAs} instead of ambiguous value
- Remove unused email field from Capture struct and query
- Rename getSourceEmail/getTargetEmail → getSourceOnBehalfOf/getTargetOnBehalfOf
- Rename test functions from preserve_email to preserve_permissioned_as

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

* fix: add permissioned_as to all test schedule INSERTs

Since the migration no longer uses DEFAULT '', all INSERTs must
explicitly provide permissioned_as. Updated test fixtures and
schedule_push tests.

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

* fix: strip permissioned_as from exports/sync, fix OpenAPI required field

- Add permissioned_as to workspace export strip list (like edited_by)
- Add permissioned_as to CLI TriggerFile Omit list
- Fix TriggerExtraProperty.required: email → permissioned_as
- Regenerate frontend and CLI types

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

* fix: remove accidentally committed generated files

These directories are gitignored and should not be tracked.

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

* chore: regenerate system prompts for permissioned_as schema changes

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

* fix: remove permissioned_as from CLI TriggerFile Omit list

Already stripped in workspace export, no need to also omit from the type.

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

* fix: optimize email cache key and revert TriggerFile Omit change

- Use single concatenated string for cache key instead of (String, String) tuple
- Remove permissioned_as from CLI TriggerFile Omit (already stripped in export)

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

* fix: zero-allocation email cache lookups using Equivalent trait

Use a borrowed EmailCacheKey(&str, &str) for cache lookups via
quick_cache's Equivalent support. Only allocates (String, String)
on cache miss for insert. This is called on every trigger fire
and schedule push.

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

* fix: add permissioned_as to Schedule required fields in OpenAPI spec

The backend always returns permissioned_as (non-optional String),
so the schema should reflect that.

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

* fix: handle group- prefix in migration UPDATE statements

edited_by can be 'group-{name}' for group-owned triggers/schedules.
The migration now correctly maps these to 'g/{name}' format instead
of incorrectly producing 'u/group-{name}'.

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

* Revert "fix: handle group- prefix in migration UPDATE statements"

This reverts commit 0971392b38.

* fix: use superadmin email to resolve permissioned_as in schedule migration

For users upgrading from older versions where edited_by may not reflect
the actual schedule owner, check if the email belongs to a superadmin
and look up their username. Otherwise fall back to edited_by.

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

* fix: fall back to superadmin email when not in workspace usr table

If the superadmin isn't a member of the workspace, use their email
as raw permissioned_as instead of falling back to edited_by.

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

* fix: always update permissioned_as and email on schedule edit

Consistent with pre-refactor behavior where email and edited_by
were always updated on every edit. permissioned_as is now always
set (to editing user or preserved value), removing the COALESCE
that previously preserved it when not provided.

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

* feat: add schedule permission tests and centralize group prefix constants

Tests: schedule create/update for normal user, workspace admin, and
superadmin not in workspace. Verifies schedule fields (email,
permissioned_as, edited_by) and pushed job fields (permissioned_as,
permissioned_as_email).

Constants: centralize "u/", "g/", "group-" as PERMISSIONED_AS_USER_PREFIX,
PERMISSIONED_AS_GROUP_PREFIX, USERNAME_GROUP_PREFIX.

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

* fix: use @unknown.windmill.dev for synthetic email fallback

Prevents privilege escalation: a user with username like
'superadmin_secret' would get superadmin via the synthetic
email matching SUPERADMIN_SECRET_EMAIL. Using a different
subdomain avoids any collision with hardcoded @windmill.dev emails.

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

* update ee ref

* sqlx

* chore: regenerate system prompts after main merge

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

* chore: update ee-repo-ref to bda51bc33bcb573659e7ff07d0a23ff6e23b8148

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

Previous ee-repo-ref: 8cf1802f8fe183f430830590b4f3172a50207843

New ee-repo-ref: bda51bc33bcb573659e7ff07d0a23ff6e23b8148

Automated by sync-ee-ref workflow.

---------

Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
2026-03-20 16:28:38 +00:00

1523 lines
68 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::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(),
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,
};
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(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);
// 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(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 wraps in SingleStepFlow
// -----------------------------------------------------------------------
#[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 = 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(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(())
}
}
// ===================================================================
// 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(migrations = "../migrations", 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, permissioned_as)
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, 'u/test-user')"
)
.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(())
}
}