From 76f0420d64d308c34b97d97e55c939c36664158f Mon Sep 17 00:00:00 2001 From: simoleo89 Date: Fri, 11 Sep 2026 00:36:22 +0200 Subject: [PATCH] feat(import): run catalog synchronizations as durable jobs --- scripts/jobs-worker.ts | 5 + .../import/official/official-sync-client.tsx | 224 +----------------- src/app/admin/import/sync/sync-all-client.tsx | 212 +---------------- .../api/admin/import/clone/sync-all/route.ts | 202 +--------------- .../admin/import/official/sync-all/route.ts | 95 +------- src/app/api/admin/studio/import-jobs/route.ts | 2 + .../admin/studio/furniture-sync-queue.tsx | 73 ++++++ src/lib/furni/import-job.ts | 5 + src/lib/services/furni-job-worker.test.ts | 28 ++- src/lib/services/furni-job-worker.ts | 176 ++++++++------ src/lib/services/furni-sync-item.test.ts | 75 ++++++ src/lib/services/furni-sync-item.ts | 36 +++ src/lib/services/furni-sync-queue.test.ts | 113 +++++++++ src/lib/services/furni-sync-queue.ts | 114 +++++++++ 14 files changed, 559 insertions(+), 801 deletions(-) create mode 100644 src/components/admin/studio/furniture-sync-queue.tsx create mode 100644 src/lib/services/furni-sync-item.test.ts create mode 100644 src/lib/services/furni-sync-item.ts create mode 100644 src/lib/services/furni-sync-queue.test.ts create mode 100644 src/lib/services/furni-sync-queue.ts diff --git a/scripts/jobs-worker.ts b/scripts/jobs-worker.ts index 1a5c4944..d92d78ec 100644 --- a/scripts/jobs-worker.ts +++ b/scripts/jobs-worker.ts @@ -1,3 +1,4 @@ +import { drainFurnitureImports } from "../src/lib/services/furni-job-worker"; import "./load-env"; import { Cron } from "croner"; import { and, eq, lt, lte, sql } from "drizzle-orm"; @@ -268,6 +269,10 @@ async function reportWorkerHeartbeat(): Promise { } async function main() { + new Cron("* * * * *", () => { + void drainFurnitureImports(); + }); + void drainFurnitureImports(); new Cron("* * * * *", () => { runCatalogExport().catch((e) => captureWorkerError(e, "Catalog export failed"), diff --git a/src/app/admin/import/official/official-sync-client.tsx b/src/app/admin/import/official/official-sync-client.tsx index 630377f1..25541ac3 100644 --- a/src/app/admin/import/official/official-sync-client.tsx +++ b/src/app/admin/import/official/official-sync-client.tsx @@ -1,225 +1,5 @@ "use client"; - -import { AlertCircle, Check, CloudDownload, Loader2 } from "lucide-react"; -import { useCallback, useRef, useState } from "react"; -import { toast } from "sonner"; -import { Button } from "@/components/ui/button"; -import { adminFetch } from "@/lib/admin-fetch"; - -interface SyncEvent { - type: string; - classname?: string; - status?: string; - succeeded?: number; - failed?: number; - message?: string; -} - -interface SyncResult { - succeeded: number; - failed: number; - message: string; -} - -/** - * Official furnidata sync client. - * Imports all furniture from official Habbo furnidata that the hotel doesn't have yet. - */ +import { FurnitureSyncQueue } from "@/components/admin/studio/furniture-sync-queue"; export function OfficialSyncClient() { - const [syncing, setSyncing] = useState(false); - const [progress, setProgress] = useState<{ - done: number; - total: number; - } | null>(null); - const [log, setLog] = useState([]); - const [result, setResult] = useState(null); - const abortRef = useRef(null); - - const startSync = useCallback(async () => { - setSyncing(true); - setProgress(null); - setLog([]); - setResult(null); - - const abort = new AbortController(); - abortRef.current = abort; - - try { - const res = await adminFetch("/api/admin/import/official/sync-all", { - method: "POST", - headers: { "Content-Type": "application/json" }, - body: JSON.stringify({}), - signal: abort.signal, - }); - - if (!res.body) { - toast.error("No response stream"); - setSyncing(false); - return; - } - - if (!res.headers.get("content-type")?.includes("text/event-stream")) { - let errMsg = "Sync failed"; - try { - const json = await res.json(); - errMsg = json.message || errMsg; - } catch { - // ignore - } - toast.error(errMsg); - setSyncing(false); - return; - } - - const reader = res.body.getReader(); - const decoder = new TextDecoder(); - let buf = ""; - let succeeded = 0; - let failed = 0; - let totalItems = 0; - - while (true) { - const { value, done: streamDone } = await reader.read(); - if (streamDone) break; - buf += decoder.decode(value, { stream: true }); - const parts = buf.split("\n\n"); - buf = parts.pop() ?? ""; - for (const part of parts) { - if (!part.startsWith("data: ")) continue; - try { - const evt = JSON.parse(part.slice(6)) as SyncEvent; - if (evt.type === "item_progress") { - const classname = evt.classname ?? ""; - if (evt.status === "done") { - succeeded++; - totalItems++; - setProgress({ done: succeeded + failed, total: totalItems }); - setLog((prev) => [...prev, `✓ ${classname}`]); - } else if (evt.status === "failed") { - failed++; - totalItems++; - setProgress({ done: succeeded + failed, total: totalItems }); - setLog((prev) => [...prev, `✗ ${classname} (failed)`]); - } - } - if (evt.type === "batch_complete") { - succeeded = evt.succeeded ?? succeeded; - failed = evt.failed ?? failed; - } - if (evt.type === "error") { - toast.error( - typeof evt.message === "string" - ? evt.message - : "Sync finished with errors", - ); - } - } catch { - // skip parse errors - } - } - } - - setResult({ - succeeded, - failed, - message: - totalItems > 0 - ? `${succeeded} items synced` - : "All items were already up to date", - }); - if (succeeded > 0) toast.success(`${succeeded} items synced`); - if (failed > 0) toast.error(`${failed} items failed`); - } catch (err) { - if ((err as Error).name === "AbortError") { - toast.info("Sync cancelled"); - } else { - toast.error("Sync request failed"); - } - } finally { - setSyncing(false); - abortRef.current = null; - } - }, []); - - return ( -
- {!result && ( -
- - {syncing && ( - - )} -
- )} - - {progress && ( -
- Processed {progress.done} items - {progress.total > 0 ? ` out of ${progress.total}` : ""} -
- )} - - {log.length > 0 && ( -
- {log.map((line) => ( -
- {line} -
- ))} -
- )} - - {result && ( -
-
- {result.failed === 0 ? ( - - ) : ( - - )} - Sync Complete -
-
-

- Succeeded:{" "} - {result.succeeded} -

-

- Failed:{" "} - {result.failed} -

-

{result.message}

-
-
- )} - - {!syncing && !result && ( -
- -

- Click the button above to scan the official Habbo furnidata and - automatically import any furniture that your hotel doesn't have yet. -

-
- )} -
- ); + return ; } diff --git a/src/app/admin/import/sync/sync-all-client.tsx b/src/app/admin/import/sync/sync-all-client.tsx index cf551c0e..9428a446 100644 --- a/src/app/admin/import/sync/sync-all-client.tsx +++ b/src/app/admin/import/sync/sync-all-client.tsx @@ -1,213 +1,5 @@ "use client"; - -import { Check, CloudDownload, Loader2 } from "lucide-react"; -import { useCallback, useRef, useState } from "react"; -import { toast } from "sonner"; -import { Button } from "@/components/ui/button"; -import { adminFetch } from "@/lib/admin-fetch"; - -interface SyncEvent { - type: string; - classname?: string; - status?: string; - succeeded?: number; - failed?: number; -} - +import { FurnitureSyncQueue } from "@/components/admin/studio/furniture-sync-queue"; export function SyncAllClient() { - const [syncing, setSyncing] = useState(false); - const [progress, setProgress] = useState<{ - done: number; - total: number; - } | null>(null); - const [log, setLog] = useState([]); - const [result, setResult] = useState<{ - succeeded: number; - failed: number; - sources: number; - } | null>(null); - const abortRef = useRef(null); - - const startSync = useCallback(async () => { - setSyncing(true); - setProgress(null); - setLog([]); - setResult(null); - - const abort = new AbortController(); - abortRef.current = abort; - - try { - const res = await adminFetch("/api/admin/import/clone/sync-all", { - method: "POST", - headers: { "Content-Type": "application/json" }, - body: JSON.stringify({}), - signal: abort.signal, - }); - - if (!res.body) { - toast.error("No response stream"); - setSyncing(false); - return; - } - - if (!res.headers.get("content-type")?.includes("text/event-stream")) { - let errMsg = "Sync failed"; - try { - const json = await res.json(); - errMsg = json.error || errMsg; - } catch { - // ignore - } - toast.error(errMsg); - setSyncing(false); - return; - } - - const reader = res.body.getReader(); - const decoder = new TextDecoder(); - let buf = ""; - let succeeded = 0; - let failed = 0; - let totalItems = 0; - - while (true) { - const { value, done: streamDone } = await reader.read(); - if (streamDone) break; - buf += decoder.decode(value, { stream: true }); - const parts = buf.split("\n\n"); - buf = parts.pop() ?? ""; - for (const part of parts) { - if (!part.startsWith("data: ")) continue; - try { - const evt = JSON.parse(part.slice(6)) as SyncEvent; - if ( - evt.type === "item_progress" && - (evt.status === "done" || evt.status === "failed") - ) { - totalItems++; - if (evt.status === "done") succeeded++; - else failed++; - setProgress({ done: succeeded + failed, total: totalItems }); - setLog((prev) => [ - ...prev, - `${evt.status === "done" ? "✓" : "✗"} ${evt.classname ?? ""}`, - ]); - } - if (evt.type === "batch_complete") { - succeeded = evt.succeeded ?? succeeded; - failed = evt.failed ?? failed; - } - } catch { - // skip parse errors - } - } - } - - // Check for skipped state (all up to date) - try { - const parsed = JSON.parse( - buf.includes("data:") ? (buf.split("data:").pop() ?? "{}") : "{}", - ); - if (parsed.skipped) { - toast.info(parsed.message || "All sources are up to date"); - setResult({ succeeded: 0, failed: 0, sources: 0 }); - setSyncing(false); - return; - } - } catch { - // not a skipped response - } - - setResult({ succeeded, failed, sources: 1 }); - if (succeeded > 0) toast.success(`${succeeded} items synced`); - if (failed > 0) toast.error(`${failed} items failed`); - } catch (err) { - if ((err as Error).name === "AbortError") { - toast.info("Sync cancelled"); - } else { - toast.error("Sync request failed"); - } - } finally { - setSyncing(false); - abortRef.current = null; - } - }, []); - - return ( -
- {!result && ( -
- - {syncing && ( - - )} -
- )} - - {progress && ( -
- Processed {progress.done} items - {progress.total > 0 ? ` out of ${progress.total}` : ""} -
- )} - - {log.length > 0 && ( -
- {log.map((line) => ( -
- {line} -
- ))} -
- )} - - {result && result.sources > 0 && ( -
-
- - Sync Complete -
-
-

- Succeeded:{" "} - {result.succeeded} -

-

- Failed:{" "} - {result.failed} -

-
-
- )} - - {!syncing && !result && ( -
- -

- Click the button above to scan all clone sources for missing - furniture and import it automatically. -

-
- )} -
- ); + return ; } diff --git a/src/app/api/admin/import/clone/sync-all/route.ts b/src/app/api/admin/import/clone/sync-all/route.ts index 3c80233b..f55586b1 100644 --- a/src/app/api/admin/import/clone/sync-all/route.ts +++ b/src/app/api/admin/import/clone/sync-all/route.ts @@ -1,205 +1,7 @@ -import { inArray } from "drizzle-orm"; import { withAdmin } from "@/lib/api-handler"; -import { db, ItemsBase } from "@/lib/db"; import { PERMS } from "@/lib/permissions"; -import { logServerError } from "@/lib/server-log"; -import { logAudit } from "@/lib/services/audit"; -import { - cloneSingleFurni, - fetchSourceFurnidata, - type SourceFurni, -} from "@/lib/services/clone-import"; -import { listSources } from "@/lib/services/clone-sources"; -import { - appendFurniEntriesBatch, - appendFurniEntry, -} from "@/lib/services/furni-data"; -import { runSseBatch } from "@/lib/services/import/core/sse-batch"; - -function sseEvent(data: unknown): string { - return `data: ${JSON.stringify(data)}\n\n`; -} - +import { enqueueFurnitureSync } from "@/lib/services/furni-sync-queue"; export const POST = withAdmin( { permission: PERMS.ASSETS_IMPORT }, - async (request, ctx) => { - const body = await request.json().catch(() => ({})); - const sourceIdFilter: string | undefined = body.sourceId || undefined; - - const allSources = await listSources(); - const sources = sourceIdFilter - ? allSources.filter((s) => s.id === sourceIdFilter) - : allSources; - - if (sources.length === 0) { - return new Response( - sseEvent({ type: "error", message: "No sources configured" }), - { - status: 400, - headers: { "Content-Type": "text/event-stream" }, - }, - ); - } - - // Pre-fetch furnidata ONCE per source and build a lookup map. - // The old code re-fetched inside the worker loop for every item, - // causing 2000+ redundant DB reads + HTTP requests. - const sourceFurniDataMap = new Map>(); - - for (const source of sources) { - try { - const entries = await fetchSourceFurnidata(source.furnidataUrl); - const byClassname = new Map(entries.map((e) => [e.classname, e])); - sourceFurniDataMap.set(source.id, byClassname); - } catch (error) { - logServerError("clone-import.furnidata_prefetch_failed", error, { - source: source.id, - }); - } - } - - // Collect all missing items using pre-fetched data - const allItems: Array<{ - classname: string; - sourceName: string; - sourceId: string; - }> = []; - - for (const source of sources) { - const byClassname = sourceFurniDataMap.get(source.id); - if (!byClassname) continue; - - const classnames = [...byClassname.keys()]; - if (classnames.length === 0) continue; - - const rows = await db - .select({ item_name: ItemsBase.itemName }) - .from(ItemsBase) - .where(inArray(ItemsBase.itemName, classnames)); - const have = new Set(rows.map((r) => r.item_name)); - - for (const classname of classnames) { - if (!have.has(classname)) { - allItems.push({ - classname, - sourceName: source.name, - sourceId: source.id, - }); - } - } - } - - if (allItems.length === 0) { - return new Response( - sseEvent({ - type: "batch_complete", - succeeded: 0, - failed: 0, - skipped: true, - message: "All sources are up to date", - }), - { headers: { "Content-Type": "text/event-stream" } }, - ); - } - - if (allItems.length > 2000) { - return new Response( - sseEvent({ - type: "error", - message: `Too many items to sync (${allItems.length}). Sync individual sources instead.`, - }), - { - status: 400, - headers: { "Content-Type": "text/event-stream" }, - }, - ); - } - - // FurnitureData entries are appended once at the end (single write) rather - // than once per item, removing the main serialization bottleneck. - const deferredEntries: Array<{ - entry: Record; - itemType: string; - }> = []; - - // Source lookup cache — pre-fetched above, no need to re-fetch per item - const sourceMap = new Map(sources.map((s) => [s.id, s])); - - return runSseBatch({ - items: allItems, - concurrency: 10, - signal: request.signal, - labelOf: (it) => `${it.sourceName}/${it.classname}`, - worker: async (it, _index, report) => { - const source = sourceMap.get(it.sourceId); - if (!source) { - return { ok: false, error: "source not found" }; - } - - // Use pre-fetched furnidata instead of re-fetching - const byClassname = sourceFurniDataMap.get(it.sourceId); - if (!byClassname) { - return { - ok: false, - error: "furnidata not available for this source", - }; - } - - const entry = byClassname.get(it.classname); - if (!entry) { - return { - ok: false, - error: "classname not found in source furnidata", - }; - } - - const result = await cloneSingleFurni({ - source, - entry, - onProgress: report, - deferFurniData: true, - }); - - if (result.ok) { - if (result.furniDataEntry) { - deferredEntries.push(result.furniDataEntry); - } - logAudit({ - userId: ctx.session.user.id, - action: "auto_sync_furni", - target: "FurniAsset", - targetId: 0, - after: { - classname: it.classname, - sourceId: it.sourceId, - sourceName: it.sourceName, - }, - }); - } - - return { - ok: result.ok, - warnings: result.warnings, - error: result.error, - }; - }, - flush: async () => { - if (deferredEntries.length === 0) return; - try { - await appendFurniEntriesBatch(deferredEntries); - } catch (err) { - for (const { entry, itemType } of deferredEntries) { - try { - await appendFurniEntry(entry, itemType); - } catch { - /* best effort */ - } - } - throw new Error( - `FurnitureData batch write failed: ${(err as Error).message}`, - ); - } - }, - }); - }, + (request, ctx) => enqueueFurnitureSync(request, ctx.session.user.id, "clone"), ); diff --git a/src/app/api/admin/import/official/sync-all/route.ts b/src/app/api/admin/import/official/sync-all/route.ts index 67dc7c74..41430705 100644 --- a/src/app/api/admin/import/official/sync-all/route.ts +++ b/src/app/api/admin/import/official/sync-all/route.ts @@ -1,97 +1,8 @@ import { withAdmin } from "@/lib/api-handler"; -import { db, ItemsBase } from "@/lib/db"; import { PERMS } from "@/lib/permissions"; -import { importSingleFurni } from "@/lib/services/furni-import"; -import { getOfficialHabboFurnidata } from "@/lib/services/habbo-furnidata-cache"; -import { runSseBatch } from "@/lib/services/import/core/sse-batch"; - +import { enqueueFurnitureSync } from "@/lib/services/furni-sync-queue"; export const POST = withAdmin( { permission: PERMS.ASSETS_IMPORT }, - async (request) => { - // Haal alle officiële furnidata - const furnidata = await getOfficialHabboFurnidata(); - - // Haal alle bestaande classnames uit de database - const existingClassnames = new Set( - await db - .select({ itemName: ItemsBase.itemName }) - .from(ItemsBase) - .then((rows) => rows.map((r) => r.itemName)), - ); - - // Filter naar alleen items die niet in de database staan - const missingItems = Array.from(furnidata.entries()) - .filter(([, entry]) => !existingClassnames.has(entry.classname)) - .map(([, entry]) => ({ - classname: entry.classname, - name: entry.name, - description: entry.description ?? "", - type: entry.category === "wallitem" ? "wallitem" : "flooritem", - revision: entry.revision ?? 0, - category: entry.category ?? "unknown", - })); - - if (missingItems.length === 0) { - return new Response( - JSON.stringify({ - success: true, - message: "All items are already in the database", - imported: 0, - }), - { - headers: { "Content-Type": "application/json" }, - }, - ); - } - - if (missingItems.length > 500) { - return new Response( - JSON.stringify({ - success: false, - message: `Too many missing items (${missingItems.length}). Consider importing in smaller batches.`, - }), - { - headers: { "Content-Type": "application/json" }, - status: 400, - }, - ); - } - - // Importeren via SSE batch - return runSseBatch({ - items: missingItems, - concurrency: 3, - signal: request.signal, - labelOf: (it) => it.classname, - worker: async (it, _index, report) => { - const result = await importSingleFurni({ - id: 0, - classname: it.classname, - name: it.name, - description: it.description, - type: it.type, - revision: it.revision, - category: it.category, - onProgress: report, - }); - - if (result.ok) { - return { - ok: true, - warnings: result.warnings, - error: result.error, - }; - } - - return { - ok: false, - warnings: result.warnings, - error: result.error, - }; - }, - flush: async () => { - // FurnitureData.json wordt automatisch geschreven door importSingleFurni - }, - }); - }, + (request, ctx) => + enqueueFurnitureSync(request, ctx.session.user.id, "official"), ); diff --git a/src/app/api/admin/studio/import-jobs/route.ts b/src/app/api/admin/studio/import-jobs/route.ts index f68fe674..e9254507 100644 --- a/src/app/api/admin/studio/import-jobs/route.ts +++ b/src/app/api/admin/studio/import-jobs/route.ts @@ -2,6 +2,7 @@ import { after } from "next/server"; import { z } from "zod"; import { apiError, apiOk } from "@/lib/api"; import { withAdmin } from "@/lib/api-handler"; +import { getRequestId } from "@/lib/foundation/request-context"; import { validateClassnames } from "@/lib/furni/studio-inspection"; import { PERMS } from "@/lib/permission-slugs"; import { redis } from "@/lib/redis"; @@ -92,6 +93,7 @@ export const POST = withAdmin( const job = await store.create({ ...body, userId: ctx.session.user.id, + operationId: getRequestId(), createdAt, updatedAt: createdAt, state: "queued", diff --git a/src/components/admin/studio/furniture-sync-queue.tsx b/src/components/admin/studio/furniture-sync-queue.tsx new file mode 100644 index 00000000..8c16b82c --- /dev/null +++ b/src/components/admin/studio/furniture-sync-queue.tsx @@ -0,0 +1,73 @@ +"use client"; +import { useTranslations } from "next-intl"; +import { useRef, useState } from "react"; +import { Button } from "@/components/ui/button"; +import { adminFetch } from "@/lib/admin-fetch"; +import { FurnitureJobHistory } from "./furniture-jobs"; +import { useFurnitureJobs } from "./use-furniture-jobs"; + +export function FurnitureSyncQueue({ kind }: { kind: "official" | "clone" }) { + const t = useTranslations("pages.admin.syncQueue"); + const history = useFurnitureJobs(() => {}, true); + const [busy, setBusy] = useState(false); + const [error, setError] = useState(""); + const [queued, setQueued] = useState(false); + const requestId = useRef(null); + const submitting = useRef(false); + async function start() { + if (submitting.current) return; + submitting.current = true; + setBusy(true); + setError(""); + setQueued(false); + // Keep the ID after an uncertain response, including a browser reload. + const key = `furniture-sync-request:${kind}`; + try { + if (!requestId.current) { + try { + requestId.current = sessionStorage.getItem(key); + } catch {} + requestId.current ||= crypto.randomUUID(); + } + try { + sessionStorage.setItem(key, requestId.current); + } catch {} + const response = await adminFetch(`/api/admin/import/${kind}/sync-all`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ id: requestId.current }), + }); + const result = await response.json().catch(() => ({})); + if (response.status >= 400 && response.status < 500) { + requestId.current = null; + try { + sessionStorage.removeItem(key); + } catch {} + } + if (!response.ok) throw Error(result.error || t("failed")); + if (result.job?.id !== requestId.current) throw Error(t("failed")); + try { + sessionStorage.removeItem(key); + } catch {} + requestId.current = null; + setQueued(true); + await history.refresh(); + } catch (error) { + setError(error instanceof Error ? error.message : t("failed")); + } finally { + submitting.current = false; + setBusy(false); + } + } + return ( +
+

{t("description")}

+ + {error &&

{error}

} + {queued &&

{t("queued")}

} + +
+ ); +} diff --git a/src/lib/furni/import-job.ts b/src/lib/furni/import-job.ts index adf2e846..1494bab3 100644 --- a/src/lib/furni/import-job.ts +++ b/src/lib/furni/import-job.ts @@ -1,3 +1,4 @@ +import type { SourceFurni } from "@/lib/services/clone-import"; export interface ImportJobItem { id: number; classname: string; @@ -7,9 +8,13 @@ export interface ImportJobItem { revision: number; category: string; attachmentId?: string; + cloneSourceId?: string; + cloneEntry?: SourceFurni; } export interface ImportJob { mode?: "repair"; + syncKind?: "official" | "clone"; + operationId?: string; retryOf?: string; id: string; userId: number; diff --git a/src/lib/services/furni-job-worker.test.ts b/src/lib/services/furni-job-worker.test.ts index 650ac113..a705d819 100644 --- a/src/lib/services/furni-job-worker.test.ts +++ b/src/lib/services/furni-job-worker.test.ts @@ -16,7 +16,7 @@ const mocks = vi.hoisted(() => ({ })); vi.mock("@/lib/redis", () => ({ redis: { set: mocks.set, eval: mocks.eval } })); vi.mock("@/lib/server-log", () => ({ logServerError: vi.fn() })); -vi.mock("./audit", () => ({ logAudit: vi.fn() })); +vi.mock("./audit", () => ({ logAudit: vi.fn().mockResolvedValue(undefined) })); vi.mock("./catalog-git-queue", () => ({ withCatalogExport: mocks.export })); vi.mock("./clone-sources", () => ({ getSource: async () => null })); vi.mock("./furni-job-store", () => ({ @@ -42,6 +42,8 @@ vi.mock("./rcon", () => ({ vi.mock("./catalog-git-export", () => ({ runCatalogExport: vi.fn() })); +vi.mock("./furni-sync-item", () => ({ runSyncJobItem: mocks.import })); + import { drainFurnitureImports } from "./furni-job-worker"; let job: ImportJob; @@ -223,3 +225,27 @@ it("does not save a translated outcome after ownership is lost during translatio expect(snapshots.at(-1)?.items[0].state).toBe("running"); expect(mocks.updateCatalog).not.toHaveBeenCalled(); }); + +it("restores the original operation context while running a durable job", async () => { + const { getOperationContext } = await import( + "@/lib/foundation/request-context" + ); + job.operationId = "original-operation"; + mocks.import.mockImplementation(async () => { + expect(getOperationContext()).toEqual({ + operationId: "original-operation", + userId: 5, + }); + return { ok: true, itemId: 900, warnings: [] }; + }); + await drainFurnitureImports(); + expect(getOperationContext()).toEqual({}); +}); + +it("keeps successful imports complete if their audit write fails", async () => { + const { logAudit } = await import("./audit"); + vi.mocked(logAudit).mockRejectedValueOnce(new Error("Audit unavailable")); + await drainFurnitureImports(); + expect(job.items[0].state).toBe("done"); + expect(job.items[0].warnings).toContain("Audit record could not be saved"); +}); diff --git a/src/lib/services/furni-job-worker.ts b/src/lib/services/furni-job-worker.ts index 07f71ccc..438c9ae8 100644 --- a/src/lib/services/furni-job-worker.ts +++ b/src/lib/services/furni-job-worker.ts @@ -1,4 +1,6 @@ import { randomUUID } from "node:crypto"; +import { createStore, runWithStore } from "@/lib/foundation/request-context"; +import type { IpAddress, RequestId, UserId } from "@/lib/foundation/types"; import type { ImportJob } from "@/lib/furni/import-job"; import { redis } from "@/lib/redis"; import { logServerError } from "@/lib/server-log"; @@ -10,6 +12,7 @@ import { readFurnitureAttachment } from "./furni-attachment"; import { patchLocalizedFurniDataEntries } from "./furni-data-i18n"; import { ensureDirectories, importSingleFurni } from "./furni-import"; import { ImportJobStore } from "./furni-job-store"; +import { runSyncJobItem } from "./furni-sync-item"; import { repairFurniture } from "./furniture-repair"; import { rcon } from "./rcon"; @@ -100,87 +103,108 @@ async function drain() { } job.state = "running"; await saveOwned(job); - await withCatalogExport(async () => { - await ensureDirectories(); - const source = job.sourceId ? await getSource(job.sourceId) : null; - for (const item of job.items) { - await requireLease(); - if (item.state !== "pending") continue; - if (await store.isCancellationRequested(job.id)) { - for (const remaining of job.items) - if (remaining.state === "pending") remaining.state = "cancelled"; - break; - } - item.state = "running"; - await saveOwned(job); - try { - if (job.sourceId && !source) - throw Error("Import source no longer exists"); - const providedNitro = item.attachmentId - ? await readFurnitureAttachment( - item.attachmentId, - item.classname, - job.userId, - ) - : undefined; + const operation = createStore("worker" as IpAddress); + operation.requestId = (job.operationId ?? job.id) as RequestId; + operation.userId = job.userId as UserId; + await runWithStore(operation, () => + withCatalogExport(async () => { + await ensureDirectories(); + const source = job.sourceId ? await getSource(job.sourceId) : null; + for (const item of job.items) { await requireLease(); - const result = await (job.mode === "repair" - ? repairFurniture - : importSingleFurni)({ - ...item, - repairExisting: true, - providedNitro, - sourceSwfBaseUrl: source?.sourceSwfBaseUrl, - nitroBaseUrl: source?.nitroBaseUrl, - iconBaseUrl: source?.iconBaseUrl, - }); - await requireLease(); - item.warnings = result.warnings; - item.itemId = result.itemId; - if (!result.ok) throw Error(result.error || "Import failed"); - item.state = "done"; - logAudit({ - userId: job.userId, - action: "furni_import", - target: "ItemsBase", - targetId: result.itemId ?? 0, - after: { - classname: item.classname, + if (item.state !== "pending") continue; + if (await store.isCancellationRequested(job.id)) { + for (const remaining of job.items) + if (remaining.state === "pending") + remaining.state = "cancelled"; + break; + } + item.state = "running"; + await saveOwned(job); + try { + if (job.sourceId && !source) + throw Error("Import source no longer exists"); + const providedNitro = item.attachmentId + ? await readFurnitureAttachment( + item.attachmentId, + item.classname, + job.userId, + ) + : undefined; + await requireLease(); + const result = job.syncKind + ? await runSyncJobItem(item) + : await (job.mode === "repair" + ? repairFurniture + : importSingleFurni)({ + ...item, + repairExisting: true, + providedNitro, + sourceSwfBaseUrl: source?.sourceSwfBaseUrl, + nitroBaseUrl: source?.nitroBaseUrl, + iconBaseUrl: source?.iconBaseUrl, + }); + await requireLease(); + item.warnings = result.warnings; + item.itemId = result.itemId; + if (!result.ok) throw Error(result.error || "Import failed"); + item.state = "done"; + await logAudit({ + userId: job.userId, + action: "furni_import", + target: "ItemsBase", + targetId: result.itemId ?? 0, + after: { + classname: item.classname, + jobId: job.id, + operationId: job.operationId, + syncKind: job.syncKind, + repairExisting: true, + }, + }).catch((auditError) => { + item.warnings ??= []; + item.warnings.push("Audit record could not be saved"); + logServerError("furni.audit_failed", auditError, { + jobId: job.id, + classname: item.classname, + }); + }); + if (job.translate && job.mode !== "repair") + try { + await patchLocalizedFurniDataEntries([item], true, job.langs); + } catch { + item.warnings.push( + "Translation failed; furniture imported successfully", + ); + } + } catch (error) { + if (error instanceof ImportLeaseLostError) throw error; + await requireLease(); + logServerError("furni.item_failed", error, { jobId: job.id, - repairExisting: true, - }, - }); - if (job.translate && job.mode !== "repair") - try { - await patchLocalizedFurniDataEntries([item], true, job.langs); - } catch { - item.warnings.push( - "Translation failed; furniture imported successfully", - ); - } + classname: item.classname, + }); + item.state = "failed"; + item.error = + error instanceof Error ? error.message : "Import failed"; + } + await saveOwned(job); + } + await requireLease(); + try { + await rcon.updateCatalog(); + await requireLease(); + await rcon.updateItems(); } catch (error) { if (error instanceof ImportLeaseLostError) throw error; - await requireLease(); - item.state = "failed"; - item.error = - error instanceof Error ? error.message : "Import failed"; + for (const item of job.items) + if (item.state === "done") { + item.warnings ??= []; + item.warnings.push("Game cache refresh failed"); + } } - await saveOwned(job); - } - await requireLease(); - try { - await rcon.updateCatalog(); - await requireLease(); - await rcon.updateItems(); - } catch (error) { - if (error instanceof ImportLeaseLostError) throw error; - for (const item of job.items) - if (item.state === "done") { - item.warnings ??= []; - item.warnings.push("Game cache refresh failed"); - } - } - }); + }), + ); await requireLease(); job.state = job.items.some((item) => item.state === "cancelled") ? "cancelled" diff --git a/src/lib/services/furni-sync-item.test.ts b/src/lib/services/furni-sync-item.test.ts new file mode 100644 index 00000000..86e50695 --- /dev/null +++ b/src/lib/services/furni-sync-item.test.ts @@ -0,0 +1,75 @@ +import { beforeEach, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + existing: vi.fn(), + clone: vi.fn(), + source: vi.fn(), + import: vi.fn(), +})); +vi.mock("@/lib/db", async () => ({ + ...(await import("@/db/schema")), + db: { + select: () => ({ + from: () => ({ where: () => ({ limit: mocks.existing }) }), + }), + }, +})); +vi.mock("./clone-import", () => ({ cloneSingleFurni: mocks.clone })); +vi.mock("./clone-sources", () => ({ getSource: mocks.source })); +vi.mock("./furni-import", () => ({ importSingleFurni: mocks.import })); + +import { runSyncJobItem } from "./furni-sync-item"; + +const item = { + id: 0, + classname: "chair", + name: "Chair", + description: "", + type: "flooritem", + revision: 1, + category: "other", +}; +beforeEach(() => { + vi.clearAllMocks(); + mocks.existing.mockResolvedValue([]); + mocks.import.mockResolvedValue({ ok: true, warnings: [] }); +}); +it("skips an item another queued request has already imported", async () => { + mocks.existing.mockResolvedValue([{ id: 3 }]); + expect(await runSyncJobItem(item)).toMatchObject({ ok: true, itemId: 3 }); + expect(mocks.import).not.toHaveBeenCalled(); + expect(mocks.clone).not.toHaveBeenCalled(); +}); +it("never repairs existing furniture during official sync", async () => { + await runSyncJobItem(item); + expect(mocks.import).toHaveBeenCalledWith({ ...item, repairExisting: false }); +}); +it("fails a deleted clone source without substituting official assets", async () => { + mocks.source.mockResolvedValue(null); + await expect( + runSyncJobItem({ ...item, cloneSourceId: "gone" }), + ).rejects.toThrow("no longer available"); + expect(mocks.import).not.toHaveBeenCalled(); +}); + +it("returns the created clone item ID for audit correlation", async () => { + mocks.existing.mockResolvedValueOnce([]).mockResolvedValueOnce([{ id: 71 }]); + mocks.source.mockResolvedValue({ id: "source" }); + mocks.clone.mockResolvedValue({ ok: true, warnings: [] }); + const cloneEntry = { + id: 1, + classname: "chair", + name: "Chair", + description: "", + xdim: 1, + ydim: 1, + canstandon: false, + cansiton: true, + canlayon: false, + customparams: "", + itemType: "s" as const, + }; + expect( + await runSyncJobItem({ ...item, cloneSourceId: "source", cloneEntry }), + ).toMatchObject({ ok: true, itemId: 71 }); +}); diff --git a/src/lib/services/furni-sync-item.ts b/src/lib/services/furni-sync-item.ts new file mode 100644 index 00000000..028db9b2 --- /dev/null +++ b/src/lib/services/furni-sync-item.ts @@ -0,0 +1,36 @@ +import { eq } from "drizzle-orm"; +import { db, ItemsBase } from "@/lib/db"; +import type { ImportJobItem } from "@/lib/furni/import-job"; +import { cloneSingleFurni } from "./clone-import"; +import { getSource } from "./clone-sources"; +import { importSingleFurni } from "./furni-import"; + +/** Called only inside the shared import lease and catalog export lock. */ +export async function runSyncJobItem(item: ImportJobItem) { + const [existing] = await db + .select({ id: ItemsBase.id }) + .from(ItemsBase) + .where(eq(ItemsBase.itemName, item.classname)) + .limit(1); + if (existing) + return { + ok: true, + itemId: existing.id, + warnings: ["Already present; synchronization skipped this item."], + }; + if (item.cloneSourceId) { + const source = await getSource(item.cloneSourceId); + if (!source || !item.cloneEntry) + throw Error("Synchronization source no longer available"); + const result = await cloneSingleFurni({ source, entry: item.cloneEntry }); + const [created] = result.ok + ? await db + .select({ id: ItemsBase.id }) + .from(ItemsBase) + .where(eq(ItemsBase.itemName, item.classname)) + .limit(1) + : []; + return { ...result, itemId: created?.id }; + } + return importSingleFurni({ ...item, repairExisting: false }); +} diff --git a/src/lib/services/furni-sync-queue.test.ts b/src/lib/services/furni-sync-queue.test.ts new file mode 100644 index 00000000..4e0088dc --- /dev/null +++ b/src/lib/services/furni-sync-queue.test.ts @@ -0,0 +1,113 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + read: vi.fn(), + create: vi.fn(), + official: vi.fn(), + sources: vi.fn(), + fetch: vi.fn(), + ping: vi.fn(), + drain: vi.fn(), + after: vi.fn(), + existing: vi.fn(), +})); +vi.mock("next/server", () => ({ after: mocks.after, NextResponse: Response })); +vi.mock("@/lib/db", () => ({ + ItemsBase: { itemName: "item_name" }, + db: { select: () => ({ from: mocks.existing }) }, +})); +vi.mock("@/lib/redis", () => ({ redis: { ping: mocks.ping } })); +vi.mock("./furni-job-store", () => ({ + validJobId: (id: unknown) => typeof id === "string" && id.length === 36, + ImportJobStore: class { + read = mocks.read; + create = mocks.create; + }, +})); +vi.mock("./furni-job-worker", () => ({ drainFurnitureImports: mocks.drain })); +vi.mock("./clone-sources", () => ({ listSources: mocks.sources })); +vi.mock("./clone-import", () => ({ fetchSourceFurnidata: mocks.fetch })); +vi.mock("./habbo-furnidata-cache", () => ({ + getOfficialHabboFurnidata: mocks.official, +})); + +import { enqueueFurnitureSync } from "./furni-sync-queue"; + +const id = "11111111-1111-4111-8111-111111111111"; +const request = (body: unknown = { id }) => + new Request("https://example.test/api/sync", { + method: "POST", + body: JSON.stringify(body), + }); +beforeEach(() => { + vi.clearAllMocks(); + mocks.read.mockRejectedValue({ code: "ENOENT" }); + mocks.create.mockImplementation(async (job) => job); + mocks.ping.mockResolvedValue("PONG"); + mocks.existing.mockResolvedValue([]); + mocks.official.mockResolvedValue(new Map()); + mocks.sources.mockResolvedValue([]); +}); +describe("durable synchronization enqueue", () => { + it("persists an empty completed job when already up to date", async () => { + const res = await enqueueFurnitureSync(request(), 7, "official"); + expect(res.status).toBe(200); + expect(mocks.create).toHaveBeenCalledWith( + expect.objectContaining({ state: "completed", items: [], userId: 7 }), + ); + }); + it("reuses an existing request before source discovery", async () => { + const job = { id, userId: 7, syncKind: "official" }; + mocks.read.mockResolvedValue(job); + expect((await enqueueFurnitureSync(request(), 7, "official")).status).toBe( + 200, + ); + expect(mocks.create).not.toHaveBeenCalled(); + expect(mocks.official).not.toHaveBeenCalled(); + }); + it("rejects other owners and changed operation kinds", async () => { + mocks.read.mockResolvedValue({ id, userId: 8, syncKind: "official" }); + expect((await enqueueFurnitureSync(request(), 7, "official")).status).toBe( + 409, + ); + mocks.read.mockResolvedValue({ id, userId: 7, syncKind: "clone" }); + expect((await enqueueFurnitureSync(request(), 7, "official")).status).toBe( + 409, + ); + }); + it("does not enqueue while redis is unavailable", async () => { + mocks.ping.mockRejectedValue(Error("offline")); + expect((await enqueueFurnitureSync(request(), 7, "official")).status).toBe( + 503, + ); + expect(mocks.create).not.toHaveBeenCalled(); + }); + it("deduplicates classnames across source snapshots", async () => { + mocks.sources.mockResolvedValue([ + { id: "a", furnidataUrl: "a" }, + { id: "b", furnidataUrl: "b" }, + ]); + mocks.fetch.mockResolvedValue([ + { + id: 1, + classname: "chair", + name: "Chair", + description: "", + itemType: "s", + }, + ]); + await enqueueFurnitureSync(request(), 7, "clone"); + const job = mocks.create.mock.calls[0][0]; + expect(job.items).toHaveLength(1); + expect(job.items[0].cloneSourceId).toBe("a"); + expect(job.items[0].cloneEntry.classname).toBe("chair"); + }); + it("never queues partial discovery after a source failure", async () => { + mocks.sources.mockResolvedValue([{ id: "a", furnidataUrl: "a" }]); + mocks.fetch.mockRejectedValue(Error("unavailable")); + await expect(enqueueFurnitureSync(request(), 7, "clone")).rejects.toThrow( + "unavailable", + ); + expect(mocks.create).not.toHaveBeenCalled(); + }); +}); diff --git a/src/lib/services/furni-sync-queue.ts b/src/lib/services/furni-sync-queue.ts new file mode 100644 index 00000000..e05a8bef --- /dev/null +++ b/src/lib/services/furni-sync-queue.ts @@ -0,0 +1,114 @@ +import { after } from "next/server"; +import { apiError, apiOk } from "@/lib/api"; +import { db, ItemsBase } from "@/lib/db"; +import { getRequestId } from "@/lib/foundation/request-context"; +import type { ImportJob, ImportJobItem } from "@/lib/furni/import-job"; +import { validateClassnames } from "@/lib/furni/studio-inspection"; +import { redis } from "@/lib/redis"; +import { fetchSourceFurnidata } from "./clone-import"; +import { listSources } from "./clone-sources"; +import { ImportJobStore, validJobId } from "./furni-job-store"; +import { drainFurnitureImports } from "./furni-job-worker"; +import { getOfficialHabboFurnidata } from "./habbo-furnidata-cache"; + +export async function enqueueFurnitureSync( + request: Request, + userId: number, + kind: "official" | "clone", +) { + const body = await request.json().catch(() => null); + if ( + !validJobId(body?.id) || + (kind === "official" && body.sourceId !== undefined) || + (body.sourceId !== undefined && + (typeof body.sourceId !== "string" || body.sourceId.length > 100)) + ) + return apiError("Invalid synchronization request", 400); + const store = new ImportJobStore(); + const existing = await store.read(body.id).catch((error) => { + if (error.code === "ENOENT") return null; + throw error; + }); + if (existing) { + if ( + existing.userId !== userId || + existing.syncKind !== kind || + existing.sourceId !== body.sourceId + ) + return apiError("Request ID already used", 409); + after(drainFurnitureImports); + return apiOk({ job: existing }); + } + if (!redis || (await redis.ping().catch(() => null)) !== "PONG") + return apiError("Background import queue is temporarily unavailable", 503); + const have = new Set( + (await db.select({ name: ItemsBase.itemName }).from(ItemsBase)).map( + (row) => row.name, + ), + ); + const items = new Map(); + if (kind === "official") { + for (const entry of (await getOfficialHabboFurnidata()).values()) { + if (have.has(entry.classname)) continue; + items.set(entry.classname, { + id: 0, + classname: entry.classname, + name: entry.name, + description: entry.description ?? "", + type: entry.category === "wallitem" ? "wallitem" : "flooritem", + revision: entry.revision ?? 0, + category: entry.category ?? "unknown", + }); + } + } else { + const sources = (await listSources()).filter( + (source) => !body.sourceId || source.id === body.sourceId, + ); + if (!sources.length) return apiError("No sources configured", 400); + // Discovery is read-only. Persist one complete snapshot, never a partial request. + for (const source of sources) + for (const entry of await fetchSourceFurnidata(source.furnidataUrl)) { + if (have.has(entry.classname) || items.has(entry.classname)) continue; + items.set(entry.classname, { + id: entry.id, + classname: entry.classname, + name: entry.name, + description: entry.description, + type: entry.itemType === "i" ? "wallitem" : "flooritem", + revision: Number(entry.revision) || 0, + category: String(entry.category ?? "unknown"), + cloneSourceId: source.id, + cloneEntry: entry, + }); + } + } + if (items.size > (kind === "official" ? 500 : 2000)) + return apiError( + "Too many missing items. Import a smaller selection from Studio.", + 400, + ); + if ([...items.keys()].some((classname) => !validateClassnames([classname]))) + return apiError("Source contains invalid furniture identifiers", 400); + const now = new Date().toISOString(); + const job: ImportJob = { + id: body.id, + userId, + syncKind: kind, + operationId: getRequestId(), + sourceId: body.sourceId, + translate: false, + createdAt: now, + updatedAt: now, + state: items.size ? "queued" : "completed", + items: [...items.values()].map((item) => ({ ...item, state: "pending" })), + }; + const saved = await store.create(job); + if ( + saved.userId !== userId || + saved.syncKind !== kind || + saved.sourceId !== body.sourceId + ) + return apiError("Request ID already used", 409); + after(drainFurnitureImports); + return apiOk({ job: saved }); +}