import { readFile } from "node:fs/promises"; import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; // Persistent state for the cron-driven sync scheduler. Unlike the in-memory job // registry (which is wiped on every server restart), this survives restarts so // failure backoff and the run log are durable. It is small (one entry per // channel plus a bounded run log) and written atomically on each tick. // // This module is pure I/O + types — the scheduling logic lives in // common/jobs/syncScheduler.ts and the orchestration in the editor tick route. export type ChannelSyncState = { // Number of consecutive scheduled syncs that ended in failure/cancel. Reset // to 0 on the first success. Drives the exponential backoff window. consecutiveFailures: number; // Epoch ms before which the channel is held back by backoff. null = eligible. nextEligibleAt: number | null; // The job id of the most recent scheduler-queued sync whose terminal outcome // has not yet been folded into the counters above. null once reconciled. lastQueuedJobId: string | null; // Epoch ms the scheduler last queued a sync for this channel. lastQueuedAt: number | null; // Outcome of the last reconciled scheduled sync, for the observability panel. lastOutcome: "ok" | "failed" | null; lastOutcomeAt: number | null; // Epoch ms the scheduler last queued a keep-latest deletion check for this // channel. Drives the keepLatestCheckIntervalMinutes cadence independently of // the sync cadence. null = never run. lastKeptCheckAt: number | null; }; export type SchedulerSkip = { slug: string; reason: string }; export type SchedulerRun = { at: number; queued: string[]; skipped: SchedulerSkip[]; }; export type SchedulerState = { channels: Record; // Newest-first, bounded to SCHEDULER_RUN_LOG_LIMIT entries. runs: SchedulerRun[]; // Epoch ms the scheduler last queued a saved-video backup (a global, not // per-channel, job). Drives the savedVideoBackup.intervalMinutes cadence. // null = never run. See editor/app/scheduler/runTick.ts. lastSavedVideoBackupAt: number | null; // Epoch ms of the last manual full "Sync all" sweep (syncAllChannelsAction). // A global freshness marker distinct from any single channel's lastSyncedAt; // surfaced by the monitor widget's last-sync readout. null = never run. lastSyncAllAt: number | null; }; export const SCHEDULER_RUN_LOG_LIMIT = 50; export function emptyChannelSyncState(): ChannelSyncState { return { consecutiveFailures: 0, nextEligibleAt: null, lastQueuedJobId: null, lastQueuedAt: null, lastOutcome: null, lastOutcomeAt: null, lastKeptCheckAt: null, }; } export function emptySchedulerState(): SchedulerState { return { channels: {}, runs: [], lastSavedVideoBackupAt: null, lastSyncAllAt: null, }; } // Read the state file, tolerating a missing/corrupt file by returning an empty // state. Defensive: any stored shape that doesn't match is coerced so a // hand-edited or partially-written file can't crash a tick. export async function readSchedulerState(paths: Paths): Promise { let raw: unknown; try { raw = JSON.parse(await readFile(paths.schedulerStateFile, "utf8")); } catch { return emptySchedulerState(); } if (!raw || typeof raw !== "object") return emptySchedulerState(); const r = raw as Record; const channels: Record = {}; if (r.channels && typeof r.channels === "object") { for (const [slug, value] of Object.entries( r.channels as Record, )) { channels[slug] = coerceChannelState(value); } } const runs: SchedulerRun[] = Array.isArray(r.runs) ? r.runs.map(coerceRun).slice(0, SCHEDULER_RUN_LOG_LIMIT) : []; return { channels, runs, lastSavedVideoBackupAt: num(r.lastSavedVideoBackupAt), lastSyncAllAt: num(r.lastSyncAllAt), }; } // Write the state atomically (tmp file + rename), creating the .scheduler dir on // first use. The run log is trimmed to the limit on the way out. export async function writeSchedulerState( paths: Paths, state: SchedulerState, ): Promise { const out: SchedulerState = { channels: state.channels, runs: state.runs.slice(0, SCHEDULER_RUN_LOG_LIMIT), lastSavedVideoBackupAt: state.lastSavedVideoBackupAt ?? null, lastSyncAllAt: state.lastSyncAllAt ?? null, }; await writeJsonAtomic(paths.schedulerStateFile, out, { mkdir: true }); } // Get (or lazily create) the mutable state entry for a channel. export function channelState( state: SchedulerState, slug: string, ): ChannelSyncState { let entry = state.channels[slug]; if (!entry) { entry = emptyChannelSyncState(); state.channels[slug] = entry; } return entry; } // Prepend a run to the log and trim to the limit. export function recordRun(state: SchedulerState, run: SchedulerRun): void { state.runs.unshift(run); if (state.runs.length > SCHEDULER_RUN_LOG_LIMIT) { state.runs.length = SCHEDULER_RUN_LOG_LIMIT; } } function num(value: unknown): number | null { return typeof value === "number" && Number.isFinite(value) ? value : null; } function coerceChannelState(value: unknown): ChannelSyncState { const base = emptyChannelSyncState(); if (!value || typeof value !== "object") return base; const r = value as Record; return { consecutiveFailures: Math.max(0, num(r.consecutiveFailures) ?? 0), nextEligibleAt: num(r.nextEligibleAt), lastQueuedJobId: typeof r.lastQueuedJobId === "string" ? r.lastQueuedJobId : null, lastQueuedAt: num(r.lastQueuedAt), lastOutcome: r.lastOutcome === "ok" || r.lastOutcome === "failed" ? r.lastOutcome : null, lastOutcomeAt: num(r.lastOutcomeAt), lastKeptCheckAt: num(r.lastKeptCheckAt), }; } function coerceRun(value: unknown): SchedulerRun { const r = (value ?? {}) as Record; const queued = Array.isArray(r.queued) ? r.queued.filter((x): x is string => typeof x === "string") : []; const skipped: SchedulerSkip[] = Array.isArray(r.skipped) ? r.skipped.flatMap((x) => { if (!x || typeof x !== "object") return []; const s = x as Record; if (typeof s.slug !== "string" || typeof s.reason !== "string") { return []; } return [{ slug: s.slug, reason: s.reason }]; }) : []; return { at: num(r.at) ?? 0, queued, skipped }; }