import { Job } from "../db/models"; import { runJob } from "./runner"; import { getConfig } from "../utils/config"; import { log } from "../utils/log"; import { computeInitialNextRun, computeNextRun } from "../utils/schedule"; import type { JobResult } from "../types"; export { computeInitialNextRun, computeNextRun }; /** * Log a finished job at a level matching its outcome. `runJob` resolves with * `status: "error"` instead of throwing, so the caller's `.catch` only ever sees * an infrastructure fault — a job that ran and failed must be branched on here * or it is indistinguishable from success in the log. */ export function logJobOutcome(result: JobResult): void { const fields = { job: result.job, status: result.status, duration: result.duration_ms }; if (result.status === "error") { log.error({ ...fields, error: result.error, terminal_reason: result.terminal_reason }, "scheduler: job failed"); return; } log.info(fields, "scheduler: job completed"); } function isWithinActiveHours(): boolean { const config = getConfig(); const { start, end } = config.activeHours; const now = new Date(); const formatter = new Intl.DateTimeFormat("en-US", { hour: "2-digit", minute: "2-digit", hour12: false, timeZone: config.timezone, }); const current = formatter.format(now).replace(/\u200e/g, ""); // Handle midnight-crossing windows (e.g. 09:00–02:00) if (start <= end) { return current >= start && current <= end; } return current >= start || current <= end; } let timer: ReturnType | null = null; const runningJobs = new Set(); async function tick(): Promise { let dueJobs: Awaited>; try { dueJobs = await Job.listDue(); } catch (err) { log.warn({ err }, "scheduler: failed to query due jobs"); return; } const config = getConfig(); for (const job of dueJobs) { if (!job.always && !isWithinActiveHours()) { // Leave next_run_at untouched so the job stays due and fires as soon as // active hours resume. Rescheduling here would advance a cron job whose // only fire time sits outside the window to the next (also-outside) // occurrence — starving it forever. log.info({ job: job.name }, "scheduler: skipping — outside active hours"); continue; } if (runningJobs.has(job.name)) { log.info({ job: job.name }, "scheduler: skipping — still running from previous invocation"); continue; } log.info({ job: job.name, type: job.scheduleType }, "scheduler: running job"); runningJobs.add(job.name); runJob(job) .then(logJobOutcome) .catch((err) => { log.error({ err, job: job.name }, "scheduler: job crashed"); }) .finally(() => { runningJobs.delete(job.name); }); let nextRun: Date | null = null; try { nextRun = computeNextRun(job.scheduleType, job.schedule, config.timezone, new Date()); } catch (err) { log.error({ err, job: job.name, schedule: job.schedule }, "scheduler: invalid schedule"); try { await Job.update(job.name, { status: "disabled" }); log.info({ job: job.name }, "scheduler: disabled job with invalid schedule"); } catch (updateErr) { log.error( { err: updateErr, job: job.name }, "scheduler: could not disable job with invalid schedule — it will keep failing every tick", ); } continue; } await Job.markRun(job.name, nextRun).catch((err) => { log.error({ err, job: job.name }, "scheduler: failed to update next_run_at"); }); // Auto-disable one-shot jobs after execution if (job.scheduleType === "once") { try { await Job.update(job.name, { status: "disabled" }); log.info({ job: job.name }, "scheduler: one-shot job completed, auto-disabled"); } catch (err) { log.error( { err, job: job.name }, "scheduler: one-shot job completed but could not be disabled — it will run again", ); } } } } export function startScheduler(): void { log.info("scheduler started (60s poll interval)"); tick(); timer = setInterval(tick, 60_000); } export function stopScheduler(): void { if (timer) { clearInterval(timer); timer = null; } } export async function recomputeAllNextRuns(): Promise { const config = getConfig(); const jobs = await Job.listEnabled(); const { getSql } = await import("../db/connection"); const sql = getSql(); for (const job of jobs) { if (job.nextRunAt) continue; const nextRun = computeInitialNextRun(job.scheduleType, job.schedule, config.timezone); await sql`UPDATE jobs SET next_run_at = ${nextRun} WHERE name = ${job.name}`; } }