Files
Epicnabbo-Catalogus-Updated…/src/lib/services/furni-job-worker.ts
T

148 lines
4.3 KiB
TypeScript

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