import { Client, Events, GatewayIntentBits, type Message, } from "discord.js"; import { closeDb } from "@/db/client"; import { getLatestLunarisVersion, processCatalogSync, } from "@/lib/catalog/sync"; import { closeRedisClient, discordWorkerHeartbeatKey, getRedisClient } from "@/lib/redis/client"; const RETRY_DELAYS_MS = [0, 30_000, 120_000, 600_000, 1_800_000] as const; function required(name: string): string { const value = process.env[name]?.trim(); if (!value) throw new Error(`${name} is required for the Discord catalog worker.`); return value; } function log(event: string, details: Record = {}): void { console.log(JSON.stringify({ component: "discord-catalog-worker", event, ...details })); } function wait(milliseconds: number, signal: AbortSignal): Promise { if (signal.aborted) return Promise.resolve(false); return new Promise((resolve) => { const timer = setTimeout(() => { signal.removeEventListener("abort", aborted); resolve(true); }, milliseconds); const aborted = () => { clearTimeout(timer); resolve(false); }; signal.addEventListener("abort", aborted, { once: true }); }); } const token = required("DISCORD_BOT_TOKEN"); const channelId = required("DISCORD_CHANNEL_ID"); const logChannelId = required("DISCORD_LOG_CHANNEL_ID"); const shutdown = new AbortController(); const client = new Client({ intents: [ GatewayIntentBits.Guilds, GatewayIntentBits.GuildMessages, GatewayIntentBits.MessageContent, ], }); let heartbeatTimer: ReturnType | undefined; async function writeHeartbeat(): Promise { if (!client.isReady()) return; try { const redis = await getRedisClient(); await redis.set(discordWorkerHeartbeatKey(), "1", "EX", 45); } catch (cause) { log("heartbeat_failed", { error: cause instanceof Error ? cause.message : String(cause) }); } } let pendingMessageUrl: string | null = null; let runner: Promise | null = null; let activeFetch: AbortController | null = null; async function notify(event: string, details: Record = {}): Promise { try { const channel = await client.channels.fetch(logChannelId); if (!channel?.isSendable()) { log("notification_failed", { event, error: "Discord log channel is not sendable." }); return; } await channel.send(`**Lunaris catalog worker - ${event}**\n\`\`\`json\n${JSON.stringify(details, null, 2)}\n\`\`\`${typeof details.messageUrl === "string" ? `\nTrigger message: ${details.messageUrl}` : ""}`); } catch (cause) { log("notification_failed", { event, error: cause instanceof Error ? cause.message : String(cause), }); } } async function syncWithRetry(messageUrl: string): Promise { const cancellation = new AbortController(); activeFetch = cancellation; try { for (let attempt = 0; attempt < RETRY_DELAYS_MS.length; attempt += 1) { if (shutdown.signal.aborted || cancellation.signal.aborted) return; if (!(await wait(RETRY_DELAYS_MS[attempt], AbortSignal.any([shutdown.signal, cancellation.signal])))) return; try { const version = await getLatestLunarisVersion(cancellation.signal); if (cancellation.signal.aborted) return; const details = { version, attempt: attempt + 1, messageUrl }; log("sync_started", details); await notify("sync started", details); const result = await processCatalogSync(version, { checkCancelled: async () => { cancellation.signal.throwIfAborted(); }, report: async () => { cancellation.signal.throwIfAborted(); }, signal: cancellation.signal, }); if (cancellation.signal.aborted) return; log(result === "synced" ? "sync_completed" : "sync_ignored", { version, reason: result === "unchanged" ? "current-or-older" : undefined, }); await notify(result === "synced" ? "sync completed" : "already current", { version, }); return; } catch (cause) { if (cancellation.signal.aborted) return; const message = cause instanceof Error ? cause.message : String(cause); if (attempt === RETRY_DELAYS_MS.length - 1) { const details = { attempts: attempt + 1, error: message }; log("sync_failed", details); await notify("sync failed", details); return; } const details = { attempt: attempt + 1, error: message }; log("sync_retrying", details); await notify("sync retrying", details); } } } finally { if (activeFetch === cancellation) activeFetch = null; } } async function drain(): Promise { while (pendingMessageUrl && !shutdown.signal.aborted) { const messageUrl = pendingMessageUrl; pendingMessageUrl = null; await syncWithRetry(messageUrl); } } function handleMessage(message: Message): void { if (message.channelId !== channelId) return; if (message.content.trim().toLowerCase() === "!lunaris stop") { pendingMessageUrl = null; activeFetch?.abort(); log("sync_stopped", { messageId: message.id }); void notify("sync stopped", { messageUrl: message.url }); return; } pendingMessageUrl = message.url; log("sync_queued", { messageId: message.id, messageUrl: message.url }); if (!runner) { runner = drain().finally(() => { runner = null; }); } } client.on(Events.MessageCreate, handleMessage); client.on(Events.Error, (cause) => log("discord_error", { error: cause.message })); client.on(Events.Warn, (message) => log("discord_warning", { message })); let stopWorker: (() => void) | undefined; const stopped = new Promise((resolve) => { stopWorker = resolve; }); const stop = () => { shutdown.abort(); activeFetch?.abort(); stopWorker?.(); }; process.once("SIGINT", stop); process.once("SIGTERM", stop); try { const ready = new Promise((resolve) => { client.once(Events.ClientReady, () => resolve()); }); await client.login(token); await Promise.race([ready, stopped]); if (!shutdown.signal.aborted) { log("discord_ready", { userId: client.user?.id, channelId }); await writeHeartbeat(); heartbeatTimer = setInterval(() => { void writeHeartbeat(); }, 15_000); await notify("worker ready", { userId: client.user?.id, channelId }); await stopped; } } finally { shutdown.abort(); if (heartbeatTimer) clearInterval(heartbeatTimer); client.destroy(); await Promise.resolve(runner).catch(() => undefined); await Promise.all([closeDb(), closeRedisClient()]); log("stopped"); }