feat : let discord fetcher can stop working and accept every message type
This commit is contained in:
@@ -108,10 +108,13 @@ refetch authoritative state; SSE never contains page content and is not the
|
||||
source of truth.
|
||||
|
||||
The Discord catalog worker watches the private channel configured by
|
||||
`DISCORD_CHANNEL_ID`. It parses new Lunaris `Version Change` announcements and
|
||||
then fetches and synchronizes the latest version reported by Lunaris. Startup
|
||||
`DISCORD_CHANNEL_ID`. Every new message in that channel, including messages you send yourself,
|
||||
triggers a check and synchronization of the latest version reported by Lunaris. Startup
|
||||
only connects the listener; it does not replay channel history or start a sync.
|
||||
Worker and sync results are sent to `DISCORD_LOG_CHANNEL_ID`. Enable the Discord
|
||||
Worker and sync results are sent to `DISCORD_LOG_CHANNEL_ID`; start notifications
|
||||
include a clickable link to the message that triggered the check. Send
|
||||
`!lunaris stop` in the configured channel to cancel the active fetch and clear
|
||||
queued fetches; a later message starts a new check. Enable the Discord
|
||||
Message Content intent and grant the bot View Channel and Send Messages in the
|
||||
configured channels.
|
||||
|
||||
|
||||
+58
-74
@@ -10,10 +10,6 @@ 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;
|
||||
@@ -28,17 +24,6 @@ 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) => {
|
||||
@@ -66,8 +51,9 @@ const client = new Client({
|
||||
],
|
||||
});
|
||||
|
||||
let pendingVersion: string | null = null;
|
||||
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 {
|
||||
@@ -76,7 +62,7 @@ async function notify(event: string, details: Record<string, unknown> = {}): Pro
|
||||
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\`\`\``);
|
||||
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,
|
||||
@@ -85,69 +71,72 @@ async function notify(event: string, details: Record<string, unknown> = {}): Pro
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
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();
|
||||
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);
|
||||
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);
|
||||
}
|
||||
const details = { announcedVersion, 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 (pendingVersion && !shutdown.signal.aborted) {
|
||||
const version = pendingVersion;
|
||||
pendingVersion = null;
|
||||
await syncWithRetry(version);
|
||||
while (pendingMessageUrl && !shutdown.signal.aborted) {
|
||||
const messageUrl = pendingMessageUrl;
|
||||
pendingMessageUrl = null;
|
||||
await syncWithRetry(messageUrl);
|
||||
}
|
||||
}
|
||||
|
||||
function enqueue(version: string): void {
|
||||
if (!pendingVersion || compareLunarisVersions(version, pendingVersion) > 0) {
|
||||
pendingVersion = version;
|
||||
log("version_queued", { version, source: "message" });
|
||||
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;
|
||||
@@ -155,12 +144,6 @@ function enqueue(version: string): void {
|
||||
}
|
||||
}
|
||||
|
||||
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 }));
|
||||
@@ -171,6 +154,7 @@ const stopped = new Promise<void>((resolve) => {
|
||||
});
|
||||
const stop = () => {
|
||||
shutdown.abort();
|
||||
activeFetch?.abort();
|
||||
stopWorker?.();
|
||||
};
|
||||
process.once("SIGINT", stop);
|
||||
|
||||
+63
-37
@@ -113,27 +113,47 @@ interface SyncedCharacter {
|
||||
baseStats: CatalogBaseStat[];
|
||||
}
|
||||
|
||||
type SyncControl = { checkCancelled: () => Promise<void>; report: (phase: string, completed: number, total: number, message: string) => Promise<void> };
|
||||
type SyncControl = { signal?: AbortSignal; checkCancelled: () => Promise<void>; report: (phase: string, completed: number, total: number, message: string) => Promise<void> };
|
||||
|
||||
async function retryTimeout<T>(label: string, task: () => Promise<T>, signal?: AbortSignal): Promise<T> {
|
||||
for (let attempt = 1; ; attempt += 1) {
|
||||
signal?.throwIfAborted();
|
||||
try {
|
||||
return await task();
|
||||
} catch (cause) {
|
||||
signal?.throwIfAborted();
|
||||
const message = cause instanceof Error ? cause.message : String(cause);
|
||||
if (attempt >= 3 || !/timed out|timeout/i.test(message)) {
|
||||
throw new Error(`${label}: ${message}`, { cause });
|
||||
}
|
||||
await new Promise((resolve) => setTimeout(resolve, attempt * 1_000));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function getJson(url: string, signal?: AbortSignal): Promise<unknown> {
|
||||
const response = await fetch(url, {
|
||||
signal: signal ?? AbortSignal.timeout(30_000),
|
||||
cache: "no-store",
|
||||
});
|
||||
if (!response.ok)
|
||||
throw new Error(`Lunaris returned ${response.status} for ${url}`);
|
||||
return response.json();
|
||||
return retryTimeout(`Lunaris request ${url}`, async () => {
|
||||
const response = await fetch(url, {
|
||||
signal: signal ? AbortSignal.any([signal, AbortSignal.timeout(60_000)]) : AbortSignal.timeout(60_000),
|
||||
cache: "no-store",
|
||||
});
|
||||
if (!response.ok)
|
||||
throw new Error(`Lunaris returned ${response.status} for ${url}`);
|
||||
return response.json();
|
||||
}, signal);
|
||||
}
|
||||
|
||||
async function getOptionalJson(url: string, signal?: AbortSignal): Promise<unknown | null> {
|
||||
const response = await fetch(url, {
|
||||
signal: signal ?? AbortSignal.timeout(30_000),
|
||||
cache: "no-store",
|
||||
});
|
||||
if (response.status === 404) return null;
|
||||
if (!response.ok)
|
||||
throw new Error(`Lunaris returned ${response.status} for ${url}`);
|
||||
return response.json();
|
||||
return retryTimeout(`Lunaris request ${url}`, async () => {
|
||||
const response = await fetch(url, {
|
||||
signal: signal ? AbortSignal.any([signal, AbortSignal.timeout(60_000)]) : AbortSignal.timeout(60_000),
|
||||
cache: "no-store",
|
||||
});
|
||||
if (response.status === 404) return null;
|
||||
if (!response.ok)
|
||||
throw new Error(`Lunaris returned ${response.status} for ${url}`);
|
||||
return response.json();
|
||||
}, signal);
|
||||
}
|
||||
|
||||
function webpName(value: string): string {
|
||||
@@ -222,27 +242,29 @@ async function uploadAssets(
|
||||
if (index >= assets.length) return;
|
||||
await control?.checkCancelled();
|
||||
const asset = assets[index];
|
||||
if (!(await storage.file(asset.target).exists())) {
|
||||
if (asset.body) {
|
||||
await storage.write(asset.target, asset.body, {
|
||||
type: asset.type ?? "image/webp",
|
||||
acl: "public-read",
|
||||
});
|
||||
} else {
|
||||
const response = await fetch(asset.source!, {
|
||||
signal: AbortSignal.timeout(30_000),
|
||||
cache: "no-store",
|
||||
});
|
||||
if (response.status === 404) missing.add(asset.target);
|
||||
else if (!response.ok)
|
||||
throw new Error(`Asset returned ${response.status}: ${asset.source}`);
|
||||
else
|
||||
await storage.write(asset.target, response, {
|
||||
await retryTimeout(`Asset ${asset.target}`, async () => {
|
||||
if (!(await storage.file(asset.target).exists())) {
|
||||
if (asset.body) {
|
||||
await storage.write(asset.target, asset.body, {
|
||||
type: asset.type ?? "image/webp",
|
||||
acl: "public-read",
|
||||
});
|
||||
} else {
|
||||
const response = await fetch(asset.source!, {
|
||||
signal: control?.signal ? AbortSignal.any([control.signal, AbortSignal.timeout(60_000)]) : AbortSignal.timeout(60_000),
|
||||
cache: "no-store",
|
||||
});
|
||||
if (response.status === 404) missing.add(asset.target);
|
||||
else if (!response.ok)
|
||||
throw new Error(`Asset returned ${response.status}: ${asset.source}`);
|
||||
else
|
||||
await storage.write(asset.target, response, {
|
||||
type: asset.type ?? "image/webp",
|
||||
acl: "public-read",
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}, control?.signal);
|
||||
completed += 1;
|
||||
await reportProgress();
|
||||
}
|
||||
@@ -256,10 +278,10 @@ async function uploadAssets(
|
||||
async function syncCatalogVersion(version: string, control?: SyncControl): Promise<void> {
|
||||
await control?.report("lists", 0, 4, "กำลังดึงรายการ Catalog จาก Lunaris");
|
||||
const [charactersRaw, weaponsRaw, artifactsRaw, materialsRaw] = await Promise.all([
|
||||
getJson(`${API_ROOT}/${version}/charlist.json`),
|
||||
getJson(`${API_ROOT}/${version}/weaponlist.json`),
|
||||
getJson(`${API_ROOT}/${version}/artifactlist.json`),
|
||||
getJson(`${API_ROOT}/${version}/materiallist.json`),
|
||||
getJson(`${API_ROOT}/${version}/charlist.json`, control?.signal),
|
||||
getJson(`${API_ROOT}/${version}/weaponlist.json`, control?.signal),
|
||||
getJson(`${API_ROOT}/${version}/artifactlist.json`, control?.signal),
|
||||
getJson(`${API_ROOT}/${version}/materiallist.json`, control?.signal),
|
||||
]);
|
||||
await control?.report("lists", 4, 4, "ได้รับรายการ Catalog แล้ว");
|
||||
const characterData = characterSchema.parse(charactersRaw);
|
||||
@@ -288,6 +310,7 @@ async function syncCatalogVersion(version: string, control?: SyncControl): Promi
|
||||
async ([key, value]) => {
|
||||
const detailRaw = await getOptionalJson(
|
||||
`${API_ROOT}/${version}/en/char/${encodeURIComponent(key)}.json`,
|
||||
control?.signal,
|
||||
);
|
||||
const detail = detailRaw
|
||||
? characterDetailSchema.parse(detailRaw)
|
||||
@@ -337,6 +360,7 @@ async function syncCatalogVersion(version: string, control?: SyncControl): Promi
|
||||
const detail = characterDetailSchema.parse(
|
||||
await getJson(
|
||||
`${API_ROOT}/${version}/en/char/${traveler.detailKey}.json`,
|
||||
control?.signal,
|
||||
),
|
||||
);
|
||||
return {
|
||||
@@ -496,6 +520,7 @@ async function syncCatalogVersion(version: string, control?: SyncControl): Promi
|
||||
weapons = weapons.filter((item) => !missingAssets.has(item.imageKey));
|
||||
artifacts = artifacts.filter((item) => !missingAssets.has(item.imageKey));
|
||||
|
||||
await control?.checkCancelled();
|
||||
const storage = await getMediaStorage();
|
||||
await Promise.all([
|
||||
storage.write(
|
||||
@@ -552,6 +577,7 @@ async function syncCatalogVersion(version: string, control?: SyncControl): Promi
|
||||
),
|
||||
]);
|
||||
|
||||
await control?.checkCancelled();
|
||||
await control?.report("database", 0, 1, "กำลังบันทึก Catalog");
|
||||
await getDb().transaction(async (tx) => {
|
||||
await tx.delete(catalogCharacters);
|
||||
|
||||
@@ -197,10 +197,14 @@ describe("environment template contract", () => {
|
||||
expect(rbac).toContain("buzz-sheet-discord-worker");
|
||||
});
|
||||
|
||||
it("syncs only for new Discord version announcements", async () => {
|
||||
it("checks Lunaris for every new message in the configured Discord channel", async () => {
|
||||
const worker = await repositoryFile("backend/discord.ts");
|
||||
|
||||
expect(worker).toContain("client.on(Events.MessageCreate, handleMessage)");
|
||||
expect(worker).toContain("if (message.channelId !== channelId) return;");
|
||||
expect(worker).not.toContain("message.author.id === client.user?.id");
|
||||
expect(worker).not.toContain("parseLunarisVersionChange");
|
||||
expect(worker).toContain("const details = { version, attempt: attempt + 1, messageUrl };");
|
||||
expect(worker).not.toContain("channel.messages.fetch");
|
||||
expect(worker).not.toContain('source: "history"');
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user