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