feat(catalog): queue furniture imports with history and matching file attachments
CI / check (push) Successful in 1m30s
CI / deploy (push) Successful in 1m26s
CI / e2e (push) Successful in 25s

This commit is contained in:
Simo committed 2026-09-05 17:02:26 +02:00
1 parent 775d14f861
commit 2619bec165
16 files changed
+1016 -232

No files matched your search

@@ -0,0 +1,37 @@
import { withAdmin } from "@/lib/api-handler";
import { apiError, apiOk } from "@/lib/api-response";
import { validateClassnames } from "@/lib/furni/studio-inspection";
import { PERMS } from "@/lib/permission-slugs";
import { stageFurnitureAttachment } from "@/lib/services/furni-attachment";
export const POST = withAdmin(
{ permission: PERMS.ASSETS_IMPORT, maxBodyBytes: 52 * 1024 * 1024 },
async (request, ctx) => {
const form = await request.formData(),
classname = form.get("classname"),
file = form.get("file");
if (
typeof classname !== "string" ||
!validateClassnames([classname]) ||
!(file instanceof File) ||
file.size === 0 ||
file.size > 50 * 1024 * 1024
)
return apiError(
"Select a .nitro file for this furniture (maximum 50 MB)",
400,
);
try {
const attachmentId = await stageFurnitureAttachment(
Buffer.from(await file.arrayBuffer()),
classname,
ctx.session.user.id,
);
return apiOk({ attachmentId });
} catch (error) {
return apiError(
error instanceof Error ? error.message : "Invalid .nitro file",
400,
);
}
},
);
@@ -0,0 +1,116 @@
import { randomUUID } from "node:crypto";
import { NextRequest } from "next/server";
import { beforeEach, expect, it, vi } from "vitest";
import { PERMS } from "@/lib/permission-slugs";
const mocks = vi.hoisted(() => ({
guard: vi.fn(),
create: vi.fn(),
list: vi.fn(),
after: vi.fn(),
ping: vi.fn(),
source: vi.fn(),
}));
vi.mock("@/lib/api-handler", () => ({
withAdmin: (options: unknown, handler: unknown) => {
mocks.guard(options);
return handler;
},
}));
vi.mock("next/server", async (original) => ({
...(await original<typeof import("next/server")>()),
after: mocks.after,
}));
vi.mock("@/lib/redis", () => ({ redis: { ping: mocks.ping } }));
vi.mock("@/lib/services/furni-job-worker", () => ({
drainFurnitureImports: vi.fn(),
}));
vi.mock("@/lib/services/furni-job-store", () => ({
ImportJobStore: class {
create = mocks.create;
list = mocks.list;
},
}));
vi.mock("@/lib/services/clone-sources", () => ({ getSource: mocks.source }));
import { GET, POST } from "./route";
const item = {
id: 1,
classname: "nft_china_light",
name: "Lamp",
description: "",
type: "wallitem",
revision: 70317,
category: "other",
};
const ctx = { session: { user: { id: 7 } } } as never;
const request = (body: unknown) =>
new NextRequest("http://localhost/api/admin/studio/import-jobs", {
method: "POST",
body: JSON.stringify(body),
});
beforeEach(() => {
mocks.create.mockClear();
mocks.list.mockClear();
mocks.after.mockClear();
mocks.ping.mockResolvedValue("PONG");
mocks.source.mockResolvedValue(null);
mocks.create.mockImplementation((job) => job);
mocks.list.mockResolvedValue([]);
});
it("requires asset import permission", () => {
expect(mocks.guard).toHaveBeenCalledWith({ permission: PERMS.ASSETS_IMPORT });
});
it("queues a durable owner-scoped job and schedules execution after responding", async () => {
const response = await POST(
request({ id: randomUUID(), items: [item, item] }),
ctx,
);
expect(response.status).toBe(200);
expect(mocks.create.mock.calls[0][0]).toMatchObject({
userId: 7,
state: "queued",
items: [{ ...item, state: "pending" }],
});
expect(mocks.after).toHaveBeenCalledOnce();
});
it("rejects traversal before saving work", async () => {
expect(
(
await POST(
request({
id: randomUUID(),
items: [{ ...item, classname: "../bad" }],
}),
ctx,
)
).status,
).toBe(400);
expect(mocks.create).not.toHaveBeenCalled();
});
it("does not accept work when the worker lease service is unavailable", async () => {
mocks.ping.mockRejectedValue(Error("offline"));
expect(
(await POST(request({ id: randomUUID(), items: [item] }), ctx)).status,
).toBe(503);
expect(mocks.create).not.toHaveBeenCalled();
});
it("rejects a missing configured source", async () => {
expect(
(
await POST(
request({ id: randomUUID(), sourceId: "gone", items: [item] }),
ctx,
)
).status,
).toBe(404);
});
it("only returns the current operators history", async () => {
mocks.list.mockResolvedValue([
{ id: "a", userId: 7 },
{ id: "b", userId: 9 },
]);
const response = await GET(request({}), ctx);
expect((await response.json()).jobs).toEqual([{ id: "a", userId: 7 }]);
});
@@ -0,0 +1,83 @@
import { after } from "next/server";
import { z } from "zod";
import { withAdmin } from "@/lib/api-handler";
import { apiError, apiOk } from "@/lib/api-response";
import { validateClassnames } from "@/lib/furni/studio-inspection";
import { PERMS } from "@/lib/permission-slugs";
import { redis } from "@/lib/redis";
import { getSource } from "@/lib/services/clone-sources";
import { ImportJobStore } from "@/lib/services/furni-job-store";
import { drainFurnitureImports } from "@/lib/services/furni-job-worker";
const schema = z.object({
id: z.uuid().refine((value) => value[14] === "4"),
sourceId: z.string().max(100).optional(),
translate: z.boolean().default(false),
langs: z
.array(z.string().regex(/^[a-z]{2,3}$/))
.max(25)
.optional(),
items: z
.array(
z.object({
id: z.number().int().nonnegative(),
classname: z.string(),
name: z.string().min(1).max(500),
description: z.string().max(5000),
type: z.enum(["flooritem", "wallitem"]),
revision: z.number().int().nonnegative(),
category: z.string().max(200),
attachmentId: z.uuid().optional(),
}),
)
.min(1)
.max(500),
});
export const GET = withAdmin(
{ permission: PERMS.ASSETS_IMPORT },
async (_request, ctx) => {
after(drainFurnitureImports);
const jobs = (await new ImportJobStore().list())
.filter((job) => job.userId === ctx.session.user.id)
.reverse()
.slice(0, 30);
return apiOk({ jobs });
},
);
export const POST = withAdmin(
{ permission: PERMS.ASSETS_IMPORT },
async (request, ctx) => {
const parsed = schema.safeParse(await request.json().catch(() => null));
if (
!parsed.success ||
!validateClassnames(parsed.data.items.map((item) => item.classname))
)
return apiError(
"Invalid furniture import request (maximum 500 items)",
400,
);
if (!redis || (await redis.ping().catch(() => null)) !== "PONG")
return apiError(
"Background import queue is temporarily unavailable. Please retry shortly.",
503,
);
if (parsed.data.sourceId && !(await getSource(parsed.data.sourceId)))
return apiError("Source not found", 404);
const store = new ImportJobStore(),
body = parsed.data;
const createdAt = new Date().toISOString();
const unique = [
...new Map(body.items.map((item) => [item.classname, item])).values(),
];
const job = await store.create({
...body,
userId: ctx.session.user.id,
createdAt,
updatedAt: createdAt,
state: "queued",
items: unique.map((item) => ({ ...item, state: "pending" })),
});
after(drainFurnitureImports);
return apiOk({ job });
},
);
@@ -0,0 +1,161 @@
"use client";
import { useCallback, useEffect, useRef, useState } from "react";
import { toast } from "sonner";
import { Button } from "@/components/ui/button";
import { adminFetch } from "@/lib/admin-fetch";
import type { ImportJob, ImportJobItem } from "@/lib/furni/import-job";
export function useFurnitureJobs(onComplete: () => void) {
const [jobs, setJobs] = useState<ImportJob[]>([]),
[busy, setBusy] = useState(false),
[error, setError] = useState("");
const completed = useRef(new Set<string>()),
callback = useRef(onComplete);
callback.current = onComplete;
const request = useRef<{ payload: string; id: string } | null>(null);
const submitting = useRef(false);
const refresh = useCallback(async () => {
try {
const response = await adminFetch("/api/admin/studio/import-jobs");
const data = await response.json();
if (!response.ok)
throw Error(
data.error || `Import history unavailable (${response.status})`,
);
setJobs(data.jobs);
setError("");
let changed = false;
for (const job of data.jobs as ImportJob[])
if (
(job.state === "completed" || job.state === "interrupted") &&
!completed.current.has(job.id)
) {
completed.current.add(job.id);
changed = true;
}
if (changed) callback.current();
} catch (error) {
setError(
error instanceof Error ? error.message : "Cannot load import history",
);
}
}, []);
useEffect(() => {
void refresh();
const timer = setInterval(() => void refresh(), 5000);
return () => clearInterval(timer);
}, [refresh]);
const submit = async (
items: ImportJobItem[],
options: { sourceId?: string; translate: boolean; langs?: string[] },
) => {
if (submitting.current) return false;
submitting.current = true;
setBusy(true);
const body = { items, ...options },
payload = JSON.stringify(body);
if (request.current?.payload !== payload)
request.current = { payload, id: crypto.randomUUID() };
try {
const response = await adminFetch("/api/admin/studio/import-jobs", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ ...body, id: request.current.id }),
});
const data = await response.json();
if (!response.ok)
throw Error(
data.error || `Could not queue import (${response.status})`,
);
request.current = null;
await refresh();
toast.success(
"Import queued. You can leave this page and return to its history.",
);
return true;
} catch (error) {
toast.error(
error instanceof Error ? error.message : "Could not queue import",
);
return false;
} finally {
submitting.current = false;
setBusy(false);
}
};
return { jobs, busy, error, submit, refresh };
}
export function FurnitureJobHistory({
state,
}: {
state: ReturnType<typeof useFurnitureJobs>;
}) {
return (
<details className="shrink-0 border-b border-[var(--admin-border)] bg-[var(--admin-surface)] px-4 py-2 text-sm">
<summary className="cursor-pointer">
Import history ·{" "}
{
state.jobs.filter(
(j) => j.state === "queued" || j.state === "running",
).length
}{" "}
active
</summary>
<div className="max-h-72 space-y-2 overflow-auto py-2">
{state.error && <p role="alert">{state.error}</p>}
{!state.jobs.length && <p>No background imports yet.</p>}
{state.jobs.map((job) => {
const done = job.items.filter((i) => i.state === "done"),
failed = job.items.filter((i) => i.state === "failed");
return (
<details
key={job.id}
className="rounded border border-[var(--admin-border)] p-2"
>
<summary className="cursor-pointer">
{new Date(job.createdAt).toLocaleString()} · {job.state} ·{" "}
{done.length}/{job.items.length} imported · {failed.length}{" "}
failed
</summary>
<ul className="space-y-1 py-2">
{job.items.map((item) => (
<li key={item.classname}>
<strong>{item.classname}</strong> · {item.state}
{item.itemId ? ` · ID ${item.itemId}` : ""}
{item.error && (
<p className="text-[var(--admin-warning)]">
{item.error}
</p>
)}
{[...new Set(item.warnings ?? [])].map((warning) => (
<p
key={warning}
className="text-xs text-[var(--admin-text-muted)]"
>
{warning}
</p>
))}
</li>
))}
</ul>
{failed.length > 0 && job.state !== "running" && (
<Button
variant="outline"
disabled={state.busy}
onClick={() =>
void state.submit(failed, {
sourceId: job.sourceId,
translate: job.translate,
langs: job.langs,
})
}
>
Retry failed items only
</Button>
)}
</details>
);
})}
</div>
</details>
);
}
+64 -3
View File
@@ -1,5 +1,6 @@
"use client";
import { useState } from "react";
import { toast } from "sonner";
import { Button } from "@/components/ui/button";
import {
Dialog,
@@ -9,6 +10,7 @@ import {
DialogHeader,
DialogTitle,
} from "@/components/ui/dialog";
import { adminFetch } from "@/lib/admin-fetch";
import { previewAutoCatalog } from "@/lib/furni/auto-catalog";
import { compareFurniture } from "@/lib/furni/studio-inspection";
import {
@@ -24,6 +26,7 @@ export function ImportReview({
source,
translation,
onCancel,
busy = false,
onConfirm,
}: {
items: FurniItem[];
@@ -31,13 +34,41 @@ export function ImportReview({
source: string;
translation: string;
onCancel: () => void;
busy?: boolean;
onConfirm: (items: FurniItem[]) => void;
}) {
const [attachments, setAttachments] = useState<Record<string, string>>({});
const [uploading, setUploading] = useState<string | null>(null);
async function attach(item: FurniItem, file: File) {
setUploading(item.classname);
try {
const form = new FormData();
form.set("classname", item.classname);
form.set("file", file);
const response = await adminFetch("/api/admin/studio/import-attachment", {
method: "POST",
body: form,
});
const data = await response.json();
if (!response.ok) throw Error(data.error || "Upload failed");
setAttachments((current) => ({
...current,
[item.classname]: data.attachmentId,
}));
toast.success("Matching .nitro attached to this import");
} catch (error) {
toast.error(error instanceof Error ? error.message : "Upload failed");
} finally {
setUploading(null);
}
}
const inspection = useFurnitureInspection(
items.map((item) => item.classname),
);
const assets = useSourceAssetChecks(items, sourceId, inspection);
const unavailable = assets.items.filter((item) => item.state === "missing");
const unavailable = assets.items.filter(
(item) => item.state === "missing" && !attachments[item.classname],
);
const [expanded, setExpanded] = useState(items[0]?.classname ?? "");
const [step, setStep] = useState<"review" | "confirm">("review");
const [page, setPage] = useState(0);
@@ -73,7 +104,7 @@ export function ImportReview({
<Dialog
open
onOpenChange={(open) => {
if (!open) onCancel();
if (!open && !busy && !uploading) onCancel();
}}
>
<DialogContent className="flex max-h-[90dvh] w-[calc(100vw-2rem)] max-w-3xl sm:max-w-3xl flex-col overflow-hidden border-[var(--admin-border)] bg-[var(--admin-surface)] text-[var(--admin-text)]">
@@ -208,6 +239,27 @@ export function ImportReview({
</button>
{expanded === item.classname && (
<div className="border-t border-[var(--admin-border)] p-3">
{local.nitro.exists !== true && (
<label className="mt-3 block text-sm">
{attachments[item.classname]
? "Matching .nitro attached"
: "Attach the original .nitro"}
<input
className="mt-2 block max-w-full"
type="file"
accept=".nitro"
disabled={busy || uploading !== null}
onChange={(event) => {
const file = event.target.files?.[0];
if (file) void attach(item, file);
event.target.value = "";
}}
/>
{uploading === item.classname && (
<span role="status">Validating file…</span>
)}
</label>
)}
<FurnitureComparison item={item} local={local} />
</div>
)}
@@ -278,6 +330,7 @@ export function ImportReview({
<DialogFooter>
<Button
variant="outline"
disabled={busy || uploading !== null}
onClick={
step === "confirm"
? () => {
@@ -295,6 +348,8 @@ export function ImportReview({
color: "var(--admin-accent-foreground)",
}}
disabled={
busy ||
uploading !== null ||
inspection.loading ||
!!inspection.error ||
!complete ||
@@ -304,7 +359,13 @@ export function ImportReview({
if (step === "review") {
setStep("confirm");
setPage(0);
} else onConfirm(ready);
} else
onConfirm(
ready.map((item) => ({
...item,
attachmentId: attachments[item.classname],
})),
);
}}
>
{step === "review"
+19 -223
View File
@@ -60,7 +60,6 @@ import {
type AutoCatalogPreview,
previewAutoCatalog,
} from "@/lib/furni/auto-catalog";
import { readImportResponse } from "@/lib/furni/import-response";
import type { FurniImportSource } from "@/lib/habbo-gamedata-hotel";
import { readSseStream } from "@/lib/sse-client";
import { cn } from "@/lib/utils";
@@ -69,6 +68,7 @@ import { BatchProgress } from "./batch-progress";
import { CatalogRail } from "./catalog-rail";
import { CheckboxDot } from "./checkbox-dot";
import { FurnitureInspector } from "./furniture-inspector";
import { FurnitureJobHistory, useFurnitureJobs } from "./furniture-jobs";
import { ImportReview } from "./import-review";
import {
getFurniImageUrl,
@@ -116,6 +116,10 @@ export function StudioClient({
single: boolean;
cloneAll?: boolean;
} | null>(null);
const jobs = useFurnitureJobs(() => {
void fetchItems(activeSearch, 1, activeSource);
void fetchStats();
});
const [preparingReview, setPreparingReview] = useState(false);
// Catalog tree
const [tree] = useState<TreeNode[]>(initialTree);
@@ -157,7 +161,7 @@ export function StudioClient({
const [detail, setDetail] = useState<FurniItem | null>(null);
// Import state
const [importingId, setImportingId] = useState<string | null>(null);
const importingId = jobs.busy ? (review?.items[0]?.classname ?? null) : null;
const [batchProgress, setBatchProgress] = useState<Map<
string,
BatchItemStatus
@@ -452,146 +456,6 @@ export function StudioClient({
setPreparingReview(false);
}
}
async function cloneAllMissing(classnames: string[]) {
if (!activeSource || !classnames.length) return;
// Initialize progress tracking as a Map for BatchProgress compatibility
const initial = new Map<string, BatchItemStatus>();
for (const cn of classnames) {
initial.set(cn, { classname: cn, status: "pending" });
}
setBatchProgress(initial);
batchDoneRef.current = false;
setBatchDone(false);
setBatchSucceeded(0);
setBatchFailed(0);
setFailedClassnames([]);
setCloneAllCancelling(false);
batchStartTimeRef.current = Date.now();
const abort = new AbortController();
cloneAllAbortRef.current = abort;
const CHUNK = 400;
let okCount = 0;
let failCount = 0;
try {
for (let i = 0; i < classnames.length; i += CHUNK) {
if (abort.signal.aborted) break;
const chunk = classnames.slice(i, i + CHUNK);
const isLast = i + CHUNK >= classnames.length;
const res = await adminFetch("/api/admin/import/clone/batch", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({
sourceId: activeSource,
classnames: chunk,
concurrency: 10,
final: isLast,
}),
signal: abort.signal,
});
if (!res.ok) {
const errMsg =
res.status === 502
? "Furnidata ophalen mislukt – controleer de hotel-URL en netwerk"
: res.status === 400
? "Ongeldige input voor batch-import"
: `Batch chunk mislukt (${res.status})`;
toast.error(errMsg);
// Toon foutmelding uit SSE-stream als die al onderweg is
if (res.body && "getReader" in res.body) {
void res.body.getReader();
toast.error(
"Foutmelding stream niet beschikbaar – de import loopt voort, fouten worden in de stream getoond",
);
} else {
toast.error(
"SSE-stream niet beschikbaar – foutmelding overgeslagen",
);
}
return;
}
await readSseStream(
res.body!,
(evt) => {
if (evt.type === "item_progress") {
const cn = String(evt.classname ?? "");
setBatchProgress((prev) => {
if (!prev) return prev;
const next = new Map(prev);
next.set(cn, {
classname: cn,
status: evt.status as BatchItemStatus["status"],
message: evt.message as string | undefined,
});
return next;
});
if (evt.status === "done" || evt.status === "failed") {
if (evt.status === "failed") {
setFailedClassnames((prev) =>
prev.includes(cn) ? prev : [...prev, cn],
);
}
}
}
if (evt.type === "batch_complete") {
okCount += Number(evt.succeeded ?? 0);
failCount += Number(evt.failed ?? 0);
}
if (evt.type === "error") {
toast.error(
typeof evt.message === "string" ? evt.message : "Import fout",
);
}
},
abort.signal,
);
}
if (!abort.signal.aborted) {
setBatchSucceeded(okCount);
setBatchFailed(failCount);
markBatchDone();
if (okCount > 0) {
toast.success(
`${okCount} items geïmporteerd${
failCount > 0 ? ` (${failCount} mislukt)` : ""
}`,
);
}
if (failCount > 0 && okCount === 0) {
toast.error(`${failCount} items mislukt`);
}
fetchItems(activeSearch, 1, activeSource);
fetchStats();
fetchCloneStats(activeSource);
// Auto-translate after clone all if translation is enabled
if (translate && okCount > 0) {
toast.info("Translating all languages...");
buildAllLanguages();
}
} else {
toast.info("Clone all geannuleerd");
}
} catch (err) {
if ((err as Error)?.name === "AbortError") {
toast.info("Clone all geannuleerd");
} else {
toast.error("Verbindingsfout tijdens clone all");
}
} finally {
cloneAllAbortRef.current = null;
setCloneAllCancelling(false);
setSelected(new Set());
}
}
useEffect(() => {
fetchItems("", 1, activeSource);
@@ -686,76 +550,6 @@ export function StudioClient({
return () => document.removeEventListener("keydown", onKeyDown);
}, [toggleAll, review]);
async function importSingle(item: FurniItem) {
setImportingId(item.classname);
try {
const res = await adminFetch("/api/admin/import/furni", {
method: "POST",
headers: {
"Content-Type": "application/json",
Accept: "text/event-stream",
},
body: JSON.stringify({
repairExisting: true,
id: item.id,
classname: item.classname,
name: item.name,
description: item.description,
type: item.type,
revision: item.revision,
category: item.category,
sourceId: activeSource || undefined,
translate,
langs:
translate && translateLangs.length > 0 ? translateLangs : undefined,
}),
});
const { status, data } = await readImportResponse(res);
if (status >= 200 && status < 300) {
setItems((prev) =>
prev.map((i) =>
i.classname === item.classname
? { ...i, alreadyImported: true, nitroExists: true }
: i,
),
);
if (detail?.classname === item.classname) {
setDetail({ ...item, alreadyImported: true, nitroExists: true });
}
fetchStats();
const parts = [`${item.classname} imported into the catalog`];
if (data.itemId && data.itemId !== item.id)
parts.push(`Local ID: ${data.itemId} (source: ${item.id})`);
if (typeof data.furniDataFixedIds === "number") {
parts.push(
`verified: ${data.offerIdsFixed} offer_id, ${data.furniDataFixedIds} furnidata id`,
);
}
if (
typeof data.furniDataMissing === "number" &&
data.furniDataMissing > 0
) {
parts.push(`${data.furniDataMissing} furnidata missing`);
}
toast.success(parts.join(" · "));
} else {
toast.error(
typeof data.error === "string"
? data.error
: `Import failed for ${item.classname}`,
);
}
} catch (error) {
toast.error(
error instanceof Error
? error.message
: "Import connection failed. Check the furniture status before retrying.",
);
} finally {
setImportingId(null);
}
}
async function importBatch(
requested?: string[],
reviewedItems?: FurniItem[],
@@ -1223,6 +1017,7 @@ export function StudioClient({
return (
<div className="flex h-full min-h-0 min-w-0 flex-col overflow-hidden">
<FurnitureJobHistory state={jobs} />
{/* ── Header ─────────────────────────────────────────── */}
<header className="flex shrink-0 flex-wrap items-center gap-3 border-b border-[var(--admin-border)] bg-[var(--admin-surface)] px-4 py-3">
<div className="flex items-center gap-2.5">
@@ -2419,17 +2214,18 @@ export function StudioClient({
: "Off"
}
onCancel={() => setReview(null)}
onConfirm={(accepted) => {
const single = review.single;
setReview(null);
if (review.cloneAll)
void cloneAllMissing(accepted.map((item) => item.classname));
else if (single && accepted[0]) void importSingle(accepted[0]);
else
void importBatch(
accepted.map((item) => item.classname),
accepted,
);
busy={jobs.busy}
onConfirm={async (accepted) => {
if (
await jobs.submit(accepted, {
sourceId: activeSource || undefined,
translate,
langs: translateLangs.length ? translateLangs : undefined,
})
) {
setReview(null);
setSelected(new Set());
}
}}
/>
)}
@@ -1,4 +1,5 @@
export interface FurniItem {
attachmentId?: string;
id: number;
classname: string;
type: string;
+15
View File
@@ -20,3 +20,18 @@ export const onRequestError: Instrumentation.onRequestError = async (
: null,
});
};
export async function register() {
if (
process.env.NEXT_RUNTIME !== "nodejs" ||
process.env.NEXT_PHASE === "phase-production-build"
)
return;
const { drainFurnitureImports } = await import(
"@/lib/services/furni-job-worker"
);
const timer = setInterval(() => {
void drainFurnitureImports();
}, 30000);
timer.unref();
}
+28
View File
@@ -0,0 +1,28 @@
export interface ImportJobItem {
id: number;
classname: string;
name: string;
description: string;
type: string;
revision: number;
category: string;
attachmentId?: string;
}
export interface ImportJob {
id: string;
userId: number;
createdAt: string;
updatedAt: string;
state: "queued" | "running" | "completed" | "interrupted";
sourceId?: string;
translate: boolean;
langs?: string[];
items: Array<
ImportJobItem & {
state: "pending" | "running" | "done" | "failed";
error?: string;
warnings?: string[];
itemId?: number;
}
>;
}
+33
View File
@@ -0,0 +1,33 @@
import { expect, it } from "vitest";
import { validateFurnitureAttachment } from "./furni-attachment";
import { createNitroBundle, encodePng } from "./swf/nitro-builder";
const png = encodePng(1, 1, Buffer.from([1, 2, 3, 255]));
it("accepts the matching bundle including color variants sharing a library", () => {
const buffer = createNitroBundle({ name: "chair" }, png, "chair");
expect(validateFurnitureAttachment(buffer, "chair*2")).toBe(buffer);
});
it("rejects a similarly named furniture bundle", () => {
expect(() =>
validateFurnitureAttachment(
createNitroBundle({ name: "china_light" }, png, "china_light"),
"nft_china_light",
),
).toThrow("does not belong");
});
it("rejects a renamed JSON entry with mismatching metadata", () => {
expect(() =>
validateFurnitureAttachment(
createNitroBundle({ name: "china_light" }, png, "nft_china_light"),
"nft_china_light",
),
).toThrow("does not belong");
});
it("rejects missing texture data", () => {
expect(() =>
validateFurnitureAttachment(
createNitroBundle({ name: "chair" }, Buffer.from("not png"), "chair"),
"chair",
),
).toThrow("not a PNG");
});
+56
View File
@@ -0,0 +1,56 @@
import { randomUUID } from "node:crypto";
import { promises as fs } from "node:fs";
import path from "node:path";
import { parseNitroBundle } from "@/lib/services/swf/nitro-builder";
import { importRoot, validJobId } from "./furni-job-store";
export function validateFurnitureAttachment(buffer: Buffer, classname: string) {
if (buffer.length > 50 * 1024 * 1024)
throw Error("Maximum file size is 50 MB");
const parsed = parseNitroBundle(buffer);
const base = classname.split("*")[0];
if (
parsed.jsonFileName !== `${base}.json` ||
(parsed.json.name && parsed.json.name !== base)
)
throw Error(`This bundle does not belong to ${classname}`);
if (
!parsed.png
.subarray(0, 8)
.equals(Buffer.from([137, 80, 78, 71, 13, 10, 26, 10]))
)
throw Error("The bundle texture is not a PNG");
return buffer;
}
export async function stageFurnitureAttachment(
buffer: Buffer,
classname: string,
userId: number,
) {
validateFurnitureAttachment(buffer, classname);
const id = randomUUID(),
root = path.join(importRoot(), "attachments");
await fs.mkdir(root, { recursive: true });
await fs.writeFile(path.join(root, `${id}.nitro`), buffer);
await fs.writeFile(
path.join(root, `${id}.json`),
JSON.stringify({ classname, userId }),
);
return id;
}
export async function readFurnitureAttachment(
id: string,
classname: string,
userId: number,
) {
if (!validJobId(id)) throw Error("Invalid attachment");
const root = path.join(importRoot(), "attachments");
const meta = JSON.parse(
await fs.readFile(path.join(root, `${id}.json`), "utf8"),
);
if (meta.userId !== userId || meta.classname !== classname)
throw Error("Attachment does not belong to this furniture import");
return validateFurnitureAttachment(
await fs.readFile(path.join(root, `${id}.nitro`)),
classname,
);
}
+28 -6
View File
@@ -37,6 +37,7 @@ import {
} from "@/lib/services/swf-to-nitro";
import { getRuntimePath } from "@/lib/utils/runtime-path";
import type { ImportSingleResult } from "@/types/furni";
import { extractFurniIconPng } from "./clone-icon";
import { reserveFurnitureId } from "./furniture-id-reservation";
// Re-export the download helpers now owned by the shared import core.
@@ -512,6 +513,8 @@ export async function importSingleFurni(params: {
skipFurniDataWrite?: boolean;
updateExisting?: boolean;
repairExisting?: boolean;
/** Validated matching bundle supplied by a queued import. */
providedNitro?: Buffer;
onProgress?: (status: string) => void;
/** Per-source SWF download base URL (e.g. "https://virtualc.nl/dcr"). */
sourceSwfBaseUrl?: string;
@@ -645,6 +648,10 @@ export async function importSingleFurni(params: {
const mirrorIconDirs = assetTargets.mirrorDirs.map((d) => d.iconDir);
const mirrorNitroDirs = assetTargets.mirrorDirs.map((d) => d.nitroDir);
if (params.providedNitro && !existsSync(nitroPath)) {
await fs.writeFile(nitroPath, params.providedNitro, { flag: "wx" });
}
// ── Download assets in parallel (with retry + validation) ────────
// Icon strategy matches Nitro's URL template:
// `${hof.furni.url}/icons/%libname%%param%_icon.png`
@@ -700,11 +707,13 @@ export async function importSingleFurni(params: {
: tryDownloadCandidates(iconUrls, iconPath, "png", (detail) => {
iconFailure = detail;
}),
preservingExisting && existsSync(swfPath)
? Promise.resolve(true)
: tryDownloadCandidates(swfUrls, swfPath, "swf", (detail) => {
swfFailure = detail;
}),
params.providedNitro
? Promise.resolve(false)
: preservingExisting && existsSync(swfPath)
? Promise.resolve(true)
: tryDownloadCandidates(swfUrls, swfPath, "swf", (detail) => {
swfFailure = detail;
}),
]);
let iconOk = iconOkResult;
@@ -744,7 +753,7 @@ export async function importSingleFurni(params: {
);
}
if (!swfOk && !nitroDownloadOk)
if (!swfOk && !nitroDownloadOk && !params.providedNitro)
warnings.push(`SWF download failed: ${swfFailure || "unknown cause"}`);
// ── Convert SWF to Nitro ──────────────────────────────────────────
@@ -781,6 +790,19 @@ export async function importSingleFurni(params: {
}
}
if (!iconOk && existsSync(nitroPath)) {
try {
const icon = extractFurniIconPng(await fs.readFile(nitroPath));
if (icon) {
await fs.writeFile(iconPath, icon);
iconOk = true;
warnings.push("Icon extracted from Nitro bundle");
}
} catch {
warnings.push("Could not extract an icon from the Nitro bundle");
}
}
// ── Rollback: if conversion/download failed AND no .nitro exists, remove DB record ──
// Also validate that the .nitro file is non-trivial (≥ 128 bytes — a bare
// Nitro bundle header is larger than that, anything smaller is corrupt/empty).
+61
View File
@@ -0,0 +1,61 @@
import { randomUUID } from "node:crypto";
import { promises as fs } from "node:fs";
import { tmpdir } from "node:os";
import path from "node:path";
import { afterEach, expect, it } from "vitest";
import type { ImportJob } from "@/lib/furni/import-job";
import { ImportJobStore } from "./furni-job-store";
const roots: string[] = [];
afterEach(async () => {
for (const root of roots.splice(0))
await fs.rm(root, { recursive: true, force: true });
});
const job = (): ImportJob => ({
id: randomUUID(),
userId: 4,
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
state: "queued",
translate: false,
items: [],
});
async function store() {
const root = await fs.mkdtemp(path.join(tmpdir(), "furni-jobs-test-"));
roots.push(root);
return new ImportJobStore(root);
}
it("persists history across store instances", async () => {
const s = await store(),
j = job();
await s.create(j);
expect(await new ImportJobStore(s.root).read(j.id)).toEqual(j);
});
it("deduplicates concurrent submissions of the same request", async () => {
const s = await store(),
j = job();
await Promise.all([s.create(j), s.create(j)]);
expect(await s.list()).toHaveLength(1);
});
it("does not overwrite a completed job on a repeated request", async () => {
const s = await store(),
j = job();
await s.create(j);
await s.save({ ...j, state: "completed" });
expect((await s.create(j)).state).toBe("completed");
});
it("rejects another user attempting to reuse an ID", async () => {
const s = await store(),
j = job();
await s.create(j);
await expect(s.create({ ...j, userId: 8 })).rejects.toThrow(
"Request ID already used",
);
});
it("rejects unsafe job paths", async () => {
const s = await store();
await expect(s.create({ ...job(), id: "../escape" })).rejects.toThrow(
"Invalid import ID",
);
await expect(s.read("../escape")).rejects.toThrow("Invalid import ID");
});
+58
View File
@@ -0,0 +1,58 @@
import { randomUUID } from "node:crypto";
import { promises as fs } from "node:fs";
import path from "node:path";
import type { ImportJob } from "@/lib/furni/import-job";
export const importRoot = () =>
path.join(process.cwd(), "storage", "furniture-imports");
export const validJobId = (value: unknown): value is string =>
typeof value === "string" &&
/^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test(
value,
);
export class ImportJobStore {
constructor(readonly root = importRoot()) {}
async save(job: ImportJob) {
await fs.mkdir(this.root, { recursive: true });
const target = path.join(this.root, `${job.id}.json`),
tmp = `${target}.${randomUUID()}.tmp`;
job.updatedAt = new Date().toISOString();
await fs.writeFile(tmp, JSON.stringify(job));
await fs.rename(tmp, target);
}
async create(job: ImportJob) {
if (!validJobId(job.id)) throw Error("Invalid import ID");
await fs.mkdir(this.root, { recursive: true });
const tmp = path.join(this.root, `${randomUUID()}.tmp`);
await fs.writeFile(tmp, JSON.stringify(job));
try {
await fs.link(tmp, path.join(this.root, `${job.id}.json`));
return job;
} catch (error) {
if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error;
const existing = await this.read(job.id);
if (existing.userId !== job.userId)
throw Error("Request ID already used");
return existing;
} finally {
await fs.unlink(tmp);
}
}
async read(id: string): Promise<ImportJob> {
if (!validJobId(id)) throw Error("Invalid import ID");
return JSON.parse(
await fs.readFile(path.join(this.root, `${id}.json`), "utf8"),
);
}
async list(): Promise<ImportJob[]> {
const files = await fs.readdir(this.root).catch((error) => {
if (error.code === "ENOENT") return [];
throw error;
});
const jobs = await Promise.all(
files
.filter((f) => f.endsWith(".json") && validJobId(f.slice(0, -5)))
.map((f) => this.read(f.slice(0, -5))),
);
return jobs.sort((a, b) => a.createdAt.localeCompare(b.createdAt));
}
}
+114
View File
@@ -0,0 +1,114 @@
import { beforeEach, expect, it, vi } from "vitest";
import type { ImportJob } from "@/lib/furni/import-job";
const mocks = vi.hoisted(() => ({
set: vi.fn(),
eval: vi.fn(),
list: vi.fn(),
save: vi.fn(),
import: vi.fn(),
attachment: vi.fn(),
export: vi.fn(),
translate: vi.fn(),
}));
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("./catalog-git-queue", () => ({ withCatalogExport: mocks.export }));
vi.mock("./clone-sources", () => ({ getSource: async () => null }));
vi.mock("./furni-job-store", () => ({
ImportJobStore: class {
list = mocks.list;
save = mocks.save;
},
}));
vi.mock("./furni-import", () => ({
ensureDirectories: async () => {},
importSingleFurni: mocks.import,
}));
vi.mock("./furni-data-i18n", () => ({
patchLocalizedFurniDataEntries: mocks.translate,
}));
vi.mock("./furni-attachment", () => ({
readFurnitureAttachment: mocks.attachment,
}));
vi.mock("./rcon", () => ({
rcon: { updateCatalog: async () => {}, updateItems: async () => {} },
}));
vi.mock("./catalog-git-export", () => ({ runCatalogExport: vi.fn() }));
import { drainFurnitureImports } from "./furni-job-worker";
let job: ImportJob;
beforeEach(() => {
vi.clearAllMocks();
mocks.set.mockResolvedValue("OK");
mocks.eval.mockResolvedValue(1);
mocks.export.mockImplementation((fn) => fn());
mocks.import.mockResolvedValue({ ok: true, itemId: 900, warnings: [] });
mocks.save.mockResolvedValue(undefined);
job = {
id: "job",
userId: 5,
createdAt: "now",
updatedAt: "now",
state: "queued",
translate: false,
items: [
{
id: 1,
classname: "chair",
name: "Chair",
description: "",
type: "flooritem",
revision: 1,
category: "other",
state: "pending",
},
],
};
mocks.list.mockResolvedValue([job]);
});
it("imports with repair enabled and saves the allocated local ID", async () => {
await drainFurnitureImports();
expect(mocks.import).toHaveBeenCalledWith(
expect.objectContaining({ classname: "chair", repairExisting: true }),
);
expect(job.state).toBe("completed");
expect(job.items[0]).toMatchObject({ state: "done", itemId: 900 });
expect(mocks.export).toHaveBeenCalledOnce();
});
it("does not run a second worker while another holds the lease", async () => {
mocks.set.mockResolvedValue(null);
await drainFurnitureImports();
expect(mocks.import).not.toHaveBeenCalled();
});
it("marks abandoned work interrupted without repeating an uncertain import", async () => {
job.state = "running";
job.items[0].state = "running";
await drainFurnitureImports();
expect(job.state).toBe("interrupted");
expect(job.items[0].state).toBe("failed");
expect(mocks.import).not.toHaveBeenCalled();
});
it("keeps errors for failed items and continues the batch", async () => {
job.items.push({ ...job.items[0], classname: "table" });
mocks.import.mockResolvedValueOnce({
ok: false,
error: "Missing asset",
warnings: [],
});
await drainFurnitureImports();
expect(job.items.map((i) => i.state)).toEqual(["failed", "done"]);
expect(job.items[0].error).toBe("Missing asset");
});
it("loads attachments using the job owner and exact classname", async () => {
job.items[0].attachmentId = "upload";
mocks.attachment.mockResolvedValue(Buffer.from("verified"));
await drainFurnitureImports();
expect(mocks.attachment).toHaveBeenCalledWith("upload", "chair", 5);
expect(mocks.import).toHaveBeenCalledWith(
expect.objectContaining({ providedNitro: Buffer.from("verified") }),
);
});
+142
View File
@@ -0,0 +1,142 @@
import { randomUUID } from "node:crypto";
import { redis } from "@/lib/redis";
import { logServerError } from "@/lib/server-log";
import { logAudit } from "./audit";
import { runCatalogExport } from "./catalog-git-export";
import { withCatalogExport } from "./catalog-git-queue";
import { getSource } from "./clone-sources";
import { readFurnitureAttachment } from "./furni-attachment";
import { patchLocalizedFurniDataEntries } from "./furni-data-i18n";
import { ensureDirectories, importSingleFurni } from "./furni-import";
import { ImportJobStore } from "./furni-job-store";
import { rcon } from "./rcon";
const LOCK = "furniture-import-worker:v1";
let running: Promise<void> | undefined;
export function drainFurnitureImports(): Promise<void> {
if (running) return running;
running = drain()
.catch((error) => logServerError("furni.worker_failed", error))
.finally(() => {
running = undefined;
});
return running;
}
async function drain() {
if (!redis) return;
const token = randomUUID();
if ((await redis.set(LOCK, token, "EX", 600, "NX")) !== "OK") return;
let lease = true;
const timer = setInterval(() => {
void redis
?.eval(
"if redis.call('get',KEYS[1]) == ARGV[1] then return redis.call('expire',KEYS[1],600) else return 0 end",
1,
LOCK,
token,
)
.then((result) => {
if (result !== 1) lease = false;
})
.catch(() => {
lease = false;
});
}, 30000);
try {
const store = new ImportJobStore();
for (const job of await store.list()) {
if (!lease) break;
if (job.state === "running") {
job.state = "interrupted";
for (const item of job.items)
if (item.state === "running" || item.state === "pending") {
item.state = "failed";
item.error =
"Server restarted during import. Review local data and retry to complete missing parts.";
}
await store.save(job);
continue;
}
if (job.state !== "queued") continue;
job.state = "running";
await store.save(job);
await withCatalogExport(async () => {
await ensureDirectories();
const source = job.sourceId ? await getSource(job.sourceId) : null;
for (const item of job.items) {
if (!lease) throw Error("Import worker lost its lease");
item.state = "running";
await store.save(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 result = await importSingleFurni({
...item,
repairExisting: true,
providedNitro,
sourceSwfBaseUrl: source?.sourceSwfBaseUrl,
nitroBaseUrl: source?.nitroBaseUrl,
iconBaseUrl: source?.iconBaseUrl,
});
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,
jobId: job.id,
repairExisting: true,
},
});
if (job.translate)
try {
await patchLocalizedFurniDataEntries([item], true, job.langs);
} catch {
item.warnings.push(
"Translation failed; furniture imported successfully",
);
}
} catch (error) {
item.state = "failed";
item.error =
error instanceof Error ? error.message : "Import failed";
}
await store.save(job);
}
try {
await rcon.updateCatalog();
await rcon.updateItems();
} catch {
for (const item of job.items)
if (item.state === "done") {
item.warnings ??= [];
item.warnings.push("Game cache refresh failed");
}
}
});
job.state = "completed";
await store.save(job);
}
await runCatalogExport();
} finally {
clearInterval(timer);
await redis.eval(
"if redis.call('get',KEYS[1]) == ARGV[1] then return redis.call('del',KEYS[1]) else return 0 end",
1,
LOCK,
token,
);
}
}