mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-03 16:02:12 +00:00
* docs: plan for detecting skipped schedule occurrences Design plan only, no implementation. Records the scheduler's re-anchoring behaviour, the measurements behind it, and the three-piece design that came out of reviewing the alternatives. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * docs: state the user-facing outcome in the schedule plan The plan described the mechanism but never what a user would see, which made it hard to judge what the work is worth. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * docs: state which cause the schedule plan catches, and correct its scope Records which of the two causes each piece covers, and corrects the overrun scope: a script schedule carrying retry or dynamic_skip is pushed as a SingleStepFlow, so it re-arms at step 0 entry and its occurrences overlap like a flow's. Resolves the no_flow_overlap question, splits the read-time work into bounded detection and editor-only counting behind measured croner costs, and fixes the delivery order. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: count the occurrences a schedule skipped A schedule that overruns its interval, or waits for a worker, silently loses the occurrences in between: the scheduler keeps one queued occurrence and re-anchors on the clock, so nothing records that a run was due and never happened. Recovers the sequence from rows that already exist rather than writing per occurrence. `push_scheduled_job` anchors on `now_from_db` inside the transaction that inserts the job, and `v2_job.created_at` defaults to that same transaction timestamp, so `scheduled_for = find_next(created_at)` holds exactly and the whole occurrence history is derivable. The schedules list reports how many of the recent runs were followed by a lost occurrence, and a new occurrences endpoint carries the per-run wait and duration behind it. Detection is one `find_next` per gap, which stays bounded on a full page; counting walks the gap and runs only for a single schedule. The one write is `occurrence_baseline_at`, advanced at create, edit, re-enable and re-arm. Gaps older than it span a pause, a cron change, a re-enable or a reconciler re-arm, none of which mean runs were lost. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: show the wait and run time behind a schedule's skipped occurrences The list badge says a schedule is losing runs; this says which of the two causes did it. A large wait means not enough workers, a long run means the job outgrew its interval, and the pair is what tells them apart. Sits under the existing upcoming-events panel, so due and overdue read together. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: flag a schedule that is running late right now Reconstruction is retrospective: a gap only appears once the next occurrence has a row, which needs the current one to finish. A schedule wedged mid-run shows nothing until it moves, which is the case an operator most wants to see. An occurrence still in flight past the time its own successor was due will cost that successor, so `now > find_next(scheduled_for)` is the signal, needing no threshold and self-calibrating across a daily and a per-minute schedule. It applies only where occurrences serialize; an overlapping schedule starts its successor on time and would flag constantly while healthy. The queue is read in one aggregating pass keyed on (trigger, runnable_path) rather than a subquery per schedule, and an overlapping schedule holds more than one root row, hence the aggregate. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: run the schedule overrun alert from the monitor pass Wires `schedule_overrun_alerts` in next to `jobs_waiting_alerts`, every 30 iterations (~5 min). Its Enterprise implementation lives in windmill-labs/windmill-ee-private#772; only the wiring and the OSS stub are here. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * chore: refresh the sqlx offline cache Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: record and alert when a schedule skips occurrences push_scheduled_job compares each chained occurrence with the slot after the previous one. A gap is written to schedule.skipped_occurrences off the push transaction, alerts once when a clean schedule starts skipping, and recovers on the next clean chain. The schedules list shows a badge. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * feat: alert only on a streak of skipping runs, keep a recent skip visible The skip state now describes the current streak and is written in the push transaction, so it commits or rolls back with the push. The alert fires once when 3 runs in a row skipped, and the list keeps a muted badge for 7 days after the latest skip. Editing or toggling a schedule resets it. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * fix: name the missed-occurrence state after what it counts, alert only once committed Renames the columns to late_run_streak, missed_occurrences and last_missed_at, keeps the missed count after a streak ends so the muted badge can show it, and rewords both badges. The alert task now reads the streak FOR SHARE, which waits for the push transaction, so a push that rolls back and retries alerts once. A failed slot count leaves the streak untouched, and a schedule deleted mid-push no longer fails it. Adds an integration test for the streak and its reset. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * fix: recover the late run alert, store the missed slot, name it missed throughout The alert now recovers (and so acknowledges itself) when a streak that alerted ends on a run on time, under the schedule:{path} resource used by the other trigger alerts. last_missed_at records the last missed cron slot rather than when the late run chained, and the counting helpers say missed, since skipped already names occurrences queued and not run. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * fix: scope the late run alert to its workspace, acknowledge it on edit, toggle and delete Recovery acknowledges alerts by resource alone, so the resource now carries the workspace. Editing, toggling or deleting a schedule clears its streak and a disabled or deleted one never chains a run on time, so those handlers acknowledge its open alert after committing. Past the 1000-slot cap, last_missed_at falls back to the detection time. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 * refactor: raise the late run alert like the other critical alerts Drops the recovery, the workspace-scoped resource and the acknowledgement on edit, toggle and delete: the alert now fires once per streak with no resource and is acknowledged from the alerts feed, as the trigger and job failure alerts are. The FOR SHARE read stays, so a push that rolls back across the flow path's retries still alerts once. Notes in openapi that past 1000 misses in one late run the count is a floor and the time approximate. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LJ9wpjWp2YgLUSqt1Ai5d6 --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
1874 lines
82 KiB
Rust
1874 lines
82 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,
|
|
late_run_streak: 0,
|
|
};
|
|
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,
|
|
build_binary_only: false,
|
|
}
|
|
}
|
|
|
|
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(())
|
|
}
|
|
|
|
// -----------------------------------------------------------------------
|
|
// push_scheduled_job: late run streak
|
|
// -----------------------------------------------------------------------
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
|
|
async fn test_push_tracks_late_run_streak(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
let mut schedule = make_schedule(|s| s.schedule = "0 0 0 * * *".to_string());
|
|
sqlx::query(
|
|
"INSERT INTO schedule (workspace_id, path, edited_by, schedule, script_path, permissioned_as)
|
|
VALUES ($1, $2, 'test-user', $3, $4, 'u/test-user')",
|
|
)
|
|
.bind(&schedule.workspace_id)
|
|
.bind(&schedule.path)
|
|
.bind(&schedule.schedule)
|
|
.bind(&schedule.script_path)
|
|
.execute(&db)
|
|
.await?;
|
|
let authed = make_authed();
|
|
let streak = || {
|
|
sqlx::query_as::<_, (i32, i32, Option<chrono::DateTime<Utc>>)>(
|
|
"SELECT late_run_streak, missed_occurrences, last_missed_at FROM schedule WHERE path = $1",
|
|
)
|
|
.bind(&schedule.path)
|
|
.fetch_one(&db)
|
|
};
|
|
let push = |schedule: Schedule, prev: chrono::DateTime<Utc>| {
|
|
let db = db.clone();
|
|
let authed = authed.clone();
|
|
async move {
|
|
sqlx::query("DELETE FROM v2_job_queue").execute(&db).await?;
|
|
let tx = db.begin().await?;
|
|
push_scheduled_job(&db, tx, &schedule, Some(&authed), Some(prev))
|
|
.await?
|
|
.commit()
|
|
.await?;
|
|
anyhow::Ok(())
|
|
}
|
|
};
|
|
|
|
let midnight = Utc::now()
|
|
.date_naive()
|
|
.and_hms_opt(0, 0, 0)
|
|
.unwrap()
|
|
.and_utc();
|
|
|
|
// The daily run due 3 days ago finished just now: the 3 midnights since were missed.
|
|
push(schedule.clone(), midnight - chrono::Duration::days(3)).await?;
|
|
let (runs, missed, at) = streak().await?;
|
|
assert_eq!((runs, missed, at), (1, 3, Some(midnight)));
|
|
|
|
// Today's run on time ends the streak but keeps what it missed, for the badge.
|
|
schedule.late_run_streak = runs;
|
|
push(schedule.clone(), midnight).await?;
|
|
assert_eq!(streak().await?, (0, 3, at));
|
|
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: a cron with no run left disables the schedule
|
|
// -----------------------------------------------------------------------
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
|
|
async fn test_cron_with_no_run_left_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, cron_version)
|
|
VALUES ('test-workspace', 'f/system/test_schedule', 'test-user', now(), '0 0 9 1 1 * 2020', 'UTC', true, 'f/system/test_script', false, 'test@windmill.dev', '{}', false, false, 'u/test-user', 'v1')"
|
|
)
|
|
.execute(&db)
|
|
.await?;
|
|
|
|
let schedule = make_schedule(|s| {
|
|
s.schedule = "0 0 9 1 1 * 2020".to_string();
|
|
s.cron_version = Some("v1".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;
|
|
assert!(err.is_none(), "completion must go through, got: {err:?}");
|
|
tx.commit().await?;
|
|
|
|
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?;
|
|
assert!(!enabled, "schedule with no run left must be disabled");
|
|
assert!(
|
|
error.as_deref().is_some_and(|e| e.contains("no run left")),
|
|
"error should say why, 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(())
|
|
}
|
|
|
|
#[sqlx::test(migrations = "../migrations", fixtures("base", "schedule_push"))]
|
|
async fn test_push_hub_script_schedule(db: Pool<Postgres>) -> anyhow::Result<()> {
|
|
// Seed the hub cache so the push resolves the script without the network.
|
|
let version = "990000001";
|
|
let hub_dir = &*windmill_common::worker::HUB_CACHE_DIR;
|
|
tokio::fs::create_dir_all(hub_dir).await?;
|
|
tokio::fs::write(
|
|
format!("{hub_dir}/{version}"),
|
|
r#"{"content":"echo hi","lockfile":null,"language":"bash","schema":{},"summary":null}"#,
|
|
)
|
|
.await?;
|
|
let hub_path = format!("hub/{version}/test/echo");
|
|
|
|
// A retry must stay a native `script_hub` job: a flow wrapper would queue the
|
|
// next tick at start and let slow runs overlap.
|
|
let retry = serde_json::json!({ "constant": { "attempts": 2, "seconds": 1 } });
|
|
for (path, retry, dynamic_skip) in [
|
|
("f/system/hub_plain", None, None),
|
|
("f/system/hub_retry", Some(retry), None),
|
|
// A skip handler needs a flow wrapper, so its metadata lookup must not
|
|
// go through the `script` table.
|
|
("f/system/hub_skip", None, Some("f/system/skip".to_string())),
|
|
] {
|
|
let has_retry = retry.is_some();
|
|
let has_skip = dynamic_skip.is_some();
|
|
let schedule = make_schedule(|s| {
|
|
s.path = path.to_string();
|
|
s.script_path = hub_path.clone();
|
|
s.retry = retry;
|
|
s.dynamic_skip = dynamic_skip;
|
|
});
|
|
let tx = db.begin().await?;
|
|
let tx = push_scheduled_job(&db, tx, &schedule, Some(&make_authed()), None).await?;
|
|
tx.commit().await?;
|
|
|
|
let (job_kind, runnable_path, language, handle) =
|
|
sqlx::query_as::<_, (String, Option<String>, Option<String>, Option<i64>)>(
|
|
"SELECT j.kind::text, j.runnable_path, j.script_lang::text, q.runnable_settings_handle
|
|
FROM v2_job j JOIN v2_job_queue q ON j.id = q.id WHERE j.trigger = $1",
|
|
)
|
|
.bind(path)
|
|
.fetch_one(&db)
|
|
.await?;
|
|
assert_eq!(runnable_path.as_deref(), Some(hub_path.as_str()));
|
|
if has_skip {
|
|
assert_eq!(job_kind, "singlestepflow");
|
|
} else {
|
|
assert_eq!(job_kind, "script_hub");
|
|
assert_eq!(language.as_deref(), Some("bash"));
|
|
assert_eq!(handle.is_some(), has_retry);
|
|
}
|
|
}
|
|
Ok(())
|
|
}
|
|
}
|