import { withAdmin } from "@/lib/api-handler"; import { apiError } from "@/lib/api-response"; import { PERMS } from "@/lib/permissions"; import { prisma } from "@/lib/prisma"; import { logAudit } from "@/lib/services/audit"; import { cloneSingleFurni, fetchSourceFurnidata, } from "@/lib/services/clone-import"; import { listSources } from "@/lib/services/clone-sources"; 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" }, }); } const allItems: Array<{ classname: string; sourceName: string; sourceId: string; }> = []; for (const source of sources) { let entries: Awaited>; try { entries = await fetchSourceFurnidata(source.furnidataUrl); } catch { continue; } const classnames = entries.map((e) => e.classname); if (classnames.length === 0) continue; const placeholders = classnames.map(() => "?").join(","); const rows = await prisma.$queryRawUnsafe>( `SELECT item_name FROM items_base WHERE item_name IN (${placeholders})`, ...classnames, ); const have = new Set(rows.map((r) => r.item_name)); for (const entry of entries) { if (!have.has(entry.classname)) { allItems.push({ classname: entry.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" }, }, ); } return runSseBatch({ items: allItems, concurrency: 2, labelOf: (it) => `${it.sourceName}/${it.classname}`, worker: async (it, _index, report) => { const allSrc = await listSources(); const source = allSrc.find((s) => s.id === it.sourceId); if (!source) { return { ok: false, error: "source not found" }; } let allEntries: Awaited>; try { allEntries = await fetchSourceFurnidata(source.furnidataUrl); } catch (err) { return { ok: false, error: `furnidata fetch failed: ${(err as Error).message}`, }; } const entry = allEntries.find((e) => e.classname === it.classname); if (!entry) { return { ok: false, error: "classname not found in source furnidata", }; } const result = await cloneSingleFurni({ source, entry, onProgress: report, }); if (result.ok) { 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, }; }, }); }, );