diff --git a/README.md b/README.md index a504708..6dbb708 100644 --- a/README.md +++ b/README.md @@ -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. diff --git a/backend/discord.ts b/backend/discord.ts index 5b160cc..fb79b75 100644 --- a/backend/discord.ts +++ b/backend/discord.ts @@ -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 = {}): 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 { 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 | null = null; +let activeFetch: AbortController | null = null; async function notify(event: string, details: Record = {}): Promise { try { @@ -76,7 +62,7 @@ async function notify(event: string, details: Record = {}): 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 = {}): Pro } } -async function syncWithRetry(announcedVersion: string): Promise { - 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 { + 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 { - 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((resolve) => { }); const stop = () => { shutdown.abort(); + activeFetch?.abort(); stopWorker?.(); }; process.once("SIGINT", stop); diff --git a/lib/catalog/sync.ts b/lib/catalog/sync.ts index 2762970..de4cfc7 100644 --- a/lib/catalog/sync.ts +++ b/lib/catalog/sync.ts @@ -113,27 +113,47 @@ interface SyncedCharacter { baseStats: CatalogBaseStat[]; } -type SyncControl = { checkCancelled: () => Promise; report: (phase: string, completed: number, total: number, message: string) => Promise }; +type SyncControl = { signal?: AbortSignal; checkCancelled: () => Promise; report: (phase: string, completed: number, total: number, message: string) => Promise }; + +async function retryTimeout(label: string, task: () => Promise, signal?: AbortSignal): Promise { + 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 { - 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 { - 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 { 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); diff --git a/tests/deployment-contract.test.ts b/tests/deployment-contract.test.ts index 8a57818..34b0332 100644 --- a/tests/deployment-contract.test.ts +++ b/tests/deployment-contract.test.ts @@ -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"'); });