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

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