feat: export Catalog Studio assets and SQL to catalog repository
This commit is contained in:
1 parent
bfbf244141
commit
9ec9ab31ad
21 files changed
+1811
-514
No files matched your search
+41
-2
@@ -1,9 +1,15 @@
|
||||
import { type NextRequest, NextResponse } from "next/server";
|
||||
import { after, type NextRequest, NextResponse } from "next/server";
|
||||
import { logAuthorizationEvent } from "@/lib/admin/authorization-events";
|
||||
import { validateCsrfToken } from "@/lib/foundation/security";
|
||||
import { canAccess, getApiAdminContext } from "@/lib/permissions";
|
||||
import { logServerError } from "@/lib/server-log";
|
||||
|
||||
import {
|
||||
catalogExportEnabled,
|
||||
catalogExportQueue,
|
||||
isCatalogMutation,
|
||||
} from "@/lib/services/catalog-git-queue";
|
||||
|
||||
const MUTATING_METHODS = new Set(["POST", "PUT", "PATCH", "DELETE"]);
|
||||
const MAX_BODY_BYTES = 10 * 1024 * 1024; // 10 MB
|
||||
|
||||
@@ -80,9 +86,42 @@ export function withAdmin(
|
||||
{ status: 403 },
|
||||
);
|
||||
}
|
||||
let finishExport: (() => Promise<void>) | undefined;
|
||||
try {
|
||||
return await handler(request, context, routeContext);
|
||||
if (
|
||||
catalogExportEnabled() &&
|
||||
isCatalogMutation(request.method, request.nextUrl.pathname)
|
||||
) {
|
||||
finishExport = await catalogExportQueue().begin();
|
||||
}
|
||||
const response = await handler(request, context, routeContext);
|
||||
if (
|
||||
finishExport &&
|
||||
response.body &&
|
||||
response.headers.get("content-type")?.includes("text/event-stream")
|
||||
) {
|
||||
const [client, completion] = response.body.tee();
|
||||
const finish = finishExport;
|
||||
after(async () => {
|
||||
const reader = completion.getReader();
|
||||
try {
|
||||
while (!(await reader.read()).done) {
|
||||
/* Wait for all import results. */
|
||||
}
|
||||
} finally {
|
||||
reader.releaseLock();
|
||||
await finish();
|
||||
}
|
||||
});
|
||||
return new Response(client, {
|
||||
status: response.status,
|
||||
headers: response.headers,
|
||||
});
|
||||
}
|
||||
await finishExport?.();
|
||||
return response;
|
||||
} catch (error) {
|
||||
await finishExport?.().catch(() => undefined);
|
||||
logServerError("admin.api_failed", error, {
|
||||
path: request.nextUrl.pathname,
|
||||
userId: context.session.user.id,
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
import { mkdtemp } from "node:fs/promises";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
import { NextRequest } from "next/server";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
after: vi.fn(),
|
||||
auth: vi.fn(),
|
||||
csrf: vi.fn(),
|
||||
}));
|
||||
vi.mock("next/server", async (original) => ({
|
||||
...(await original<typeof import("next/server")>()),
|
||||
after: mocks.after,
|
||||
}));
|
||||
vi.mock("@/lib/admin/authorization-events", () => ({
|
||||
logAuthorizationEvent: vi.fn(),
|
||||
}));
|
||||
vi.mock("@/lib/foundation/security", () => ({ validateCsrfToken: mocks.csrf }));
|
||||
vi.mock("@/lib/permissions", () => ({
|
||||
getApiAdminContext: mocks.auth,
|
||||
canAccess: () => true,
|
||||
}));
|
||||
vi.mock("@/lib/server-log", () => ({ logServerError: vi.fn() }));
|
||||
|
||||
import { withAdmin } from "./api-handler";
|
||||
import { catalogExportQueue } from "./services/catalog-git-queue";
|
||||
|
||||
describe("catalog API completion tracking", () => {
|
||||
beforeEach(async () => {
|
||||
vi.stubEnv("CATALOG_GIT_CHECKOUT", "/catalog");
|
||||
vi.stubEnv(
|
||||
"CATALOG_GIT_STATE_DIR",
|
||||
await mkdtemp(path.join(os.tmpdir(), "catalog-api-")),
|
||||
);
|
||||
mocks.auth.mockResolvedValue({
|
||||
session: { user: { id: 1, rank: 7 } },
|
||||
permissions: [],
|
||||
});
|
||||
mocks.csrf.mockResolvedValue(true);
|
||||
mocks.after.mockReset();
|
||||
});
|
||||
afterEach(() => vi.unstubAllEnvs());
|
||||
it("waits for the last SSE event before making changes publishable", async () => {
|
||||
let end: (() => void) | undefined;
|
||||
const response = new Response(
|
||||
new ReadableStream({
|
||||
start(controller) {
|
||||
controller.enqueue(
|
||||
new TextEncoder().encode('data: {"type":"progress"}\n\n'),
|
||||
);
|
||||
end = () => controller.close();
|
||||
},
|
||||
}),
|
||||
{ headers: { "content-type": "text/event-stream" } },
|
||||
);
|
||||
const route = withAdmin({}, async () => response);
|
||||
const returned = await route(
|
||||
new NextRequest("http://localhost/api/admin/import/furni/batch", {
|
||||
method: "POST",
|
||||
}),
|
||||
);
|
||||
expect(await catalogExportQueue().batch()).toBeNull();
|
||||
const complete = mocks.after.mock.calls[0][0]();
|
||||
end?.();
|
||||
await returned.text();
|
||||
await complete;
|
||||
expect(await catalogExportQueue().batch()).toHaveLength(1);
|
||||
});
|
||||
it("does not queue unauthenticated requests or reads", async () => {
|
||||
const handler = vi.fn(async () => new Response("{}"));
|
||||
const route = withAdmin({}, handler);
|
||||
mocks.auth.mockResolvedValueOnce(null);
|
||||
expect(
|
||||
(
|
||||
await route(
|
||||
new NextRequest("http://localhost/api/admin/import/furni", {
|
||||
method: "POST",
|
||||
}),
|
||||
)
|
||||
).status,
|
||||
).toBe(401);
|
||||
expect(handler).not.toHaveBeenCalled();
|
||||
await route(new NextRequest("http://localhost/api/admin/import/furni"));
|
||||
expect(await catalogExportQueue().batch()).toHaveLength(0);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,103 @@
|
||||
import { execFileSync } from "node:child_process";
|
||||
import { mkdir, mkdtemp, readFile, writeFile } from "node:fs/promises";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import {
|
||||
CatalogExportQueue,
|
||||
publishCatalogFiles,
|
||||
recoverCatalogQueue,
|
||||
sqlValue,
|
||||
} from "./catalog-git-core";
|
||||
|
||||
describe("catalog export", () => {
|
||||
it("recovers queue entries owned by a terminated local process", async () => {
|
||||
const root = await mkdtemp(path.join(os.tmpdir(), "catalog-recovery-"));
|
||||
const deadPid = Number(
|
||||
execFileSync(
|
||||
process.execPath,
|
||||
["-e", "process.stdout.write(String(process.pid))"],
|
||||
{ encoding: "utf8" },
|
||||
),
|
||||
);
|
||||
const owner = JSON.stringify({ pid: deadPid, host: os.hostname() });
|
||||
await writeFile(path.join(root, "12345678.active"), owner);
|
||||
await writeFile(path.join(root, "worker.lock"), owner);
|
||||
await recoverCatalogQueue(root);
|
||||
const q = new CatalogExportQueue(root);
|
||||
expect(await q.batch()).toEqual(["12345678.pending"]);
|
||||
expect(await q.entries()).not.toContain("worker.lock");
|
||||
});
|
||||
|
||||
it("keeps new requests pending when an older snapshot completes", async () => {
|
||||
const root = await mkdtemp(path.join(os.tmpdir(), "catalog-queue-"));
|
||||
const q = new CatalogExportQueue(root);
|
||||
const finish = await q.begin();
|
||||
expect(await q.batch()).toBeNull();
|
||||
await finish();
|
||||
const batch = await q.batch();
|
||||
expect(batch).toHaveLength(1);
|
||||
await q.request();
|
||||
if (!batch) throw new Error("Expected pending batch");
|
||||
await q.complete(batch);
|
||||
expect(await q.batch()).toHaveLength(1);
|
||||
});
|
||||
it("encodes SQL strings without quote or backslash ambiguity", () => {
|
||||
expect(sqlValue("Valentine's \\ Day\n")).toBe(
|
||||
"CONVERT(X'56616c656e74696e652773205c204461790a' USING utf8mb4)",
|
||||
);
|
||||
expect(sqlValue(null)).toBe("NULL");
|
||||
expect(sqlValue(42)).toBe("42");
|
||||
expect(() => sqlValue(Number.NaN)).toThrow();
|
||||
});
|
||||
it("publishes to a bare remote, retries a failed push and avoids empty commits", async () => {
|
||||
const root = await mkdtemp(path.join(os.tmpdir(), "catalog-git-"));
|
||||
const remote = path.join(root, "remote.git");
|
||||
const checkout = path.join(root, "checkout");
|
||||
const git = (cwd: string, ...args: string[]) =>
|
||||
execFileSync("git", args, { cwd, encoding: "utf8" }).trim();
|
||||
git(root, "init", "--bare", remote);
|
||||
git(root, "clone", remote, checkout);
|
||||
git(checkout, "checkout", "-b", "Beta-3");
|
||||
git(checkout, "config", "user.name", "Catalog test");
|
||||
git(checkout, "config", "user.email", "[email protected]");
|
||||
await writeFile(path.join(checkout, "readme.md"), "preserved");
|
||||
git(checkout, "add", "readme.md");
|
||||
git(checkout, "commit", "-m", "initial");
|
||||
git(checkout, "push", "-u", "origin", "Beta-3");
|
||||
const source = path.join(root, "chair.nitro");
|
||||
await writeFile(source, "bundle");
|
||||
const files = [
|
||||
{ source, target: "Gamedata/bundled/furniture/chair.nitro" },
|
||||
];
|
||||
const options = { checkout, remote, branch: "Beta-3", files };
|
||||
const first = await publishCatalogFiles(options);
|
||||
expect(await publishCatalogFiles(options)).toBe(first);
|
||||
await writeFile(source, "updated");
|
||||
git(
|
||||
checkout,
|
||||
"remote",
|
||||
"set-url",
|
||||
"--push",
|
||||
"origin",
|
||||
path.join(root, "missing.git"),
|
||||
);
|
||||
await expect(publishCatalogFiles(options)).rejects.toThrow();
|
||||
git(checkout, "remote", "set-url", "--push", "origin", remote);
|
||||
const retried = await publishCatalogFiles(options);
|
||||
expect(retried).not.toBe(first);
|
||||
expect(git(remote, "rev-parse", "Beta-3")).toBe(retried);
|
||||
expect(await readFile(path.join(checkout, "readme.md"), "utf8")).toBe(
|
||||
"preserved",
|
||||
);
|
||||
await expect(
|
||||
publishCatalogFiles({
|
||||
...options,
|
||||
files: [{ source, target: "../escape" }],
|
||||
}),
|
||||
).rejects.toThrow();
|
||||
await mkdir(path.join(checkout, "unrelated"));
|
||||
await writeFile(path.join(checkout, "unrelated", "local.txt"), "local");
|
||||
await expect(publishCatalogFiles(options)).rejects.toThrow(/clean/);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,230 @@
|
||||
import { execFile } from "node:child_process";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { promises as fs } from "node:fs";
|
||||
import { hostname } from "node:os";
|
||||
import path from "node:path";
|
||||
import { promisify } from "node:util";
|
||||
|
||||
const exec = promisify(execFile);
|
||||
export const CATALOG_REMOTE =
|
||||
"https://gitlab.epicnabbo.nl/remco/Epicnabbo-Catalogus-Updated-Daily.git";
|
||||
export const CATALOG_BRANCH = "Beta-3";
|
||||
export const CATALOG_SQL_ROOT = "catalogue version 2 ( Final (Dev)/sqls";
|
||||
export const CATALOG_LANGUAGE_ROOT =
|
||||
"catalogue version 2 ( Final (Dev)/langs furnidata";
|
||||
|
||||
export class CatalogExportQueue {
|
||||
constructor(readonly root: string) {}
|
||||
async entries() {
|
||||
await fs.mkdir(this.root, { recursive: true });
|
||||
return (await fs.readdir(this.root)).sort();
|
||||
}
|
||||
async begin() {
|
||||
await this.entries();
|
||||
const id = randomUUID();
|
||||
await fs.writeFile(
|
||||
path.join(this.root, `${id}.active`),
|
||||
JSON.stringify({
|
||||
pid: process.pid,
|
||||
host: hostname(),
|
||||
startedAt: new Date().toISOString(),
|
||||
}),
|
||||
{ flag: "wx" },
|
||||
);
|
||||
let finished = false;
|
||||
return async () => {
|
||||
if (finished) return;
|
||||
await fs.rename(
|
||||
path.join(this.root, `${id}.active`),
|
||||
path.join(this.root, `${id}.pending`),
|
||||
);
|
||||
finished = true;
|
||||
};
|
||||
}
|
||||
async request() {
|
||||
await (await this.begin())();
|
||||
}
|
||||
async batch(): Promise<string[] | null> {
|
||||
const entries = await this.entries();
|
||||
if (entries.some((f) => f.endsWith(".active"))) return null;
|
||||
return entries.filter((f) => f.endsWith(".pending"));
|
||||
}
|
||||
async complete(batch: string[]) {
|
||||
for (const file of batch) {
|
||||
if (!/^[a-f0-9-]+\.pending$/.test(file))
|
||||
throw new Error("Invalid queue entry");
|
||||
await fs.rm(path.join(this.root, file), { force: true });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export function sqlValue(value: unknown): string {
|
||||
if (value === null || value === undefined) return "NULL";
|
||||
if (typeof value === "number") {
|
||||
if (!Number.isFinite(value)) throw new Error("Invalid SQL number");
|
||||
return String(value);
|
||||
}
|
||||
if (typeof value === "bigint") return String(value);
|
||||
if (typeof value === "boolean") return value ? "1" : "0";
|
||||
if (Buffer.isBuffer(value)) return `X'${value.toString("hex")}'`;
|
||||
const text =
|
||||
value instanceof Date
|
||||
? value.toISOString().slice(0, 23).replace("T", " ")
|
||||
: String(value);
|
||||
return text.length
|
||||
? `CONVERT(X'${Buffer.from(text, "utf8").toString("hex")}' USING utf8mb4)`
|
||||
: "''";
|
||||
}
|
||||
|
||||
export interface CatalogFile {
|
||||
source: string;
|
||||
target: string;
|
||||
}
|
||||
|
||||
export function catalogTarget(root: string, target: string) {
|
||||
if (
|
||||
target.includes("\\") ||
|
||||
target.split("/").some((s) => s === ".." || s === ".git") ||
|
||||
path.isAbsolute(target)
|
||||
) {
|
||||
throw new Error("Invalid catalog target");
|
||||
}
|
||||
const result = path.resolve(root, target);
|
||||
if (!result.startsWith(`${path.resolve(root)}${path.sep}`))
|
||||
throw new Error("Invalid catalog target");
|
||||
if (
|
||||
!target.startsWith("Gamedata/") &&
|
||||
!target.startsWith(`${CATALOG_SQL_ROOT}/`) &&
|
||||
!target.startsWith(`${CATALOG_LANGUAGE_ROOT}/`)
|
||||
) {
|
||||
throw new Error("Unmanaged catalog target");
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
export async function publishCatalogFiles(options: {
|
||||
checkout: string;
|
||||
remote: string;
|
||||
branch: string;
|
||||
files: CatalogFile[];
|
||||
}) {
|
||||
const { checkout, remote, branch, files } = options;
|
||||
const git = async (...args: string[]) => {
|
||||
const { stdout } = await exec("git", args, {
|
||||
cwd: checkout,
|
||||
timeout: 120_000,
|
||||
maxBuffer: 16 * 1024 * 1024,
|
||||
env: { ...process.env, GIT_TERMINAL_PROMPT: "0" },
|
||||
});
|
||||
return stdout.trim();
|
||||
};
|
||||
for (const file of files) catalogTarget(checkout, file.target);
|
||||
if ((await git("remote", "get-url", "origin")) !== remote)
|
||||
throw new Error("Unexpected catalog origin");
|
||||
if ((await git("branch", "--show-current")) !== branch)
|
||||
throw new Error("Unexpected catalog branch");
|
||||
if (await git("status", "--porcelain"))
|
||||
throw new Error("Catalog checkout must be clean");
|
||||
try {
|
||||
await git("pull", "--rebase", "origin", branch);
|
||||
} catch (error) {
|
||||
await git("rebase", "--abort").catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
const created: string[] = [];
|
||||
let committed = false;
|
||||
try {
|
||||
for (const file of files) {
|
||||
const target = catalogTarget(checkout, file.target);
|
||||
// Reject symlink ancestors, including the target itself.
|
||||
let ancestor = target;
|
||||
while (ancestor !== path.resolve(checkout)) {
|
||||
const stat = await fs
|
||||
.lstat(ancestor)
|
||||
.catch((error: NodeJS.ErrnoException) => {
|
||||
if (error.code === "ENOENT") return null;
|
||||
throw error;
|
||||
});
|
||||
if (stat?.isSymbolicLink())
|
||||
throw new Error("Symlink in catalog target");
|
||||
ancestor = path.dirname(ancestor);
|
||||
}
|
||||
const exists = await fs.stat(target).then(
|
||||
() => true,
|
||||
() => false,
|
||||
);
|
||||
if (!exists) created.push(target);
|
||||
await fs.mkdir(path.dirname(target), { recursive: true });
|
||||
await fs.copyFile(file.source, target);
|
||||
}
|
||||
// Only explicitly exported paths are eligible for staging.
|
||||
for (let i = 0; i < files.length; i += 50) {
|
||||
await git(
|
||||
"--literal-pathspecs",
|
||||
"add",
|
||||
"--",
|
||||
...files.slice(i, i + 50).map((f) => f.target),
|
||||
);
|
||||
}
|
||||
if (await git("diff", "--cached", "--name-only")) {
|
||||
await git(
|
||||
"-c",
|
||||
"user.name=Catalog Studio",
|
||||
"-c",
|
||||
"[email protected]",
|
||||
"commit",
|
||||
"-m",
|
||||
"Update catalog assets and SQL from Catalog Studio",
|
||||
);
|
||||
}
|
||||
committed = true;
|
||||
// Also retries commits retained after a previous failed push.
|
||||
await git("push", "origin", `HEAD:refs/heads/${branch}`);
|
||||
return await git("rev-parse", "HEAD");
|
||||
} catch (error) {
|
||||
if (!committed) {
|
||||
// The checkout was clean at entry; restore only this attempt's changes.
|
||||
await git("restore", "--staged", "--worktree", ".");
|
||||
for (const file of created) await fs.rm(file, { force: true });
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export async function recoverCatalogQueue(root: string) {
|
||||
for (const name of await fs.readdir(root)) {
|
||||
if (!name.endsWith(".active") && name !== "worker.lock") continue;
|
||||
const file = path.join(root, name);
|
||||
let owner: { pid: number; host: string };
|
||||
try {
|
||||
owner = JSON.parse(await fs.readFile(file, "utf8"));
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
if (
|
||||
owner.host !== hostname() ||
|
||||
!Number.isInteger(owner.pid) ||
|
||||
owner.pid <= 0
|
||||
)
|
||||
continue;
|
||||
try {
|
||||
process.kill(owner.pid, 0);
|
||||
continue;
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== "ESRCH") continue;
|
||||
}
|
||||
// Recheck ownership before recovering a dead process's entry.
|
||||
if (
|
||||
(await fs.readFile(file, "utf8").catch(() => "")) !==
|
||||
JSON.stringify(owner)
|
||||
)
|
||||
continue;
|
||||
if (name === "worker.lock") await fs.rm(file, { force: true });
|
||||
else
|
||||
await fs
|
||||
.rename(file, path.join(root, name.replace(/\.active$/, ".pending")))
|
||||
.catch((error: NodeJS.ErrnoException) => {
|
||||
if (error.code !== "ENOENT") throw error;
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,57 @@
|
||||
import { mkdtemp, readFile } from "node:fs/promises";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
const mocks = vi.hoisted(() => ({ snapshot: vi.fn(), publish: vi.fn() }));
|
||||
vi.mock("./catalog-git-snapshot", () => ({
|
||||
createCatalogSnapshot: mocks.snapshot,
|
||||
}));
|
||||
vi.mock("./catalog-git-core", async (original) => ({
|
||||
...(await original<typeof import("./catalog-git-core")>()),
|
||||
publishCatalogFiles: mocks.publish,
|
||||
}));
|
||||
|
||||
import { runCatalogExport } from "./catalog-git-export";
|
||||
import { catalogExportQueue, withCatalogExport } from "./catalog-git-queue";
|
||||
|
||||
describe("export worker", () => {
|
||||
beforeEach(async () => {
|
||||
vi.stubEnv("CATALOG_GIT_CHECKOUT", "/configured/checkout");
|
||||
vi.stubEnv(
|
||||
"CATALOG_GIT_STATE_DIR",
|
||||
await mkdtemp(path.join(os.tmpdir(), "catalog-worker-")),
|
||||
);
|
||||
mocks.snapshot.mockReset().mockResolvedValue([]);
|
||||
mocks.publish.mockReset().mockResolvedValue("abc123");
|
||||
});
|
||||
afterEach(() => vi.unstubAllEnvs());
|
||||
it("retains failed work, retries it and clears only completed events", async () => {
|
||||
const q = catalogExportQueue();
|
||||
await q.request();
|
||||
mocks.publish.mockRejectedValueOnce(new Error("credential secret"));
|
||||
await runCatalogExport();
|
||||
expect(await q.batch()).toHaveLength(1);
|
||||
expect(
|
||||
await readFile(path.join(q.root, "status.json"), "utf8"),
|
||||
).not.toContain("credential secret");
|
||||
await runCatalogExport();
|
||||
expect(await q.batch()).toHaveLength(0);
|
||||
expect(mocks.publish).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
it("does not publish during an operation or if one starts while capturing", async () => {
|
||||
const q = catalogExportQueue();
|
||||
await q.request();
|
||||
await withCatalogExport(async () => {
|
||||
await runCatalogExport();
|
||||
});
|
||||
expect(mocks.snapshot).not.toHaveBeenCalled();
|
||||
mocks.snapshot.mockImplementationOnce(async () => {
|
||||
await q.request();
|
||||
return [];
|
||||
});
|
||||
await runCatalogExport();
|
||||
expect(mocks.publish).not.toHaveBeenCalled();
|
||||
expect(await q.batch()).toHaveLength(3);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,82 @@
|
||||
import { promises as fs } from "node:fs";
|
||||
import { hostname } from "node:os";
|
||||
import path from "node:path";
|
||||
import {
|
||||
CATALOG_BRANCH,
|
||||
CATALOG_REMOTE,
|
||||
publishCatalogFiles,
|
||||
recoverCatalogQueue,
|
||||
} from "./catalog-git-core";
|
||||
import { catalogExportEnabled, catalogExportQueue } from "./catalog-git-queue";
|
||||
import { createCatalogSnapshot } from "./catalog-git-snapshot";
|
||||
|
||||
export async function runCatalogExport() {
|
||||
const checkout = process.env.CATALOG_GIT_CHECKOUT?.trim();
|
||||
if (!catalogExportEnabled() || !checkout) return;
|
||||
const queue = catalogExportQueue();
|
||||
await queue.entries();
|
||||
await recoverCatalogQueue(queue.root);
|
||||
const lock = await fs
|
||||
.open(path.join(queue.root, "worker.lock"), "wx")
|
||||
.catch((error: NodeJS.ErrnoException) => {
|
||||
if (error.code === "EEXIST") return null;
|
||||
throw error;
|
||||
});
|
||||
if (!lock) return;
|
||||
let snapshot: string | undefined;
|
||||
const status = async (value: Record<string, unknown>) => {
|
||||
const temp = path.join(queue.root, "status.tmp");
|
||||
let previous: Record<string, unknown> = {};
|
||||
try {
|
||||
previous = JSON.parse(
|
||||
await fs.readFile(path.join(queue.root, "status.json"), "utf8"),
|
||||
);
|
||||
} catch {
|
||||
/* First attempt. */
|
||||
}
|
||||
await fs.writeFile(
|
||||
temp,
|
||||
JSON.stringify({ ...previous, error: undefined, ...value }),
|
||||
);
|
||||
await fs.rename(temp, path.join(queue.root, "status.json"));
|
||||
};
|
||||
try {
|
||||
await lock.writeFile(
|
||||
JSON.stringify({ pid: process.pid, host: hostname() }),
|
||||
);
|
||||
const batch = await queue.batch();
|
||||
if (!batch?.length) return;
|
||||
snapshot = await fs.mkdtemp(path.join(queue.root, "snapshot-"));
|
||||
const files = await createCatalogSnapshot(snapshot);
|
||||
// Any operation started/completed during capture invalidates the snapshot.
|
||||
const current = await queue.batch();
|
||||
if (!current || JSON.stringify(current) !== JSON.stringify(batch)) return;
|
||||
const commit = await publishCatalogFiles({
|
||||
checkout,
|
||||
remote: CATALOG_REMOTE,
|
||||
branch: CATALOG_BRANCH,
|
||||
files,
|
||||
});
|
||||
await queue.complete(batch);
|
||||
await status({ commit, finishedAt: new Date().toISOString() });
|
||||
} catch {
|
||||
// Do not expose credentials or Git stderr through the admin API.
|
||||
await status({
|
||||
error:
|
||||
"Export failed. Check repository access, branch, clean checkout and asset/database availability. Pending changes will be retried.",
|
||||
finishedAt: new Date().toISOString(),
|
||||
});
|
||||
} finally {
|
||||
try {
|
||||
if (
|
||||
snapshot &&
|
||||
path.resolve(snapshot).startsWith(path.resolve(queue.root) + path.sep)
|
||||
) {
|
||||
await fs.rm(snapshot, { recursive: true, force: true });
|
||||
}
|
||||
} finally {
|
||||
await lock.close();
|
||||
await fs.rm(path.join(queue.root, "worker.lock"), { force: true });
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
import { promises as fs } from "node:fs";
|
||||
import path from "node:path";
|
||||
import { CatalogExportQueue } from "./catalog-git-core";
|
||||
|
||||
export function catalogExportEnabled() {
|
||||
return Boolean(process.env.CATALOG_GIT_CHECKOUT?.trim());
|
||||
}
|
||||
export function catalogExportQueue() {
|
||||
return new CatalogExportQueue(
|
||||
process.env.CATALOG_GIT_STATE_DIR ||
|
||||
path.join(process.cwd(), "storage", "catalog-git"),
|
||||
);
|
||||
}
|
||||
export function isCatalogMutation(method: string, pathname: string) {
|
||||
return (
|
||||
["POST", "PUT", "PATCH", "DELETE"].includes(method) &&
|
||||
(pathname.startsWith("/api/admin/import/") ||
|
||||
pathname.startsWith("/api/admin/catalog/") ||
|
||||
pathname.startsWith("/api/admin/furni/") ||
|
||||
pathname.startsWith("/api/admin/furniture/"))
|
||||
);
|
||||
}
|
||||
export async function catalogExportStatus() {
|
||||
const q = catalogExportQueue();
|
||||
if (!catalogExportEnabled()) return { enabled: false, pending: 0, active: 0 };
|
||||
const entries = await q.entries();
|
||||
let last: { commit?: string; finishedAt?: string; error?: string } = {};
|
||||
try {
|
||||
last = JSON.parse(
|
||||
await fs.readFile(path.join(q.root, "status.json"), "utf8"),
|
||||
);
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error;
|
||||
}
|
||||
return {
|
||||
enabled: true,
|
||||
pending: entries.filter((f) => f.endsWith(".pending")).length,
|
||||
active: entries.filter((f) => f.endsWith(".active")).length,
|
||||
running: entries.includes("worker.lock"),
|
||||
...last,
|
||||
};
|
||||
}
|
||||
|
||||
export async function withCatalogExport<T>(
|
||||
operation: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
if (!catalogExportEnabled()) return operation();
|
||||
const finish = await catalogExportQueue().begin();
|
||||
try {
|
||||
return await operation();
|
||||
} finally {
|
||||
await finish();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,161 @@
|
||||
vi.mock("@/lib/services/figuredata", () => ({
|
||||
getFigureDataPath: async () =>
|
||||
path.join(await mocks.gamedata(), "custom/FigureData.json"),
|
||||
}));
|
||||
vi.mock("@/lib/services/figuremap", () => ({
|
||||
getFigureMapPath: async () =>
|
||||
path.join(await mocks.gamedata(), "custom/FigureMap.json"),
|
||||
}));
|
||||
vi.mock("@/lib/services/effectmap", () => ({
|
||||
getEffectMapPath: async () =>
|
||||
path.join(await mocks.gamedata(), "custom/EffectMap.json"),
|
||||
}));
|
||||
|
||||
import { mkdir, mkdtemp, readFile, writeFile } from "node:fs/promises";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
const mocks = vi.hoisted(() => ({
|
||||
query: vi.fn(),
|
||||
commit: vi.fn(),
|
||||
rollback: vi.fn(),
|
||||
release: vi.fn(),
|
||||
gamedata: vi.fn(),
|
||||
dirs: vi.fn(),
|
||||
data: vi.fn(),
|
||||
settings: vi.fn(),
|
||||
}));
|
||||
vi.mock("@/lib/db", () => ({
|
||||
db: { $client: { getConnection: async () => mocks } },
|
||||
}));
|
||||
vi.mock("@/lib/services/furni-asset-dirs", () => ({
|
||||
getGamedataRoot: mocks.gamedata,
|
||||
getFurniAssetDirs: mocks.dirs,
|
||||
}));
|
||||
vi.mock("@/lib/services/furni-data", () => ({
|
||||
getFurnitureDataPath: mocks.data,
|
||||
}));
|
||||
vi.mock("@/lib/services/site-settings", () => ({
|
||||
siteSettings: { get: mocks.settings },
|
||||
}));
|
||||
|
||||
import { createCatalogSnapshot } from "./catalog-git-snapshot";
|
||||
|
||||
describe("catalog snapshot", () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
mocks.settings.mockResolvedValue("");
|
||||
mocks.query.mockImplementation(async (sql: string) => {
|
||||
if (sql === "SHOW TABLES")
|
||||
return [
|
||||
[
|
||||
{ name: "items_base" },
|
||||
{ name: "catalog_pages" },
|
||||
{ name: "catalog_items" },
|
||||
],
|
||||
];
|
||||
if (sql.startsWith("SHOW CREATE"))
|
||||
return [
|
||||
[
|
||||
{
|
||||
"Create Table":
|
||||
"CREATE TABLE `test` (`id` int PRIMARY KEY, `name` text) ENGINE=InnoDB",
|
||||
},
|
||||
],
|
||||
];
|
||||
if (sql.startsWith("SHOW COLUMNS"))
|
||||
return [
|
||||
[
|
||||
{ Field: "id", Extra: "" },
|
||||
{ Field: "name", Extra: "" },
|
||||
],
|
||||
];
|
||||
if (sql.startsWith("SELECT"))
|
||||
return [[{ id: 1, name: "Valentine's chair" }]];
|
||||
return [[]];
|
||||
});
|
||||
});
|
||||
async function fixture() {
|
||||
const root = await mkdtemp(path.join(os.tmpdir(), "catalog-snapshot-"));
|
||||
for (const folder of [
|
||||
"config",
|
||||
"custom",
|
||||
"icons",
|
||||
"bundled/furniture",
|
||||
"c_images/catalogue",
|
||||
"output",
|
||||
])
|
||||
await mkdir(path.join(root, folder), { recursive: true });
|
||||
await writeFile(
|
||||
path.join(root, "config/FurnitureData.json"),
|
||||
'{"roomitemtypes":{"furnitype":[]}}',
|
||||
);
|
||||
await writeFile(path.join(root, "config/FurnitureData_it.json"), "{}");
|
||||
await writeFile(
|
||||
path.join(root, "config/private-cache.json"),
|
||||
'{"secret":"not exported"}',
|
||||
);
|
||||
await writeFile(path.join(root, "icons/chair_icon.png"), "image");
|
||||
await writeFile(path.join(root, "bundled/furniture/chair.nitro"), "bundle");
|
||||
await writeFile(
|
||||
path.join(root, "c_images/catalogue/icon_1.png"),
|
||||
"catalog icon",
|
||||
);
|
||||
await writeFile(
|
||||
path.join(root, "custom/FigureData.json"),
|
||||
'{"configured":true}',
|
||||
);
|
||||
await writeFile(path.join(root, "custom/FigureMap.json"), "{}");
|
||||
await writeFile(path.join(root, "custom/EffectMap.json"), "{}");
|
||||
mocks.gamedata.mockResolvedValue(root);
|
||||
mocks.dirs.mockResolvedValue({
|
||||
nitroDir: path.join(root, "bundled/furniture"),
|
||||
iconDir: path.join(root, "icons"),
|
||||
});
|
||||
mocks.data.mockResolvedValue(path.join(root, "config/FurnitureData.json"));
|
||||
return root;
|
||||
}
|
||||
it("includes bundles, both icon types, translations and SQL but excludes private caches", async () => {
|
||||
const root = await fixture();
|
||||
const files = await createCatalogSnapshot(path.join(root, "output"));
|
||||
const targets = files.map((f) => f.target);
|
||||
expect(targets).toContain("Gamedata/bundled/furniture/chair.nitro");
|
||||
expect(targets).toContain("Gamedata/icons/chair_icon.png");
|
||||
expect(targets).toContain("Gamedata/c_images/catalogue/icon_1.png");
|
||||
expect(targets).toContain(
|
||||
"catalogue version 2 ( Final (Dev)/langs furnidata/FurnitureData_it.json",
|
||||
);
|
||||
expect(targets.some((t) => t.includes("private-cache"))).toBe(false);
|
||||
expect(targets.filter((t) => t.endsWith(".sql"))).toHaveLength(3);
|
||||
const sqlFile = files.find((f) => f.target.endsWith("/items_base.sql"));
|
||||
expect(sqlFile).toBeDefined();
|
||||
expect(await readFile(sqlFile?.source ?? "", "utf8")).toContain(
|
||||
"ON DUPLICATE KEY UPDATE",
|
||||
);
|
||||
const figure = files.find(
|
||||
(f) => f.target === "Gamedata/config/FigureData.json",
|
||||
);
|
||||
expect(figure).toBeDefined();
|
||||
expect(await readFile(figure?.source ?? "", "utf8")).toContain(
|
||||
'"configured":true',
|
||||
);
|
||||
expect(targets).toContain("Gamedata/config/FigureMap.json");
|
||||
expect(targets).toContain("Gamedata/config/EffectMap.json");
|
||||
expect(mocks.commit).toHaveBeenCalledOnce();
|
||||
expect(mocks.release).toHaveBeenCalledOnce();
|
||||
});
|
||||
it("rejects invalid JSON before publishing and releases failed SQL snapshots", async () => {
|
||||
const root = await fixture();
|
||||
mocks.query.mockRejectedValueOnce(new Error("DB unavailable"));
|
||||
await expect(
|
||||
createCatalogSnapshot(path.join(root, "output")),
|
||||
).rejects.toThrow("DB unavailable");
|
||||
expect(mocks.rollback).toHaveBeenCalledOnce();
|
||||
expect(mocks.release).toHaveBeenCalledOnce();
|
||||
await writeFile(path.join(root, "config/FurnitureData.json"), "{");
|
||||
await expect(
|
||||
createCatalogSnapshot(path.join(root, "output")),
|
||||
).rejects.toThrow();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,216 @@
|
||||
import { promises as fs } from "node:fs";
|
||||
import path from "node:path";
|
||||
import type { RowDataPacket } from "mysql2";
|
||||
import { db } from "@/lib/db";
|
||||
import { getEffectMapPath } from "@/lib/services/effectmap";
|
||||
import { getFigureDataPath } from "@/lib/services/figuredata";
|
||||
import { getFigureMapPath } from "@/lib/services/figuremap";
|
||||
import {
|
||||
getFurniAssetDirs,
|
||||
getGamedataRoot,
|
||||
} from "@/lib/services/furni-asset-dirs";
|
||||
import { getFurnitureDataPath } from "@/lib/services/furni-data";
|
||||
import { siteSettings } from "@/lib/services/site-settings";
|
||||
import {
|
||||
CATALOG_LANGUAGE_ROOT,
|
||||
CATALOG_SQL_ROOT,
|
||||
type CatalogFile,
|
||||
sqlValue,
|
||||
} from "./catalog-git-core";
|
||||
|
||||
const CONFIG_FILES = [
|
||||
"FurnitureData.json",
|
||||
"EffectMap.json",
|
||||
"ExternalTexts.json",
|
||||
"FigureData.json",
|
||||
"FigureMap.json",
|
||||
"HabboAvatarActions.json",
|
||||
"ProductData.json",
|
||||
];
|
||||
const TABLES = ["items_base", "catalog_pages", "catalog_items"] as const;
|
||||
const identifier = (name: string) => `\`${name.replaceAll("`", "``")}\``;
|
||||
|
||||
export async function createCatalogSnapshot(
|
||||
directory: string,
|
||||
): Promise<CatalogFile[]> {
|
||||
const sources = new Map<string, string>();
|
||||
const addTree = async (root: string, target: string, extensions: RegExp) => {
|
||||
async function walk(source: string, relative = "") {
|
||||
const entries = await fs
|
||||
.readdir(source, { withFileTypes: true })
|
||||
.catch((error: NodeJS.ErrnoException) => {
|
||||
if (error.code === "ENOENT") return [];
|
||||
throw error;
|
||||
});
|
||||
for (const entry of entries) {
|
||||
if (entry.isSymbolicLink())
|
||||
throw new Error("Symlink in catalog source");
|
||||
const local = path.join(source, entry.name);
|
||||
const dest = relative ? `${relative}/${entry.name}` : entry.name;
|
||||
if (entry.isDirectory()) await walk(local, dest);
|
||||
else if (entry.isFile() && extensions.test(entry.name))
|
||||
sources.set(`${target}/${dest}`, local);
|
||||
}
|
||||
}
|
||||
await walk(root);
|
||||
};
|
||||
const gamedata = await getGamedataRoot();
|
||||
const furni = await getFurniAssetDirs();
|
||||
const furnitureData = await getFurnitureDataPath();
|
||||
// Default files first; configured live locations take precedence.
|
||||
await addTree(
|
||||
path.join(process.cwd(), "public/swf/c_images"),
|
||||
"Gamedata/c_images",
|
||||
/\.(png|gif|webp|jpe?g)$/i,
|
||||
);
|
||||
await addTree(
|
||||
path.join(process.cwd(), "public/nitro-assets/bundled"),
|
||||
"Gamedata/bundled",
|
||||
/\.nitro$/i,
|
||||
);
|
||||
if (gamedata) {
|
||||
await addTree(
|
||||
path.join(gamedata, "bundled"),
|
||||
"Gamedata/bundled",
|
||||
/\.nitro$/i,
|
||||
);
|
||||
await addTree(
|
||||
path.join(gamedata, "icons"),
|
||||
"Gamedata/icons",
|
||||
/\.(png|gif|webp|jpe?g)$/i,
|
||||
);
|
||||
await addTree(
|
||||
path.join(gamedata, "c_images"),
|
||||
"Gamedata/c_images",
|
||||
/\.(png|gif|webp|jpe?g)$/i,
|
||||
);
|
||||
}
|
||||
await addTree(furni.nitroDir, "Gamedata/bundled/furniture", /\.nitro$/i);
|
||||
await addTree(furni.iconDir, "Gamedata/icons", /\.(png|gif|webp|jpe?g)$/i);
|
||||
for (const type of ["figure", "effect"]) {
|
||||
const configured = await siteSettings.get(`${type}_nitro_dir`, "");
|
||||
if (configured?.trim())
|
||||
await addTree(configured, `Gamedata/bundled/${type}`, /\.nitro$/i);
|
||||
}
|
||||
// Pets currently write directly to the CMS directory.
|
||||
await addTree(
|
||||
path.join(process.cwd(), "public/nitro-assets/bundled/pet"),
|
||||
"Gamedata/bundled/pet",
|
||||
/\.nitro$/i,
|
||||
);
|
||||
const configRoot = gamedata
|
||||
? path.join(gamedata, "config")
|
||||
: path.dirname(furnitureData);
|
||||
const configuredJson: Record<string, string> = {
|
||||
FurnitureData: furnitureData,
|
||||
FigureData: await getFigureDataPath(),
|
||||
FigureMap: await getFigureMapPath(),
|
||||
EffectMap: await getEffectMapPath(),
|
||||
};
|
||||
for (const name of CONFIG_FILES) {
|
||||
const source =
|
||||
configuredJson[name.replace(/\.json$/, "")] ??
|
||||
path.join(configRoot, name);
|
||||
try {
|
||||
await fs.access(source);
|
||||
sources.set(`Gamedata/config/${name}`, source);
|
||||
} catch (error) {
|
||||
if (
|
||||
name === "FurnitureData.json" ||
|
||||
(error as NodeJS.ErrnoException).code !== "ENOENT"
|
||||
)
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
for (const name of await fs.readdir(path.dirname(furnitureData))) {
|
||||
if (/^FurnitureData_[a-z]{2}\.json$/.test(name)) {
|
||||
sources.set(
|
||||
`${CATALOG_LANGUAGE_ROOT}/${name}`,
|
||||
path.join(path.dirname(furnitureData), name),
|
||||
);
|
||||
}
|
||||
}
|
||||
const files: CatalogFile[] = [];
|
||||
const fingerprints = new Map<string, string>();
|
||||
const fingerprint = async (source: string) => {
|
||||
const stat = await fs.lstat(source);
|
||||
if (!stat.isFile() || stat.isSymbolicLink())
|
||||
throw new Error("Invalid catalog source");
|
||||
return `${stat.size}:${stat.mtimeMs}:${stat.ctimeMs}`;
|
||||
};
|
||||
for (const [target, source] of sources) {
|
||||
fingerprints.set(source, await fingerprint(source));
|
||||
const copy = path.join(directory, String(files.length));
|
||||
await fs.copyFile(source, copy);
|
||||
if (target.endsWith(".json")) JSON.parse(await fs.readFile(copy, "utf8"));
|
||||
files.push({ source: copy, target });
|
||||
}
|
||||
const connection = await db.$client.getConnection();
|
||||
try {
|
||||
await connection.query("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ");
|
||||
await connection.query(
|
||||
"START TRANSACTION WITH CONSISTENT SNAPSHOT, READ ONLY",
|
||||
);
|
||||
const [available] = await connection.query<RowDataPacket[]>("SHOW TABLES");
|
||||
const tableNames = new Set(
|
||||
available.map((row) => String(Object.values(row)[0])),
|
||||
);
|
||||
const tables = [
|
||||
...TABLES,
|
||||
...["catalog_pages_bc", "catalog_items_bc"].filter((name) =>
|
||||
tableNames.has(name),
|
||||
),
|
||||
];
|
||||
for (const table of tables) {
|
||||
const [schemaRows] = await connection.query<RowDataPacket[]>(
|
||||
`SHOW CREATE TABLE ${identifier(table)}`,
|
||||
);
|
||||
const definition = String(schemaRows[0]["Create Table"]);
|
||||
if (!/ENGINE=InnoDB\b/i.test(definition))
|
||||
throw new Error("Catalog SQL snapshots require InnoDB tables");
|
||||
const [columns] = await connection.query<RowDataPacket[]>(
|
||||
`SHOW COLUMNS FROM ${identifier(table)}`,
|
||||
);
|
||||
const names = columns
|
||||
.filter(
|
||||
(c) => !/\b(?:VIRTUAL|STORED) GENERATED\b/.test(String(c.Extra)),
|
||||
)
|
||||
.map((c) => String(c.Field));
|
||||
const quoted = names.map(identifier).join(", ");
|
||||
const [rows] = await connection.query<RowDataPacket[]>(
|
||||
`SELECT ${quoted} FROM ${identifier(table)} ORDER BY id`,
|
||||
);
|
||||
const source = path.join(directory, `${table}.sql`);
|
||||
const output = await fs.open(source, "wx");
|
||||
try {
|
||||
await output.writeFile(
|
||||
`-- Catalog Studio snapshot. Existing rows are updated; absent rows are not deleted.\nSET NAMES utf8mb4;\n${definition.replace("CREATE TABLE", "CREATE TABLE IF NOT EXISTS")};\nSTART TRANSACTION;\n`,
|
||||
);
|
||||
const update = names
|
||||
.filter((n) => n !== "id")
|
||||
.map((n) => `${identifier(n)}=VALUES(${identifier(n)})`)
|
||||
.join(", ");
|
||||
for (const row of rows) {
|
||||
await output.writeFile(
|
||||
`INSERT INTO ${identifier(table)} (${quoted}) VALUES (${names.map((n) => sqlValue(row[n])).join(", ")}) ON DUPLICATE KEY UPDATE ${update};\n`,
|
||||
);
|
||||
}
|
||||
await output.writeFile("COMMIT;\n");
|
||||
} finally {
|
||||
await output.close();
|
||||
}
|
||||
files.push({ source, target: `${CATALOG_SQL_ROOT}/${table}.sql` });
|
||||
}
|
||||
await connection.commit();
|
||||
} catch (error) {
|
||||
await connection.rollback();
|
||||
throw error;
|
||||
} finally {
|
||||
connection.release();
|
||||
}
|
||||
for (const [source, stamp] of fingerprints) {
|
||||
if ((await fingerprint(source)) !== stamp)
|
||||
throw new Error("Assets changed during export; retry required");
|
||||
}
|
||||
return files;
|
||||
}
|
||||
Reference in new issue
Block a user