diff --git a/backend/.sqlx/query-7d67aebe32cc5310a92f510bd2dd742198c8a492783348d457e6811601b21371.json b/backend/.sqlx/query-7d67aebe32cc5310a92f510bd2dd742198c8a492783348d457e6811601b21371.json deleted file mode 100644 index f183dc3267..0000000000 --- a/backend/.sqlx/query-7d67aebe32cc5310a92f510bd2dd742198c8a492783348d457e6811601b21371.json +++ /dev/null @@ -1,57 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT\n schedule.path, schedule.schedule, schedule.timezone,\n schedule.cron_version, t.jobs, p.push_times FROM schedule,\n LATERAL(SELECT ARRAY(\n SELECT json_build_object('id', id, 'success', status = 'success', 'duration_ms', duration_ms)\n FROM v2_job_completed c JOIN v2_job j USING (id)\n WHERE trigger_kind = 'schedule'\n AND trigger = schedule.path\n AND c.workspace_id = $1\n AND j.workspace_id = $1\n AND parent_job IS NULL AND runnable_path = schedule.script_path\n AND status <> 'skipped'\n ORDER BY completed_at DESC\n LIMIT 20\n ) AS jobs) t,\n LATERAL(SELECT ARRAY(\n SELECT created_at FROM (\n SELECT created_at, trigger, trigger_kind\n FROM v2_job\n WHERE schedule.enabled\n AND v2_job.workspace_id = $1\n AND parent_job IS NULL AND runnable_path = schedule.script_path\n AND created_at > schedule.edited_at\n ORDER BY created_at DESC\n LIMIT $6\n ) recent\n WHERE trigger_kind = 'schedule' AND trigger = schedule.path\n ORDER BY created_at DESC\n LIMIT $5\n ) AS push_times) p\n WHERE workspace_id = $1 AND NOT starts_with(schedule.path, $4)\n ORDER BY edited_at DESC\n LIMIT $2 OFFSET $3", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "path", - "type_info": "Varchar" - }, - { - "ordinal": 1, - "name": "schedule", - "type_info": "Varchar" - }, - { - "ordinal": 2, - "name": "timezone", - "type_info": "Varchar" - }, - { - "ordinal": 3, - "name": "cron_version", - "type_info": "Text" - }, - { - "ordinal": 4, - "name": "jobs", - "type_info": "JsonArray" - }, - { - "ordinal": 5, - "name": "push_times", - "type_info": "TimestamptzArray" - } - ], - "parameters": { - "Left": [ - "Text", - "Int8", - "Int8", - "Text", - "Int8", - "Int8" - ] - }, - "nullable": [ - false, - false, - false, - true, - null, - null - ] - }, - "hash": "7d67aebe32cc5310a92f510bd2dd742198c8a492783348d457e6811601b21371" -} diff --git a/backend/.sqlx/query-a9db62e3a7465feec93ff9b09142807ed1db0bb40ca8a7d345afbaaf15ef5d16.json b/backend/.sqlx/query-a9db62e3a7465feec93ff9b09142807ed1db0bb40ca8a7d345afbaaf15ef5d16.json new file mode 100644 index 0000000000..070d5c5f44 --- /dev/null +++ b/backend/.sqlx/query-a9db62e3a7465feec93ff9b09142807ed1db0bb40ca8a7d345afbaaf15ef5d16.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT created_at\n FROM trigger_history\n WHERE workspace_id = $1 AND trigger_kind = $2 AND path = $3\n AND (changes ?| array['schedule', 'timezone', 'cron_version']\n OR changes -> 'truncated_fields' ?| array['schedule', 'timezone', 'cron_version'])\n ORDER BY id DESC\n LIMIT 1", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "created_at", + "type_info": "Timestamptz" + } + ], + "parameters": { + "Left": [ + "Text", + "Text", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "a9db62e3a7465feec93ff9b09142807ed1db0bb40ca8a7d345afbaaf15ef5d16" +} diff --git a/backend/.sqlx/query-c2aa552cf5f74d1a1305709c576425c3749b9c66cef9a4c4675ed6f6dfe72d6f.json b/backend/.sqlx/query-c2aa552cf5f74d1a1305709c576425c3749b9c66cef9a4c4675ed6f6dfe72d6f.json deleted file mode 100644 index 69f90f0079..0000000000 --- a/backend/.sqlx/query-c2aa552cf5f74d1a1305709c576425c3749b9c66cef9a4c4675ed6f6dfe72d6f.json +++ /dev/null @@ -1,30 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT DISTINCT ON (path) path, created_at\n FROM trigger_history\n WHERE workspace_id = $1 AND trigger_kind = $2 AND path = ANY($3)\n AND (changes ?| array['schedule', 'timezone', 'cron_version']\n OR changes -> 'truncated_fields' ?| array['schedule', 'timezone', 'cron_version'])\n ORDER BY path, id DESC", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "path", - "type_info": "Varchar" - }, - { - "ordinal": 1, - "name": "created_at", - "type_info": "Timestamptz" - } - ], - "parameters": { - "Left": [ - "Text", - "Text", - "TextArray" - ] - }, - "nullable": [ - false, - false - ] - }, - "hash": "c2aa552cf5f74d1a1305709c576425c3749b9c66cef9a4c4675ed6f6dfe72d6f" -} diff --git a/backend/windmill-api-schedule/src/lib.rs b/backend/windmill-api-schedule/src/lib.rs index aeac916edc..ee036338a5 100644 --- a/backend/windmill-api-schedule/src/lib.rs +++ b/backend/windmill-api-schedule/src/lib.rs @@ -1000,8 +1000,6 @@ async fn list_schedule( pub struct ScheduleWJobs { pub path: String, pub jobs: Option>, - #[serde(skip_serializing_if = "Option::is_none")] - pub interval_drift: Option, } /// How often a schedule's runs are actually landing, reported only while that @@ -1025,9 +1023,11 @@ const DRIFT_SAMPLE_SIZE: i64 = DRIFT_MIN_MISSED_SLOTS as i64 + 1; /// Root jobs of the runnable the sample may walk before giving up. There is no /// index on `trigger`, so the sample rides the one on `runnable_path` and finds /// the schedule's own runs by filter: a runnable that is also busy outside this -/// schedule would otherwise be read back as far as its fourth scheduled run. -/// Past the budget the schedule reports nothing rather than costing the walk. -const DRIFT_SCAN_BUDGET: i64 = 100; +/// schedule would otherwise be read back as far as its fourth scheduled run, +/// which for a slow schedule on a busy script is unbounded. This covers that +/// case with room to spare, and is affordable because the sample is only taken +/// when someone opens one schedule, never once per row of a listing. +const DRIFT_SCAN_BUDGET: i64 = 2000; /// `push_times` are the moments the last runs were queued, newest first. The /// pusher anchors `find_next` on exactly that moment, so recomputing from it @@ -1082,8 +1082,7 @@ fn detect_interval_drift( mildest } -/// When each schedule's cron expression last changed, for the paths asked -/// about and only those that ever recorded such a change. +/// When this schedule's cron expression last changed, if it ever has. /// /// Runs queued under a previous expression measure as drift against the one /// that replaced it, so a sample reaching back past the change has to be @@ -1091,66 +1090,47 @@ fn detect_interval_drift( /// after the insert. The trigger history can, and is only read once a sample /// has already come out as drifting. /// -/// Compare strictly: the edit queues a tick of its own in the same +/// Compare strictly against it: the edit queues a tick of its own in the same /// transaction, so that first run under the new expression carries the very /// timestamp recorded here. /// -/// Reads on the privileged pool: the paths reaching here have already come -/// back from an RLS read of their own schedule row and passed the caller's -/// read scope. -async fn cron_changed_at( - db: &DB, - w_id: &str, - paths: &[String], -) -> Result)>> { - let rows = sqlx::query!( +/// Reads on the privileged pool, which the caller has earned with an RLS read +/// of the schedule row itself. +async fn cron_changed_at(db: &DB, w_id: &str, path: &str) -> Result>> { + let changed_at = sqlx::query_scalar!( // An entry too large to record field by field keeps only the names, so both // shapes have to be asked. - "SELECT DISTINCT ON (path) path, created_at + "SELECT created_at FROM trigger_history - WHERE workspace_id = $1 AND trigger_kind = $2 AND path = ANY($3) + WHERE workspace_id = $1 AND trigger_kind = $2 AND path = $3 AND (changes ?| array['schedule', 'timezone', 'cron_version'] OR changes -> 'truncated_fields' ?| array['schedule', 'timezone', 'cron_version']) - ORDER BY path, id DESC", + ORDER BY id DESC + LIMIT 1", w_id, SCHEDULE_TRIGGER_KIND, - paths + path ) - .fetch_all(db) + .fetch_optional(db) .await?; - Ok(rows.into_iter().map(|r| (r.path, r.created_at)).collect()) + Ok(changed_at) } async fn list_schedule_with_jobs( authed: ApiAuthed, Extension(user_db): Extension, - Extension(db): Extension, Path(w_id): Path, Query(pagination): Query, ) -> JsonResult> { let mut tx = user_db.begin(&authed).await?; let (per_page, offset) = paginate(pagination); - struct ScheduleWJobsRow { - path: String, - jobs: Option>, - push_times: Option>>, - schedule: String, - timezone: String, - cron_version: Option, - } - let rows = sqlx::query_as!(ScheduleWJobsRow, + let rows = sqlx::query_as!(ScheduleWJobs, // Query plan: // - use of the `ix_completed_job_workspace_id_started_at_new_2` index first, then; // - use of the `ix_v2_job_root_by_path` index; hence the `parent_job IS NULL` clause. // - both `workspace_id = $1` checks are required to hit both indexes. - // - `push_times` rides the `(workspace_id, runnable_path, created_at DESC)` index, - // hence its own `parent_job IS NULL` clause, and is capped by DRIFT_SCAN_BUDGET. - // The `edited_at` bound is not that cap: it drops the runs of whatever schedule - // last held this path, which carry the same `trigger`. Nothing writes `edited_at` - // after the insert, so no run of this schedule can predate it. "SELECT - schedule.path, schedule.schedule, schedule.timezone, - schedule.cron_version, t.jobs, p.push_times FROM schedule, + schedule.path, t.jobs FROM schedule, LATERAL(SELECT ARRAY( SELECT json_build_object('id', id, 'success', status = 'success', 'duration_ms', duration_ms) FROM v2_job_completed c JOIN v2_job j USING (id) @@ -1162,71 +1142,22 @@ async fn list_schedule_with_jobs( AND status <> 'skipped' ORDER BY completed_at DESC LIMIT 20 - ) AS jobs) t, - LATERAL(SELECT ARRAY( - SELECT created_at FROM ( - SELECT created_at, trigger, trigger_kind - FROM v2_job - WHERE schedule.enabled - AND v2_job.workspace_id = $1 - AND parent_job IS NULL AND runnable_path = schedule.script_path - AND created_at > schedule.edited_at - ORDER BY created_at DESC - LIMIT $6 - ) recent - WHERE trigger_kind = 'schedule' AND trigger = schedule.path - ORDER BY created_at DESC - LIMIT $5 - ) AS push_times) p + ) AS jobs) t WHERE workspace_id = $1 AND NOT starts_with(schedule.path, $4) ORDER BY edited_at DESC LIMIT $2 OFFSET $3", w_id, per_page as i64, offset as i64, - windmill_common::workspaces::DUCKLAKE_MAINTENANCE_PATH_PREFIX, - DRIFT_SAMPLE_SIZE, - DRIFT_SCAN_BUDGET + windmill_common::workspaces::DUCKLAKE_MAINTENANCE_PATH_PREFIX ) .fetch_all(&mut *tx) .await?; tx.commit().await?; let allowed = build_scope_path_predicate(&authed, "schedules", "read"); - - let mut out: Vec = Vec::with_capacity(rows.len()); - // Oldest run each reported drift rests on, kept to check it against the - // cron's own history below. - let mut drifting: Vec<(String, DateTime)> = Vec::new(); - for r in rows.into_iter().filter(|r| allowed(&r.path)) { - let sample = r.push_times.unwrap_or_default(); - let interval_drift = - detect_interval_drift(&sample, &r.schedule, r.cron_version.as_deref(), &r.timezone); - if interval_drift.is_some() { - if let Some(oldest) = sample.last() { - drifting.push((r.path.clone(), *oldest)); - } - } - out.push(ScheduleWJobs { path: r.path, jobs: r.jobs, interval_drift }); - } - - if !drifting.is_empty() { - let paths = drifting - .iter() - .map(|(path, _)| path.clone()) - .collect::>(); - for (path, changed_at) in cron_changed_at(&db, &w_id, &paths).await? { - if drifting - .iter() - .any(|(p, oldest)| *p == path && *oldest < changed_at) - { - if let Some(row) = out.iter_mut().find(|r| r.path == path) { - row.interval_drift = None; - } - } - } - } - - Ok(Json(out)) + Ok(Json( + rows.into_iter().filter(|r| allowed(&r.path)).collect(), + )) } // SELECT id, title AS item_title, t.tag_array @@ -1324,8 +1255,8 @@ async fn fetch_interval_drift( ); if drift.is_some() { if let Some(oldest) = push_times.last() { - let changed = cron_changed_at(db, w_id, std::slice::from_ref(&schedule.path)).await?; - if changed.iter().any(|(_, changed_at)| *oldest < *changed_at) { + let changed_at = cron_changed_at(db, w_id, &schedule.path).await?; + if changed_at.is_some_and(|changed_at| *oldest < changed_at) { return Ok(None); } } diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index b931657d8c..a0fc15f1af 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -29775,8 +29775,6 @@ components: - id - success - duration_ms - interval_drift: - $ref: "#/components/schemas/ScheduleIntervalDrift" ScheduleIntervalDrift: type: object diff --git a/frontend/src/lib/components/triggers/schedules/ScheduleEditorInner.svelte b/frontend/src/lib/components/triggers/schedules/ScheduleEditorInner.svelte index 7ffa28cb2b..ec103e8862 100644 --- a/frontend/src/lib/components/triggers/schedules/ScheduleEditorInner.svelte +++ b/frontend/src/lib/components/triggers/schedules/ScheduleEditorInner.svelte @@ -511,6 +511,9 @@ async function loadSchedule( defaultCfg?: Record ): Promise<{ overlay: Record | undefined; noDeployed: boolean }> { + // Reading the deployed row is the only thing that carries a measurement, and + // it omits the field when the schedule is keeping up. + intervalDrift = undefined if (defaultCfg) { await loadScheduleCfg(defaultCfg) return { overlay: undefined, noDeployed: false } @@ -539,7 +542,11 @@ async function loadScheduleCfg(cfg: Record): Promise { loading = true - intervalDrift = cfg.interval_drift ?? undefined + // Config snapshots from draftSync carry no measurement, and applying one + // leaves the deployed schedule running exactly as it was. + if ('interval_drift' in cfg) { + intervalDrift = cfg.interval_drift + } cronVersion = cfg.cron_version ?? 'v2' initialCronVersion = cronVersion isLatestCron = cronVersion == 'v2' diff --git a/frontend/src/routes/(root)/(logged)/schedules/+page.svelte b/frontend/src/routes/(root)/(logged)/schedules/+page.svelte index cad67f5858..e403c981de 100644 --- a/frontend/src/routes/(root)/(logged)/schedules/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/schedules/+page.svelte @@ -6,14 +6,7 @@ type WorkspaceDeployUISettings, WorkspaceService } from '$lib/gen' - import { - canWrite, - displayDate, - getLocalSetting, - msToReadableTime, - msToReadableTimeShort, - storeLocalSetting - } from '$lib/utils' + import { canWrite, displayDate, getLocalSetting, storeLocalSetting } from '$lib/utils' import { withForkConflictRetry } from '$lib/utils/forkConflict' import { base } from '$app/paths' import CenteredPage from '$lib/components/CenteredPage.svelte' @@ -156,7 +149,6 @@ for (let schedule of schedules) { if (schedulesWithJobsByPath[schedule.path]) { schedule.jobs = schedulesWithJobsByPath[schedule.path].jobs - schedule.interval_drift = schedulesWithJobsByPath[schedule.path].interval_drift } } loadingSchedulesWithJobStats = false @@ -408,7 +400,7 @@ {/if} {:else if items?.length}
- {#each items.slice(0, nbDisplayed) as { path, error, summary, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, extra_perms, canWrite, jobs, interval_drift, paused_until, labels, inherited_labels, draft_only, is_draft } (path)} + {#each items.slice(0, nbDisplayed) as { path, error, summary, edited_by, edited_at, schedule, timezone, enabled, script_path, is_flow, extra_perms, canWrite, jobs, paused_until, labels, inherited_labels, draft_only, is_draft } (path)} {@const hasDraft = getLocalDraftHint($workspaceStore, 'trigger_schedule', path) ?? is_draft} {@const href = `${is_flow ? '/flows/get' : '/scripts/get'}/${script_path}`} @@ -474,23 +466,6 @@