feat(catalog): queue furniture imports with history and matching file attachments
This commit is contained in:
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>
|
||||
);
|
||||
}
|
||||
@@ -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"
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
>;
|
||||
}
|
||||
@@ -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");
|
||||
});
|
||||
@@ -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,
|
||||
);
|
||||
}
|
||||
@@ -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).
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
@@ -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));
|
||||
}
|
||||
}
|
||||
@@ -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") }),
|
||||
);
|
||||
});
|
||||
@@ -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,
|
||||
);
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user