import { Cron } from "croner"; import { and, eq, isNotNull, lt, or } from "drizzle-orm"; import { db } from "@/lib/db/client"; import { availabilityOffers, seasons, titles, userTitleStatus, } from "@/lib/db/schema"; import { createLogger } from "@/lib/logger"; import { refreshAvailability } from "@/lib/services/availability"; import { createBackup, ensureBackupDir, pruneBackups, } from "@/lib/services/backup"; import { cacheEpisodeStills, cacheImagesForTitle, cacheProviderLogos, imageCacheEnabled, } from "@/lib/services/image-cache"; import { refreshRecommendations, refreshTitle, refreshTvChildren, } from "@/lib/services/metadata"; import { getSetting } from "@/lib/services/settings"; import { getTvDetails } from "@/lib/tmdb/client"; const log = createLogger("cron"); const DAY = 24 * 60 * 60 * 1000; const RATE_LIMIT_MS = 300; const globalForJobs = globalThis as unknown as { _jobs: Map | undefined; }; if (!globalForJobs._jobs) { globalForJobs._jobs = new Map(); } const jobs = globalForJobs._jobs; function schedule(name: string, cron: string, handler: () => Promise) { jobs.set( name, new Cron( cron, { name, protect: true, catch: (err: unknown) => { log.error(`Job ${name} failed:`, err); }, }, async () => { log.info(`Running job: ${name}`); await handler(); log.info(`Completed job: ${name}`); }, ), ); } function delay(ms: number) { return new Promise((r) => setTimeout(r, ms)); } function getLibraryTitleIds(): string[] { const rows = db .select({ titleId: userTitleStatus.titleId }) .from(userTitleStatus) .groupBy(userTitleStatus.titleId) .all(); return rows.map((r) => r.titleId); } // Refresh titles where lastFetchedAt is stale async function nightlyRefreshLibrary() { const libraryIds = getLibraryTitleIds(); const libraryStale = new Date(Date.now() - 7 * DAY); const nonLibraryStale = new Date(Date.now() - 30 * DAY); // Library titles: 7 days for (const titleId of libraryIds) { const t = db .select() .from(titles) .where( and(eq(titles.id, titleId), lt(titles.lastFetchedAt, libraryStale)), ) .get(); if (t) { await refreshTitle(titleId); await delay(RATE_LIMIT_MS); } } // Non-library titles: 30 days const nonLibrary = db .select() .from(titles) .where( and( isNotNull(titles.lastFetchedAt), lt(titles.lastFetchedAt, nonLibraryStale), ), ) .limit(50) .all(); for (const t of nonLibrary) { if (!libraryIds.includes(t.id)) { await refreshTitle(t.id); await delay(RATE_LIMIT_MS); } } } // Refresh availability for library titles where stale async function refreshAvailabilityJob() { const libraryIds = getLibraryTitleIds(); const stale = new Date(Date.now() - DAY); for (const titleId of libraryIds) { // Check if any offer is stale const offer = db .select() .from(availabilityOffers) .where( and( eq(availabilityOffers.titleId, titleId), lt(availabilityOffers.lastFetchedAt, stale), ), ) .get(); // Also handle titles with no offers yet const anyOffer = db .select() .from(availabilityOffers) .where(eq(availabilityOffers.titleId, titleId)) .get(); if (offer || !anyOffer) { await refreshAvailability(titleId); await delay(RATE_LIMIT_MS); } } } // Refresh recommendations for recently active titles async function refreshRecommendationsJob() { const libraryIds = getLibraryTitleIds(); for (const titleId of libraryIds) { await refreshRecommendations(titleId); await delay(RATE_LIMIT_MS); } } // Refresh TV episodes for returning shows async function refreshTvChildrenJob() { const returningStatuses = ["Returning Series", "In Production"]; const stale = new Date(Date.now() - 7 * DAY); const tvShows = db .select() .from(titles) .where( and( eq(titles.type, "tv"), isNotNull(titles.lastFetchedAt), or(...returningStatuses.map((s) => eq(titles.status, s))), ), ) .all(); for (const show of tvShows) { // Check if seasons are stale const staleSeason = db .select() .from(seasons) .where( and(eq(seasons.titleId, show.id), lt(seasons.lastFetchedAt, stale)), ) .get(); if (staleSeason) { const details = await getTvDetails(show.tmdbId); await refreshTvChildren(show.id, show.tmdbId, details.number_of_seasons); await delay(RATE_LIMIT_MS); } } } // Cache images for all library titles (posters, backdrops, stills, logos) async function cacheImagesJob() { if (!imageCacheEnabled()) return; const libraryIds = getLibraryTitleIds(); for (const titleId of libraryIds) { try { await cacheImagesForTitle(titleId); await cacheEpisodeStills(titleId); await cacheProviderLogos(titleId); } catch { // Continue with remaining titles } await delay(RATE_LIMIT_MS); } } async function scheduledBackupJob() { const enabled = getSetting("scheduledBackups"); if (enabled !== "true") { log.debug("Scheduled backups disabled, skipping"); return; } ensureBackupDir(); createBackup(); const maxStr = getSetting("maxBackupRetention"); const max = maxStr ? Number.parseInt(maxStr, 10) : 7; pruneBackups(max); } export function startJobs() { if (jobs.size > 0) return; schedule("scheduledBackup", "0 2 * * *", scheduledBackupJob); schedule("nightlyRefreshLibrary", "0 3 * * *", nightlyRefreshLibrary); schedule("refreshAvailability", "0 */6 * * *", refreshAvailabilityJob); schedule("refreshRecommendations", "0 */12 * * *", refreshRecommendationsJob); schedule("refreshTvChildren", "30 */12 * * *", refreshTvChildrenJob); schedule("cacheImages", "0 1,13 * * *", cacheImagesJob); log.info(`Started ${jobs.size} jobs`); } export function stopJobs() { for (const job of jobs.values()) { job.stop(); } }