feat(catalog) : sync Lunaris updates from Discord
This commit is contained in:
@@ -0,0 +1,176 @@
|
||||
import {
|
||||
ChannelType,
|
||||
Client,
|
||||
Events,
|
||||
GatewayIntentBits,
|
||||
type Message,
|
||||
} from "discord.js";
|
||||
|
||||
import { closeDb } from "@/db/client";
|
||||
import { 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 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 syncWithRetry(version: string): Promise<void> {
|
||||
for (let attempt = 0; attempt < RETRY_DELAYS_MS.length; attempt += 1) {
|
||||
if (shutdown.signal.aborted) return;
|
||||
if (
|
||||
pendingVersion &&
|
||||
compareLunarisVersions(pendingVersion, version) > 0
|
||||
) {
|
||||
log("sync_superseded", { version, pendingVersion });
|
||||
return;
|
||||
}
|
||||
if (!(await wait(RETRY_DELAYS_MS[attempt], shutdown.signal))) return;
|
||||
if (
|
||||
pendingVersion &&
|
||||
compareLunarisVersions(pendingVersion, version) > 0
|
||||
) {
|
||||
log("sync_superseded", { version, pendingVersion });
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
log("sync_started", { version, attempt: attempt + 1 });
|
||||
const result = await processCatalogSync(version);
|
||||
log(result === "synced" ? "sync_completed" : "sync_ignored", {
|
||||
version,
|
||||
reason: result === "unchanged" ? "current-or-older" : undefined,
|
||||
});
|
||||
return;
|
||||
} catch (cause) {
|
||||
const message = cause instanceof Error ? cause.message : String(cause);
|
||||
if (attempt === RETRY_DELAYS_MS.length - 1) {
|
||||
log("sync_failed", { version, attempts: attempt + 1, error: message });
|
||||
return;
|
||||
}
|
||||
log("sync_retrying", { version, attempt: attempt + 1, error: message });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function drain(): Promise<void> {
|
||||
while (pendingVersion && !shutdown.signal.aborted) {
|
||||
const version = pendingVersion;
|
||||
pendingVersion = null;
|
||||
await syncWithRetry(version);
|
||||
}
|
||||
}
|
||||
|
||||
function enqueue(version: string, source: "history" | "message"): void {
|
||||
if (!pendingVersion || compareLunarisVersions(version, pendingVersion) > 0) {
|
||||
pendingVersion = version;
|
||||
log("version_queued", { version, source });
|
||||
}
|
||||
if (!runner) {
|
||||
runner = drain().finally(() => {
|
||||
runner = null;
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
function handleMessage(message: Message, source: "history" | "message"): void {
|
||||
if (message.channelId !== channelId || message.author.id === client.user?.id) return;
|
||||
const change = parseLunarisVersionChange(messageText(message));
|
||||
if (change) enqueue(change.targetVersion, source);
|
||||
}
|
||||
|
||||
client.on(Events.MessageCreate, (message) => handleMessage(message, "message"));
|
||||
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 });
|
||||
|
||||
const channel = await client.channels.fetch(channelId);
|
||||
if (
|
||||
!channel ||
|
||||
(channel.type !== ChannelType.GuildText &&
|
||||
channel.type !== ChannelType.GuildAnnouncement)
|
||||
) {
|
||||
throw new Error("DISCORD_CHANNEL_ID must identify a guild text channel.");
|
||||
}
|
||||
const messages = await channel.messages.fetch({ limit: 25 });
|
||||
for (const message of messages.values()) handleMessage(message, "history");
|
||||
|
||||
await stopped;
|
||||
}
|
||||
} finally {
|
||||
shutdown.abort();
|
||||
client.destroy();
|
||||
await Promise.resolve(runner).catch(() => undefined);
|
||||
await Promise.all([closeDb(), closeRedisClient()]);
|
||||
log("stopped");
|
||||
}
|
||||
Reference in New Issue
Block a user