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`; } 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}`, ); } }, }); }, );