197 lines
6.0 KiB
TypeScript
197 lines
6.0 KiB
TypeScript
import {
|
|
Client,
|
|
Events,
|
|
GatewayIntentBits,
|
|
type Message,
|
|
} from "discord.js";
|
|
|
|
import { closeDb } from "@/db/client";
|
|
import {
|
|
getLatestLunarisVersion,
|
|
processCatalogSync,
|
|
} from "@/lib/catalog/sync";
|
|
import {
|
|
compareLunarisVersions,
|
|
parseLunarisVersionChange,
|
|
} from "@/lib/catalog/version";
|
|
import { closeRedisClient } 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 messageText(message: Message): string {
|
|
return [
|
|
message.content,
|
|
...message.embeds.flatMap((embed) => [
|
|
embed.title ?? "",
|
|
embed.description ?? "",
|
|
...embed.fields.flatMap((field) => [field.name, field.value]),
|
|
]),
|
|
].join("\n");
|
|
}
|
|
|
|
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 pendingVersion: string | null = null;
|
|
let runner: Promise<void> | 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\`\`\``);
|
|
} catch (cause) {
|
|
log("notification_failed", {
|
|
event,
|
|
error: cause instanceof Error ? cause.message : String(cause),
|
|
});
|
|
}
|
|
}
|
|
|
|
async function syncWithRetry(announcedVersion: string): Promise<void> {
|
|
for (let attempt = 0; attempt < RETRY_DELAYS_MS.length; attempt += 1) {
|
|
if (shutdown.signal.aborted) return;
|
|
if (
|
|
pendingVersion &&
|
|
compareLunarisVersions(pendingVersion, announcedVersion) > 0
|
|
) {
|
|
log("sync_superseded", { announcedVersion, pendingVersion });
|
|
return;
|
|
}
|
|
if (!(await wait(RETRY_DELAYS_MS[attempt], shutdown.signal))) return;
|
|
if (
|
|
pendingVersion &&
|
|
compareLunarisVersions(pendingVersion, announcedVersion) > 0
|
|
) {
|
|
log("sync_superseded", { announcedVersion, pendingVersion });
|
|
return;
|
|
}
|
|
|
|
try {
|
|
const version = await getLatestLunarisVersion();
|
|
const details = { announcedVersion, version, attempt: attempt + 1 };
|
|
log("sync_started", details);
|
|
await notify("sync started", details);
|
|
const result = await processCatalogSync(version);
|
|
log(result === "synced" ? "sync_completed" : "sync_ignored", {
|
|
announcedVersion,
|
|
version,
|
|
reason: result === "unchanged" ? "current-or-older" : undefined,
|
|
});
|
|
await notify(result === "synced" ? "sync completed" : "already current", {
|
|
announcedVersion,
|
|
version,
|
|
});
|
|
return;
|
|
} catch (cause) {
|
|
const message = cause instanceof Error ? cause.message : String(cause);
|
|
if (attempt === RETRY_DELAYS_MS.length - 1) {
|
|
const details = { announcedVersion, attempts: attempt + 1, error: message };
|
|
log("sync_failed", details);
|
|
await notify("sync failed", details);
|
|
return;
|
|
}
|
|
const details = { announcedVersion, attempt: attempt + 1, error: message };
|
|
log("sync_retrying", details);
|
|
await notify("sync retrying", details);
|
|
}
|
|
}
|
|
}
|
|
|
|
async function drain(): Promise<void> {
|
|
while (pendingVersion && !shutdown.signal.aborted) {
|
|
const version = pendingVersion;
|
|
pendingVersion = null;
|
|
await syncWithRetry(version);
|
|
}
|
|
}
|
|
|
|
function enqueue(version: string): void {
|
|
if (!pendingVersion || compareLunarisVersions(version, pendingVersion) > 0) {
|
|
pendingVersion = version;
|
|
log("version_queued", { version, source: "message" });
|
|
}
|
|
if (!runner) {
|
|
runner = drain().finally(() => {
|
|
runner = null;
|
|
});
|
|
}
|
|
}
|
|
|
|
function handleMessage(message: Message): void {
|
|
if (message.channelId !== channelId || message.author.id === client.user?.id) return;
|
|
const change = parseLunarisVersionChange(messageText(message));
|
|
if (change) enqueue(change.targetVersion);
|
|
}
|
|
|
|
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();
|
|
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 notify("worker ready", { userId: client.user?.id, channelId });
|
|
await stopped;
|
|
}
|
|
} finally {
|
|
shutdown.abort();
|
|
client.destroy();
|
|
await Promise.resolve(runner).catch(() => undefined);
|
|
await Promise.all([closeDb(), closeRedisClient()]);
|
|
log("stopped");
|
|
}
|