Files
gunshiz bec24ad00b
CI / Verify (push) Successful in 2m34s
CI / Build immutable images and deploy (push) Successful in 4m1s
feat : dd up up
2026-10-07 18:15:27 +07:00

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");
}