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 { repairFurniture } from "./furniture-repair"; import { rcon } from "./rcon"; const LOCK = "furniture-import-worker:v1"; let running: Promise | undefined; export function drainFurnitureImports(): Promise { 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 (job.mode === "repair" ? repairFurniture : 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 && job.mode !== "repair") 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, ); } }