/** * Standalone cron worker for AtomCMS-Next. * * Replaces Laravel's scheduler (app/Console/Kernel.php + queue:work) with a * single long-running process driven by `croner`. Run it OUTSIDE the Next.js * request lifecycle (e.g. a separate container / pm2 process): * * pnpm jobs:worker # -> tsx scripts/jobs-worker.ts * * Each job body is wrapped in its own try/catch so one failing tick never * tears down the scheduler; a thrown error or a falsy result is logged and the * worker keeps ticking. Croner runs ticks serially per-job (protect: true) so a * slow run can't overlap itself. * * DEPENDENCIES THIS FILE EXPECTS (orchestrator must provide them — noted here * so the build isn't silently broken): * - `croner` package (not yet in package.json): `pnpm add croner`. * - `@/lib/services/alert` exporting `alert.emulatorOffline()` — a thin * notifier (email/Discord) for emulator-down events. Not present yet; this * worker is the first consumer. */ import "dotenv/config"; import { copyFile, mkdir, readdir, stat, unlink } from "node:fs/promises"; import { join } from "node:path"; import { Cron } from "croner"; import { prisma } from "@/lib/prisma"; import { emulatorOffline } from "@/lib/services/alert"; import { fetchNowPlaying } from "@/lib/services/radio"; import { rcon } from "@/lib/services/rcon"; import { siteSettings } from "@/lib/services/site-settings"; // --- small logging helper (timestamped, namespaced) ------------------------ function log(scope: string, msg: string, ...rest: unknown[]): void { console.log(`[jobs:${scope}] ${new Date().toISOString()} ${msg}`, ...rest); } function logErr(scope: string, msg: string, e: unknown): void { console.error(`[jobs:${scope}] ${new Date().toISOString()} ${msg}`, e instanceof Error ? e.message : e); } // --- (a) emulator health ping — every minute ------------------------------- // // Our RconClient has no `isConnected`: the TCP transport is fire-and-forget and // resolves `false` on a dead/unreachable socket (it never throws for that). So // we treat BOTH a thrown error AND a falsy `send()` result as "offline" and // fire the alert. `data: null` matches the no-payload command shape AtomCMS // uses for keep-alive style commands. async function pingEmulator(): Promise { let online = false; try { online = await rcon.send("ping", null); } catch (e) { logErr("ping", "rcon.send threw", e); online = false; } if (online) { log("ping", "emulator reachable"); return; } log("ping", "emulator unreachable — raising offline alert"); try { await emulatorOffline(); } catch (e) { // Never let the alert channel failing crash the tick. logErr("ping", "emulatorOffline failed", e); } } // --- (b) maintenance:check — every minute ---------------------------------- // // website_maintenance_tasks has no scheduled-time column (id, userId, task, // completed, timestamps), so "scheduled maintenance" == at least one row with // completed = false. When such a task exists we flip the `maintenance_enabled` // website_setting on; otherwise we clear it. The setting value is the canonical // '1' / '0' string the rest of the CMS reads (see SiteSettings.getBool). // // If the model isn't present in the running schema (older DB), the prisma call // throws and we skip — the catch keeps the worker alive. async function maintenanceCheck(): Promise { let pending = 0; try { pending = await prisma.websiteMaintenanceTasks.count({ where: { completed: false }, }); } catch (e) { // Table absent / DB unreachable — skip this tick rather than thrashing. logErr("maintenance", "could not read website_maintenance_tasks — skipping", e); return; } const desired = pending > 0 ? "1" : "0"; try { const current = await prisma.websiteSetting.findUnique({ where: { key: "maintenance_enabled" }, select: { value: true }, }); if (current?.value === desired) { log("maintenance", `no change (maintenance_enabled=${desired}, pending=${pending})`); return; } await prisma.websiteSetting.upsert({ where: { key: "maintenance_enabled" }, update: { value: desired }, create: { key: "maintenance_enabled", value: desired, comment: "Toggled by jobs-worker" }, }); log("maintenance", `maintenance_enabled -> ${desired} (pending tasks: ${pending})`); } catch (e) { logErr("maintenance", "failed to toggle maintenance_enabled", e); } } // --- (c) bans cleanup — hourly --------------------------------------------- // // Bans expire purely by the query-time `ban_expire > now()` filter applied // wherever bans are read, so there's nothing to delete on a schedule. We keep // the hourly tick for observability (and an obvious hook if a future hard-purge // is ever wanted). async function bansCleanup(): Promise { try { const now = Math.floor(Date.now() / 1000); const active = await prisma.ban.count({ where: { banExpire: { gt: now } } }); log("bans", `cleanup is automatic via ban_expire>now filter — ${active} active ban(s), nothing to purge`); } catch (e) { logErr("bans", "count failed (non-fatal)", e); } } // --- (d) radio: record song plays — every 30s ------------------------------ // // Polls the configured now-playing endpoint (AtomCMS radio:record-songs). On a // track CHANGE it appends a row to radio_song_plays so the public history + // admin pages have data. De-dups against the most recent recorded play. let lastRecordedTitle: string | null = null; async function recordSongPlay(): Promise { let np: Awaited>; try { np = await fetchNowPlaying(); } catch (e) { logErr("radio", "now-playing fetch threw", e); return; } if (!np || !np.title) return; if (np.title === lastRecordedTitle) return; try { const latest = await prisma.radioSongPlays.findFirst({ orderBy: { id: "desc" }, select: { title: true }, }); if (latest?.title === np.title) { lastRecordedTitle = np.title; return; } await prisma.radioSongPlays.create({ data: { title: np.title, artist: np.artist, playedAt: new Date(), createdAt: new Date() }, }); lastRecordedTitle = np.title; log("radio", `recorded play: ${np.artist ? `${np.artist} - ` : ""}${np.title}`); } catch (e) { logErr("radio", "could not record song play", e); } } // --- (e) radio: auto-DJ rotation — every minute ---------------------------- // // When no live DJ is broadcasting (radio_current_dj_id empty) and auto-DJ is // enabled, advances the radio_auto_dj_playlist by sort_order and writes the // current track to the radio_now_playing setting (read by the player/API). async function autoDj(): Promise { try { if (!(await siteSettings.getBool("radio_auto_dj_enabled", false))) return; const liveDj = (await siteSettings.get("radio_current_dj_id", "")) ?? ""; if (liveDj) return; // a real DJ is on air — don't override const tracks = await prisma.radioAutoDjPlaylist.findMany({ where: { isActive: true }, orderBy: [{ sortOrder: "asc" }, { id: "asc" }], select: { id: true, title: true, artist: true, playCount: true }, }); if (tracks.length === 0) return; // Pick the least-recently-played (lowest playCount, then lowest id). const next = tracks.slice().sort((a, b) => a.playCount - b.playCount || Number(a.id - b.id))[0]; await prisma.radioAutoDjPlaylist.update({ where: { id: next.id }, data: { playCount: { increment: 1 }, lastPlayedAt: new Date() }, }); const label = `${next.artist ? `${next.artist} - ` : ""}${next.title}`; await prisma.websiteSetting.upsert({ where: { key: "radio_now_playing" }, update: { value: label }, create: { key: "radio_now_playing", value: label, comment: "Auto-DJ now playing" }, }); log("radio", `auto-DJ now playing: ${label}`); } catch (e) { logErr("radio", "auto-DJ tick failed", e); } } // --- (f) github: update check — hourly ------------------------------------- // // Compares the latest commit on the configured GitHub repo against the last one // seen, and stores update_available + update_latest_sha settings the admin // dashboard can surface. No-ops unless github_repo (owner/repo) is set. async function githubUpdateCheck(): Promise { try { const repo = (await siteSettings.get("github_repo", "")) ?? ""; if (!/^[\w.-]+\/[\w.-]+$/.test(repo)) return; const branch = (await siteSettings.get("github_branch", "main")) ?? "main"; const res = await fetch(`https://api.github.com/repos/${repo}/commits/${branch}`, { headers: { accept: "application/vnd.github+json", "user-agent": "atomcms-next" }, cache: "no-store", }); if (!res.ok) return; const data = (await res.json()) as { sha?: string }; const sha = data?.sha; if (!sha) return; const known = (await siteSettings.get("update_current_sha", "")) ?? ""; const available = known ? known !== sha ? "1" : "0" : "0"; for (const [key, value] of [ ["update_latest_sha", sha], ["update_available", available], ] as const) { await prisma.websiteSetting.upsert({ where: { key }, update: { value }, create: { key, value, comment: "GitHub update check" }, }); } log("github", `latest ${sha.slice(0, 7)} (update ${available === "1" ? "AVAILABLE" : "none"})`); } catch (e) { logErr("github", "update check failed", e); } } // --- (g) emulator: JAR backup — daily -------------------------------------- // // AtomCMS's backup tooling lives in the scheduler (host-side), which is exactly // what this worker is — so it CAN do the filesystem copy the Next request tier // cannot. Copies EMULATOR_JAR_PATH into EMULATOR_BACKUP_DIR with a timestamped // name, keeping the newest EMULATOR_BACKUP_KEEP (default 7). No-ops cleanly when // the env paths aren't set (e.g. on a dev box). async function emulatorBackup(): Promise { const jar = process.env.EMULATOR_JAR_PATH; const dir = process.env.EMULATOR_BACKUP_DIR; if (!jar || !dir) return; // not configured — nothing to do try { await stat(jar); // ensure the source exists await mkdir(dir, { recursive: true }); const stamp = new Date().toISOString().replace(/[:.]/g, "-"); await copyFile(jar, join(dir, `emulator-${stamp}.jar`)); // Prune to the newest N backups. const keep = Number(process.env.EMULATOR_BACKUP_KEEP ?? "7") || 7; const files = (await readdir(dir)) .filter((f) => f.startsWith("emulator-") && f.endsWith(".jar")) .sort() .reverse(); for (const old of files.slice(keep)) { await unlink(join(dir, old)).catch(() => {}); } log("backup", `emulator JAR backed up (${files.length + 1} kept, pruning to ${keep})`); } catch (e) { logErr("backup", "emulator backup failed", e); } } // --- scheduler wiring ------------------------------------------------------ const jobs: Cron[] = [ new Cron("* * * * *", { name: "emulator-ping", protect: true }, pingEmulator), new Cron("* * * * *", { name: "maintenance-check", protect: true }, maintenanceCheck), new Cron("0 * * * *", { name: "bans-cleanup", protect: true }, bansCleanup), new Cron("*/30 * * * * *", { name: "radio-record-songs", protect: true }, recordSongPlay), new Cron("* * * * *", { name: "radio-auto-dj", protect: true }, autoDj), new Cron("0 * * * *", { name: "github-update-check", protect: true }, githubUpdateCheck), new Cron("0 4 * * *", { name: "emulator-backup", protect: true }, emulatorBackup), ]; log("worker", `started — ${jobs.length} scheduled job(s): ${jobs.map((j) => j.name).join(", ")}`); // Run once immediately on boot so we don't wait up to a minute for first signal. void pingEmulator(); void maintenanceCheck(); // --- graceful shutdown ----------------------------------------------------- async function shutdown(signal: string): Promise { log("worker", `received ${signal} — stopping jobs and disconnecting`); for (const j of jobs) j.stop(); try { await prisma.$disconnect(); } catch { // ignore — we're exiting anyway } process.exit(0); } process.on("SIGINT", () => void shutdown("SIGINT")); process.on("SIGTERM", () => void shutdown("SIGTERM")); // Don't let an unexpected async throw kill the whole worker. process.on("unhandledRejection", (reason) => { logErr("worker", "unhandledRejection", reason); });