import { withAdmin } from "@/lib/api-handler"; import { apiError } from "@/lib/api-response"; import { PERMS } from "@/lib/permissions"; import { logAudit } from "@/lib/services/audit"; import { getSource } from "@/lib/services/clone-sources"; import { appendFurniEntriesBatch, rebuildCatalogOfferIds, reconcileFurniDataWithItemsBase, verifyAndFixSpriteIds, } from "@/lib/services/furni-data"; import { buildLocalizedFurniDataFiles } from "@/lib/services/furni-data-i18n"; import { ensureDirectories, ensureFurniOwnership, fixDatabaseConsistencyAfterImport, importSingleFurni, reconcileImportedOfferIds, syncAssetsToGamedataBundle, verifyAndFixInteractionModesCount, } from "@/lib/services/furni-import"; import { rcon } from "@/lib/services/rcon"; import type { ImportSingleResult } from "@/types/furni"; interface BatchItem { id: number; classname: string; name: string; description: string; type: string; revision: number; category: string; } /** * SSE batch import endpoint. * Processes multiple furniture items with configurable concurrency, * streaming real-time progress events to the client. */ export const POST = withAdmin( { permission: PERMS.ASSETS_IMPORT }, async (request, ctx) => { const body = await request.json(); const rawItems: BatchItem[] = body.items || []; const concurrency = Math.min(Math.max(body.concurrency || 1, 1), 3); const sourceId: string | undefined = body.sourceId; if (rawItems.length === 0) return apiError("No items to import", 400); // Look up clone source for SWF/nitro/icon URL fallbacks let source: Awaited> | null = null; if (sourceId) { source = await getSource(sourceId); } // Deduplicate by (spriteId, classname) to avoid concurrent INSERT races on // items_base.id when the same item is listed twice in the request body. const seen = new Set(); const items: BatchItem[] = []; for (const it of rawItems) { const key = `${it.id}:${it.classname}`; if (seen.has(key)) continue; seen.add(key); items.push(it); } await ensureDirectories(); const encoder = new TextEncoder(); // Wire client disconnect to abort controller so we stop processing // when the user navigates away or closes the browser. const abortController = new AbortController(); let aborted = false; request.signal.addEventListener("abort", () => { aborted = true; abortController.abort(); }); const stream = new ReadableStream({ async start(controller) { const send = (data: unknown) => { if (aborted) return; try { controller.enqueue( encoder.encode(`data: ${JSON.stringify(data)}\n\n`), ); } catch { /* stream closed by client */ } }; const startTime = Date.now(); send({ type: "batch_start", total: items.length, concurrency }); let succeeded = 0; let failed = 0; let withWarnings = 0; const furniDataEntries: Array<{ entry: Record; itemType: string; }> = []; // Process in chunks of `concurrency` for (let i = 0; i < items.length; i += concurrency) { if (aborted) break; const chunk = items.slice(i, i + concurrency); const promises = chunk.map(async (item, chunkIdx) => { const index = i + chunkIdx; send({ type: "item_progress", classname: item.classname, status: "started", index, }); try { const result: ImportSingleResult = await importSingleFurni({ id: item.id ?? 0, classname: item.classname, name: item.name, description: item.description ?? "", type: item.type ?? "flooritem", revision: item.revision ?? 0, category: item.category ?? "unknown", skipFurniDataWrite: false, repairExisting: body.repairExisting === true, sourceSwfBaseUrl: source?.sourceSwfBaseUrl, nitroBaseUrl: source?.nitroBaseUrl, iconBaseUrl: source?.iconBaseUrl, onProgress: (status: string) => { send({ type: "item_progress", classname: item.classname, status, index, }); }, }); if (result.ok) { succeeded++; if (result.warnings.length > 0) withWarnings++; if (result.furniDataEntry) furniDataEntries.push(result.furniDataEntry); send({ type: "item_progress", classname: item.classname, status: "done", index, itemId: result.itemId, warnings: result.warnings.length > 0 ? result.warnings : undefined, }); logAudit({ userId: ctx.session.user.id, action: "furni_import", target: "ItemsBase", targetId: result.itemId ?? 0, after: { classname: item.classname, name: item.name, type: item.type === "wallitem" ? "i" : "s", catalogItemId: result.catalogItemId, dimensions: result.dimensions, batch: true, }, }); } else { failed++; send({ type: "item_progress", classname: item.classname, status: "failed", index, message: result.error, }); } } catch (err) { failed++; send({ type: "item_progress", classname: item.classname, status: "failed", index, message: (err as Error).message, }); } }); await Promise.allSettled(promises); await new Promise((r) => setImmediate(r)); } if (aborted) { send({ type: "batch_complete", succeeded, failed, warnings: withWarnings, aborted: true, duration: Date.now() - startTime, }); controller.close(); return; } // Post-import steps: send progress so the user knows what's happening send({ type: "post_import", status: "Writing FurnitureData.json..." }); // Batch write FurnitureData.json once at the end — BEFORE RCON so the // emulator never sees items that the Nitro client can't render. let furniDataWriteOk = true; if (furniDataEntries.length > 0) { try { await appendFurniEntriesBatch(furniDataEntries); } catch (err) { furniDataWriteOk = false; send({ type: "item_progress", classname: "_batch_", status: "failed", index: -1, message: `FurnitureData.json batch write failed: ${(err as Error).message}`, }); } } // Refresh game server caches once — skip if FurnitureData write failed if (furniDataWriteOk && !aborted) { send({ type: "post_import", status: "Refreshing emulator caches...", }); try { const okCat = await rcon.updateCatalog(); const okItems = await rcon.updateItems(); if (!okCat || !okItems) { send({ type: "item_progress", classname: "_rcon_", status: "failed", index: -1, message: "RCON refresh failed — emulator cache may be stale", }); } } catch (err) { console.warn( "[import-furni] RCON update failed after batch:", (err as Error).message, ); } } if (!aborted) { send({ type: "post_import", status: "Reconciling offer IDs & ownership...", }); // Post-import reconciliation — wrapped in try/catch so one // failure doesn't block the rest. let offerIdsFixed = 0; let ownershipFixed: string[] = []; let nitrosSynced: string[] = []; let iconsSynced: string[] = []; let furniDataFixedIds = 0; let furniDataFixedOfferIds = 0; let furniDataMissing = 0; let furniDataConflicts = 0; let offerRebuildFixed = 0; let spriteFixed = 0; let interactionFixed = 0; let languages: Awaited< ReturnType > = []; try { offerIdsFixed = (await reconcileImportedOfferIds()).fixed; ownershipFixed = (await ensureFurniOwnership()).fixed; const assetsSynced = await syncAssetsToGamedataBundle(); nitrosSynced = assetsSynced.copiedNitros; iconsSynced = assetsSynced.copiedIcons; const furniReconcile = await reconcileFurniDataWithItemsBase(); furniDataFixedIds = furniReconcile.fixedIds; furniDataFixedOfferIds = furniReconcile.fixedOfferIds; furniDataMissing = furniReconcile.missing; furniDataConflicts = furniReconcile.conflicts; offerRebuildFixed = (await rebuildCatalogOfferIds()).fixed; spriteFixed = (await verifyAndFixSpriteIds()).fixed; if (!aborted) { send({ type: "post_import", status: "Verifying interaction modes...", }); interactionFixed = (await verifyAndFixInteractionModesCount()) .fixed; } try { languages = await buildLocalizedFurniDataFiles(); } catch (err) { console.warn( "[import-furni] Localized furnidata build failed:", (err as Error).message, ); } } catch (err) { console.warn( "[import-furni] Post-import reconcile/ownership failed:", (err as Error).message, ); } send({ type: "post_import", status: "Fixing database consistency...", }); const { catalogNameFixed, haveOfferFixed, costCreditsFixed } = await fixDatabaseConsistencyAfterImport(); send({ type: "batch_complete", succeeded, failed, warnings: withWarnings, offerIdsFixed, ownershipFixed, nitrosSynced, iconsSynced, furniDataFixedIds, furniDataFixedOfferIds, furniDataMissing, furniDataConflicts, offerRebuildFixed, spriteFixed, interactionFixed, catalogNameFixed, haveOfferFixed, costCreditsFixed, languages, duration: Date.now() - startTime, }); } else { send({ type: "batch_complete", succeeded, failed, warnings: withWarnings, aborted: true, duration: Date.now() - startTime, }); } controller.close(); }, }); return new Response(stream, { headers: { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", Connection: "keep-alive", }, }); }, );