195 lines
6.5 KiB
TypeScript
195 lines
6.5 KiB
TypeScript
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<string, unknown> = {}): void {
|
|
console.log(JSON.stringify({ component: "discord-catalog-worker", event, ...details }));
|
|
}
|
|
|
|
function wait(milliseconds: number, signal: AbortSignal): Promise<boolean> {
|
|
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<typeof setInterval> | undefined;
|
|
|
|
async function writeHeartbeat(): Promise<void> {
|
|
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<void> | null = null;
|
|
let activeFetch: AbortController | null = null;
|
|
|
|
async function notify(event: string, details: Record<string, unknown> = {}): Promise<void> {
|
|
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<void> {
|
|
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<void> {
|
|
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<void>((resolve) => {
|
|
stopWorker = resolve;
|
|
});
|
|
const stop = () => {
|
|
shutdown.abort();
|
|
activeFetch?.abort();
|
|
stopWorker?.();
|
|
};
|
|
process.once("SIGINT", stop);
|
|
process.once("SIGTERM", stop);
|
|
|
|
try {
|
|
const ready = new Promise<void>((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");
|
|
}
|