feat(cms): improve catalog, editorial recovery and operations
CI / check (push) Successful in 52s
CI / deploy (push) Successful in 2m10s
CI / publish-container (push) Failing after 1m18s

This commit is contained in:
Simo committed 2026-09-09 19:36:15 +02:00
1 parent 2fead134e6
commit c389c3893d
122 files changed
+8137 -703

No files matched your search

+62
View File
@@ -0,0 +1,62 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
const state = vi.hoisted(() => ({
total: 55,
queries: [] as { sql: string; params: unknown[] }[],
}));
vi.mock("@/lib/db", async () => {
const schema = await import("@/db/schema");
const { drizzle } = await import("drizzle-orm/mysql-proxy");
return {
...schema,
db: drizzle(async (sql, params) => {
state.queries.push({ sql, params });
return { rows: sql.includes("count(*)") ? [[state.total]] : [] };
}),
};
});
import { loadArticleList } from "./article-list";
describe("article listing", () => {
beforeEach(() => {
state.queries = [];
state.total = 55;
});
it("filters drafts before pagination and binds search text", async () => {
const result = await loadArticleList({
status: "draft",
search: "needle'",
page: 2,
});
expect(result.page).toBe(2);
for (const query of state.queries) {
expect(query.sql).toContain("`website_articles`.`status` = ?");
expect(query.sql).not.toContain("needle");
expect(query.params.slice(0, 3)).toEqual([
"draft",
"%needle'%",
"%needle'%",
]);
}
expect(state.queries[1].params.slice(-2)).toEqual([24, 24]);
});
it("clamps an out-of-range page to the last filtered page", async () => {
const result = await loadArticleList({ status: "scheduled", page: 999 });
expect(result).toMatchObject({ page: 3, lastPage: 3, total: 55 });
expect(state.queries[1].params.slice(-2)).toEqual([24, 48]);
});
it("bounds invalid input and keeps an empty collection on page one", async () => {
state.total = 0;
const result = await loadArticleList({ status: "unexpected", page: NaN });
expect(result).toMatchObject({
status: "all",
page: 1,
lastPage: 1,
total: 0,
});
expect(state.queries[1].sql).not.toContain(
"`website_articles`.`status` = ?",
);
});
});
+51
View File
@@ -0,0 +1,51 @@
import "server-only";
import { and, count, desc, eq, like, or } from "drizzle-orm";
import { db, User, WebsiteArticles } from "@/lib/db";
export async function loadArticleList(
options: { status?: string; search?: string; page?: number } = {},
) {
const status = ["draft", "published", "scheduled"].includes(
options.status ?? "",
)
? options.status
: "all";
const search = (options.search ?? "").trim().slice(0, 191);
const where = and(
status !== "all" ? eq(WebsiteArticles.status, status ?? "all") : undefined,
search
? or(
like(WebsiteArticles.title, `%${search}%`),
like(WebsiteArticles.shortStory, `%${search}%`),
)
: undefined,
);
const totals = await db
.select({ value: count() })
.from(WebsiteArticles)
.where(where);
const total = Number(totals[0]?.value ?? 0);
const perPage = 24;
const lastPage = Math.max(1, Math.ceil(total / perPage));
const page = Math.min(
lastPage,
Number.isFinite(options.page)
? Math.max(1, Math.trunc(options.page ?? 1))
: 1,
);
const rows = await db
.select({
id: WebsiteArticles.id,
title: WebsiteArticles.title,
image: WebsiteArticles.image,
createdAt: WebsiteArticles.createdAt,
status: WebsiteArticles.status,
author: User.username,
})
.from(WebsiteArticles)
.leftJoin(User, eq(User.id, WebsiteArticles.userId))
.where(where)
.orderBy(desc(WebsiteArticles.createdAt), desc(WebsiteArticles.id))
.limit(perPage)
.offset((page - 1) * perPage);
return { rows, page, perPage, total, lastPage, status, search };
}
+67
View File
@@ -0,0 +1,67 @@
type Change = { key: string; from: unknown; to: unknown };
const sensitive =
/password|secret|token|otp|recovery|authTicket|two_factor|api_key/i;
function record(value: unknown): value is Record<string, unknown> {
return value !== null && typeof value === "object" && !Array.isArray(value);
}
function redact(value: unknown, depth = 0): unknown {
if (depth > 10) return "[…]";
if (Array.isArray(value)) return value.map((item) => redact(item, depth + 1));
if (!record(value)) return value;
return Object.fromEntries(
Object.entries(value).map(([key, item]) => [
key,
sensitive.test(key) ? "[Redacted]" : redact(item, depth + 1),
]),
);
}
function parse(value: string) {
if (value.length > 100000) throw new Error("oversize");
const parsed: unknown = JSON.parse(value);
if (!record(parsed)) throw new Error("invalid");
return parsed;
}
export function readAuditChanges(
diff: string | null,
before?: string | null,
after?: string | null,
) {
try {
let changes: Change[];
if (diff) {
const parsed = parse(diff);
changes = Object.entries(parsed).map(([key, value]) => {
if (!record(value) || !("from" in value || "to" in value))
throw new Error("invalid");
return { key, from: value.from, to: value.to };
});
} else {
const previous = before ? parse(before) : {};
const next = after ? parse(after) : {};
changes = [...new Set([...Object.keys(previous), ...Object.keys(next)])]
.filter(
(key) => JSON.stringify(previous[key]) !== JSON.stringify(next[key]),
)
.map((key) => ({ key, from: previous[key], to: next[key] }));
}
return {
invalid: false,
changes: changes.map(({ key, from, to }) => ({
key,
from: sensitive.test(key) ? "[Redacted]" : redact(from),
to: sensitive.test(key) ? "[Redacted]" : redact(to),
})),
};
} catch {
return { invalid: true, changes: [] as Change[] };
}
}
export function formatAuditValue(value: unknown) {
const text =
value === undefined
? "∅"
: typeof value === "string"
? value
: JSON.stringify(value, null, 2);
return text.length > 4000 ? `${text.slice(0, 4000)}…` : text;
}
+39
View File
@@ -0,0 +1,39 @@
export interface AuditFilters {
search?: string;
actor?: string;
action?: string;
from?: string;
to?: string;
page?: number;
perPage?: number;
}
function day(value?: string) {
if (!value || !/^\d{4}-\d{2}-\d{2}$/.test(value)) return "";
const date = new Date(`${value}T00:00:00Z`);
return Number.isFinite(date.getTime()) &&
date.toISOString().slice(0, 10) === value
? value
: "";
}
export function normalizeAuditFilters(options: AuditFilters = {}) {
const from = day(options.from);
const to = day(options.to);
return {
search: (options.search ?? "").trim().slice(0, 191),
actor: (options.actor ?? "").trim().slice(0, 191),
action: (options.action ?? "").trim().slice(0, 191),
from,
to,
until: to
? new Date(new Date(`${to}T00:00:00Z`).getTime() + 86400000)
.toISOString()
.slice(0, 10)
: "",
page: Number.isFinite(options.page)
? Math.min(1000000, Math.max(1, Math.trunc(options.page ?? 1)))
: 1,
perPage: Number.isFinite(options.perPage)
? Math.min(100, Math.max(1, Math.trunc(options.perPage ?? 20)))
: 20,
};
}
+56
View File
@@ -0,0 +1,56 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
const state = vi.hoisted(() => ({
queries: [] as { sql: string; params: unknown[] }[],
}));
vi.mock("@/lib/db", async () => {
const schema = await import("@/db/schema");
const { drizzle } = await import("drizzle-orm/mysql-proxy");
return {
...schema,
db: drizzle(async (sql, params) => {
state.queries.push({ sql, params });
return { rows: [] };
}),
};
});
import { getAuditLogs } from "./audit";
describe("audit query filters", () => {
beforeEach(() => {
state.queries = [];
});
it("binds actor/action/date/search values in both list and count queries", async () => {
await getAuditLogs({
actor: "Alice' OR 1=1 --",
action: "update",
search: "user",
from: "2026-09-09",
to: "2026-09-09",
perPage: 10,
page: 2,
});
expect(state.queries).toHaveLength(2);
for (const query of state.queries) {
expect(query.sql).not.toContain("Alice");
expect(query.sql).toContain("`admin_audit_log`.`user_id` in (select");
expect(query.sql).toContain("`admin_audit_log`.`created_at` >= ?");
expect(query.sql).toContain("`admin_audit_log`.`created_at` < ?");
expect(query.params.slice(0, 6)).toEqual([
"%user%",
"%user%",
"update",
"Alice' OR 1=1 --",
"2026-09-09",
"2026-09-10",
]);
}
expect(state.queries[0].params.slice(-2)).toEqual([10, 10]);
});
it("uses an exact numeric actor ID and caps row count", async () => {
await getAuditLogs({ actor: "42", perPage: 1000, page: -2 });
expect(state.queries[0].params).toEqual([42, 100]);
expect(state.queries[0].sql).not.toContain("in (select");
});
});
+70
View File
@@ -0,0 +1,70 @@
import { describe, expect, it } from "vitest";
import { formatAuditValue, readAuditChanges } from "./audit-diff";
import { normalizeAuditFilters } from "./audit-filters";
describe("audit filter bounds", () => {
it("normalizes pagination and rejects invalid calendar dates", () => {
expect(
normalizeAuditFilters({
page: NaN,
perPage: Infinity,
from: "2026-02-30",
to: "x",
}),
).toMatchObject({ page: 1, perPage: 20, from: "", to: "" });
expect(normalizeAuditFilters({ page: -3, perPage: 1000 })).toMatchObject({
page: 1,
perPage: 100,
});
});
it("caps strings and preserves exact actor/action and UTC day boundaries", () => {
expect(
normalizeAuditFilters({
actor: " Alice ",
action: " update ",
from: "2026-09-09",
to: "2026-09-09",
}),
).toMatchObject({
actor: "Alice",
action: "update",
from: "2026-09-09",
until: "2026-09-10",
});
expect(
normalizeAuditFilters({ search: "x".repeat(500) }).search,
).toHaveLength(191);
});
});
describe("legacy audit diffs", () => {
it.each(["null", "[]", "1", "bad JSON", '{"name":null}', '{"name":"value"}'])(
"handles %s safely",
(diff) => {
expect(readAuditChanges(diff).invalid).toBe(true);
},
);
it("keeps objects and explicit null distinct from absent values", () => {
const result = readAuditChanges(
'{"settings":{"from":{"enabled":false},"to":null},"added":{"to":3}}',
);
expect(result.changes).toHaveLength(2);
expect(formatAuditValue(result.changes[0].from)).toContain(
'"enabled": false',
);
expect(formatAuditValue(result.changes[0].to)).toBe("null");
expect(formatAuditValue(result.changes[1].from)).toBe("∅");
});
it("shows create/delete snapshots when a legacy record has no diff", () => {
expect(readAuditChanges(null, null, '{"name":"New"}').changes).toEqual([
{ key: "name", from: undefined, to: "New" },
]);
});
it("redacts historical sensitive nested fields before rendering", () => {
expect(
formatAuditValue(
readAuditChanges('{"settings":{"from":{"password":"secret"},"to":{}}}')
.changes[0].from,
),
).not.toContain("secret");
});
});
+21
View File
@@ -187,3 +187,24 @@ describe("getAuditLogs", () => {
expect(result.rows[0].username).toBe("User #99");
});
});
describe("audit response privacy", () => {
it("does not serialize legacy snapshot secrets to the client", async () => {
selectRows.mockResolvedValue([
{
id: 1,
userId: 1,
diff: null,
before: null,
after: '{"name":"Alice","password":"secret-value"}',
},
]);
selectCount.mockResolvedValue([{ value: 1 }]);
selectUsers.mockResolvedValue([]);
const result = await getAuditLogs();
expect(JSON.stringify(result.rows)).not.toContain("secret-value");
expect(result.rows[0].before).toBeNull();
expect(result.rows[0].after).toBeNull();
expect(result.rows[0].diff).toContain("Alice");
});
});
+49 -21
View File
@@ -1,5 +1,7 @@
import { count, desc, inArray, like, or } from "drizzle-orm";
import { and, count, desc, eq, gte, inArray, like, lt, or } from "drizzle-orm";
import { AdminAuditLog, db, User } from "@/lib/db";
import { readAuditChanges } from "./audit-diff";
import { type AuditFilters, normalizeAuditFilters } from "./audit-filters";
interface AuditEntry {
userId: number;
@@ -68,23 +70,32 @@ export async function logAudit(entry: AuditEntry): Promise<void> {
});
}
interface GetLogsOptions {
search?: string;
page?: number;
perPage?: number;
}
export async function getAuditLogs(options: GetLogsOptions = {}) {
const { search, page = 1, perPage = 20 } = options;
export async function getAuditLogs(options: AuditFilters = {}) {
const { search, actor, action, from, until, page, perPage } =
normalizeAuditFilters(options);
const skip = (page - 1) * perPage;
const where = search
? or(
like(AdminAuditLog.action, `%${search}%`),
like(AdminAuditLog.target, `%${search}%`),
)
: undefined;
const where = and(
search
? or(
like(AdminAuditLog.action, `%${search}%`),
like(AdminAuditLog.target, `%${search}%`),
)
: undefined,
action ? eq(AdminAuditLog.action, action) : undefined,
actor
? /^[1-9]\d*$/.test(actor) && Number.isSafeInteger(Number(actor))
? eq(AdminAuditLog.userId, Number(actor))
: inArray(
AdminAuditLog.userId,
db
.select({ id: User.id })
.from(User)
.where(eq(User.username, actor)),
)
: undefined,
from ? gte(AdminAuditLog.createdAt, from) : undefined,
until ? lt(AdminAuditLog.createdAt, until) : undefined,
);
const [rows, totalResult] = await Promise.all([
db
.select({
@@ -94,6 +105,8 @@ export async function getAuditLogs(options: GetLogsOptions = {}) {
target: AdminAuditLog.target,
targetId: AdminAuditLog.targetId,
diff: AdminAuditLog.diff,
before: AdminAuditLog.before,
after: AdminAuditLog.after,
createdAt: AdminAuditLog.createdAt,
})
.from(AdminAuditLog)
@@ -116,10 +129,25 @@ export async function getAuditLogs(options: GetLogsOptions = {}) {
: [];
const userMap = new Map(users.map((u) => [u.id, u.username]));
const enrichedRows = rows.map((r) => ({
...r,
username: userMap.get(r.userId) ?? `User #${r.userId}`,
}));
const enrichedRows = rows.map((r) => {
const details = readAuditChanges(r.diff, r.before, r.after);
return {
...r,
// Only sanitized changes may cross the server/client boundary.
before: null,
after: null,
diff: details.invalid
? "[Unavailable legacy details]"
: details.changes.length
? JSON.stringify(
Object.fromEntries(
details.changes.map(({ key, from, to }) => [key, { from, to }]),
),
)
: null,
username: userMap.get(r.userId) ?? `User #${r.userId}`,
};
});
return {
rows: enrichedRows,
@@ -0,0 +1,53 @@
import { describe, expect, it } from "vitest";
import { getDashboardEventPhase } from "./dashboard-event-phase";
const now = new Date("2026-09-09T12:00:00Z");
const event = (
startsAt: string,
endsAt: string | null = null,
status = "published",
) => ({
startsAt: new Date(startsAt),
endsAt: endsAt ? new Date(endsAt) : null,
status,
});
describe("dashboard event timing", () => {
it("shows a future published event without requiring an end", () => {
expect(getDashboardEventPhase(event("2026-09-09T13:00:00Z"), now)).toBe(
"upcoming",
);
});
it("marks the start boundary and known duration as ongoing", () => {
expect(
getDashboardEventPhase(
event("2026-09-09T12:00:00Z", "2026-09-09T13:00:00Z"),
now,
),
).toBe("ongoing");
});
it("excludes ended events including the exact end boundary", () => {
expect(
getDashboardEventPhase(
event("2026-09-09T11:00:00Z", "2026-09-09T12:00:00Z"),
now,
),
).toBeNull();
});
it("does not invent a duration for old events without an end", () => {
expect(
getDashboardEventPhase(event("2026-09-08T12:00:00Z"), now),
).toBeNull();
});
it.each(["draft", "cancelled", "completed"])(
"excludes %s events",
(status) => {
expect(
getDashboardEventPhase(
event("2026-09-09T13:00:00Z", null, status),
now,
),
).toBeNull();
},
);
});
@@ -0,0 +1,8 @@
export function getDashboardEventPhase(
event: { status: string; startsAt: Date; endsAt: Date | null },
now: Date,
): "upcoming" | "ongoing" | null {
if (event.status !== "published") return null;
if (event.startsAt > now) return "upcoming";
return event.endsAt && event.endsAt > now ? "ongoing" : null;
}
+68
View File
@@ -0,0 +1,68 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
const state = vi.hoisted(() => ({
rows: [] as unknown[],
fail: false,
query: "",
params: [] as unknown[],
}));
vi.mock("@/lib/db", async () => {
const schema = await import("@/db/schema");
const { drizzle } = await import("drizzle-orm/mysql-proxy");
return {
...schema,
db: drizzle(async (query, params) => {
state.query = query;
state.params = params;
if (state.fail) throw new Error("offline");
return { rows: state.rows };
}),
};
});
import { loadDashboardEvent } from "./dashboard-event";
describe("dashboard event loading", () => {
beforeEach(() => {
state.rows = [];
state.fail = false;
});
it("limits to one published ongoing or future event in date order", async () => {
await loadDashboardEvent(new Date("2026-09-09T12:00:00Z"));
expect(state.query).toContain("`website_events`.`status` = ?");
expect(state.query).toContain(
"(`website_events`.`starts_at` > ? or `website_events`.`ends_at` > ?)",
);
expect(state.query).toContain(
"order by `website_events`.`starts_at` asc, `website_events`.`id` asc limit ?",
);
expect(state.params).toEqual([
"published",
"2026-09-09 12:00:00.000",
"2026-09-09 12:00:00.000",
1,
]);
});
it("distinguishes an empty calendar from a failed query", async () => {
expect(await loadDashboardEvent()).toEqual({
event: null,
unavailable: false,
});
state.fail = true;
expect(await loadDashboardEvent()).toEqual({
event: null,
unavailable: true,
});
});
it("returns the next event with the phase derived from its real dates", async () => {
state.rows = [
[3, "Quiz", "published", "2026-09-09 13:00:00.000", null, "Games"],
];
expect(
await loadDashboardEvent(new Date("2026-09-09T12:00:00Z")),
).toMatchObject({
unavailable: false,
event: { id: 3, title: "Quiz", phase: "upcoming" },
});
});
});
+36
View File
@@ -0,0 +1,36 @@
import "server-only";
import { and, asc, eq, gt, or } from "drizzle-orm";
import { db, WebsiteEvent, WebsiteEventType } from "@/lib/db";
import { getDashboardEventPhase } from "./dashboard-event-phase";
export async function loadDashboardEvent(now = new Date()) {
try {
const rows = await db
.select({
id: WebsiteEvent.id,
title: WebsiteEvent.title,
status: WebsiteEvent.status,
startsAt: WebsiteEvent.startsAt,
endsAt: WebsiteEvent.endsAt,
typeName: WebsiteEventType.name,
})
.from(WebsiteEvent)
.innerJoin(WebsiteEventType, eq(WebsiteEvent.typeId, WebsiteEventType.id))
.where(
and(
eq(WebsiteEvent.status, "published"),
or(gt(WebsiteEvent.startsAt, now), gt(WebsiteEvent.endsAt, now)),
),
)
.orderBy(asc(WebsiteEvent.startsAt), asc(WebsiteEvent.id))
.limit(1);
const event = rows[0];
const phase = event ? getDashboardEventPhase(event, now) : null;
return {
event: event && phase ? { ...event, phase } : null,
unavailable: false,
};
} catch {
return { event: null, unavailable: true };
}
}
+41
View File
@@ -59,3 +59,44 @@ it("rejects unsafe job paths", async () => {
);
await expect(s.read("../escape")).rejects.toThrow("Invalid import ID");
});
it("cancellation survives a stale worker save", async () => {
const s = await store(),
j = job();
await s.create(j);
expect((await s.requestCancellation(j.id, j.userId))?.cancelRequested).toBe(
true,
);
await s.save({ ...j, state: "running" });
expect(await s.isCancellationRequested(j.id)).toBe(true);
expect((await s.read(j.id)).cancelRequested).toBe(true);
});
it("rejects cancellation by a different operator and is idempotent for the owner", async () => {
const s = await store(),
j = job();
await s.create(j);
expect(await s.requestCancellation(j.id, 99)).toBeNull();
expect(await s.isCancellationRequested(j.id)).toBe(false);
await Promise.all([
s.requestCancellation(j.id, j.userId),
s.requestCancellation(j.id, j.userId),
]);
expect(await s.isCancellationRequested(j.id)).toBe(true);
});
it("does not alter completed jobs when cancellation arrives late", async () => {
const s = await store(),
j = { ...job(), state: "completed" as const };
await s.create(j);
expect((await s.requestCancellation(j.id, j.userId))?.state).toBe(
"completed",
);
expect(await s.isCancellationRequested(j.id)).toBe(false);
});
it("loads bounded history for the requesting owner", async () => {
const s = await store();
for (let i = 0; i < 5; i++)
await s.create({ ...job(), userId: i % 2 === 0 ? 4 : 8 });
const result = await s.list({ userId: 4, limit: 2 });
expect(result).toHaveLength(2);
expect(result.every((job) => job.userId === 4)).toBe(true);
});
+55 -7
View File
@@ -39,20 +39,68 @@ export class ImportJobStore {
}
async read(id: string): Promise<ImportJob> {
if (!validJobId(id)) throw Error("Invalid import ID");
return JSON.parse(
const job: ImportJob = JSON.parse(
await fs.readFile(path.join(this.root, `${id}.json`), "utf8"),
);
if (await this.isCancellationRequested(id)) job.cancelRequested = true;
return job;
}
async list(): Promise<ImportJob[]> {
async isCancellationRequested(id: string): Promise<boolean> {
if (!validJobId(id)) throw Error("Invalid import ID");
try {
await fs.access(path.join(this.root, `${id}.cancel`));
return true;
} catch (error) {
if ((error as NodeJS.ErrnoException).code === "ENOENT") return false;
throw error;
}
}
/** Cancellation has its own durable marker: a worker save cannot overwrite it. */
async requestCancellation(
id: string,
userId: number,
): Promise<ImportJob | null> {
const job = await this.read(id).catch((error) => {
if (error.code === "ENOENT") return null;
throw error;
});
if (!job || job.userId !== userId) return null;
if (job.state !== "queued" && job.state !== "running") return job;
await fs
.writeFile(path.join(this.root, `${id}.cancel`), "cancel", { flag: "wx" })
.catch((error) => {
if (error.code !== "EEXIST") throw error;
});
return this.read(id);
}
async list(options?: {
userId: number;
limit: number;
}): 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))),
const candidates: Array<{ id: string; modified: number }> = [];
// Metadata and payload reads are sequential, avoiding thousands of concurrent opens.
for (const file of files) {
if (!file.endsWith(".json") || !validJobId(file.slice(0, -5))) continue;
const stat = await fs.stat(path.join(this.root, file));
candidates.push({ id: file.slice(0, -5), modified: stat.mtimeMs });
}
candidates.sort((a, b) =>
options ? b.modified - a.modified : a.modified - b.modified,
);
return jobs.sort((a, b) => a.createdAt.localeCompare(b.createdAt));
const jobs: ImportJob[] = [];
const limit = options ? Math.min(30, Math.max(1, options.limit)) : Infinity;
for (const candidate of candidates) {
const job = await this.read(candidate.id);
if (options && job.userId !== options.userId) continue;
jobs.push(job);
if (jobs.length >= limit) break;
}
return options
? jobs
: jobs.sort((a, b) => a.createdAt.localeCompare(b.createdAt));
}
}
+113 -2
View File
@@ -5,11 +5,14 @@ const mocks = vi.hoisted(() => ({
set: vi.fn(),
eval: vi.fn(),
list: vi.fn(),
cancel: vi.fn(),
save: vi.fn(),
import: vi.fn(),
attachment: vi.fn(),
export: vi.fn(),
translate: vi.fn(),
updateCatalog: vi.fn(),
updateItems: vi.fn(),
}));
vi.mock("@/lib/redis", () => ({ redis: { set: mocks.set, eval: mocks.eval } }));
vi.mock("@/lib/server-log", () => ({ logServerError: vi.fn() }));
@@ -19,6 +22,7 @@ vi.mock("./clone-sources", () => ({ getSource: async () => null }));
vi.mock("./furni-job-store", () => ({
ImportJobStore: class {
list = mocks.list;
isCancellationRequested = mocks.cancel;
save = mocks.save;
},
}));
@@ -33,7 +37,7 @@ vi.mock("./furni-attachment", () => ({
readFurnitureAttachment: mocks.attachment,
}));
vi.mock("./rcon", () => ({
rcon: { updateCatalog: async () => {}, updateItems: async () => {} },
rcon: { updateCatalog: mocks.updateCatalog, updateItems: mocks.updateItems },
}));
vi.mock("./catalog-git-export", () => ({ runCatalogExport: vi.fn() }));
@@ -44,10 +48,14 @@ let job: ImportJob;
beforeEach(() => {
vi.clearAllMocks();
mocks.set.mockResolvedValue("OK");
mocks.cancel.mockResolvedValue(false);
mocks.eval.mockResolvedValue(1);
mocks.export.mockImplementation((fn) => fn());
mocks.import.mockResolvedValue({ ok: true, itemId: 900, warnings: [] });
mocks.save.mockResolvedValue(undefined);
mocks.updateCatalog.mockResolvedValue(undefined);
mocks.updateItems.mockResolvedValue(undefined);
mocks.translate.mockResolvedValue(undefined);
job = {
id: "job",
userId: 5,
@@ -89,7 +97,7 @@ it("marks abandoned work interrupted without repeating an uncertain import", asy
job.items[0].state = "running";
await drainFurnitureImports();
expect(job.state).toBe("interrupted");
expect(job.items[0].state).toBe("failed");
expect(job.items[0].state).toBe("interrupted");
expect(mocks.import).not.toHaveBeenCalled();
});
it("keeps errors for failed items and continues the batch", async () => {
@@ -112,3 +120,106 @@ it("loads attachments using the job owner and exact classname", async () => {
expect.objectContaining({ providedNitro: Buffer.from("verified") }),
);
});
it("cancels queued work without importing any items", async () => {
mocks.cancel.mockResolvedValue(true);
await drainFurnitureImports();
expect(job.state).toBe("cancelled");
expect(job.items[0].state).toBe("cancelled");
expect(mocks.import).not.toHaveBeenCalled();
});
it("finishes the in-flight item and cancels only the remaining items", async () => {
job.items.push({ ...job.items[0], classname: "table" });
mocks.import.mockImplementationOnce(async () => {
mocks.cancel.mockResolvedValue(true);
return { ok: true, itemId: 900, warnings: [] };
});
await drainFurnitureImports();
expect(job.state).toBe("cancelled");
expect(job.items.map((item) => item.state)).toEqual(["done", "cancelled"]);
expect(mocks.import).toHaveBeenCalledOnce();
});
it("reports completion if cancellation arrives during the final item", async () => {
mocks.import.mockImplementationOnce(async () => {
mocks.cancel.mockResolvedValue(true);
return { ok: true, itemId: 900, warnings: [] };
});
await drainFurnitureImports();
expect(job.state).toBe("completed");
});
it("preserves pending items when recovery finds an uncertain in-flight import", async () => {
job.state = "running";
job.items[0].state = "running";
job.items.push({ ...job.items[0], classname: "table", state: "pending" });
await drainFurnitureImports();
expect(job.items.map((item) => item.state)).toEqual([
"interrupted",
"pending",
]);
expect(mocks.import).not.toHaveBeenCalled();
});
it("does not overwrite another worker's interruption after losing ownership during the final import", async () => {
let persisted: ImportJob | undefined;
mocks.save.mockImplementation(async (value: ImportJob) => {
persisted = structuredClone(value);
});
mocks.import.mockImplementationOnce(async () => {
persisted = {
...structuredClone(job),
state: "interrupted",
items: job.items.map((item) => ({ ...item, state: "interrupted" })),
};
mocks.eval.mockResolvedValue(0);
return { ok: true, itemId: 900, warnings: [] };
});
await drainFurnitureImports();
expect(mocks.import).toHaveBeenCalledOnce();
expect(mocks.save).toHaveBeenCalledTimes(2);
expect(persisted?.state).toBe("interrupted");
expect(persisted?.items[0].state).toBe("interrupted");
expect(mocks.updateCatalog).not.toHaveBeenCalled();
});
it("keeps an uncertain import running on disk when lease verification fails", async () => {
const snapshots: ImportJob[] = [];
mocks.save.mockImplementation(async (value: ImportJob) => {
snapshots.push(structuredClone(value));
});
mocks.import.mockImplementationOnce(async () => {
mocks.eval.mockRejectedValue(new Error("Redis unavailable"));
throw new Error("SQL outcome unknown");
});
await drainFurnitureImports();
expect(snapshots).toHaveLength(2);
expect(snapshots.at(-1)?.items[0].state).toBe("running");
expect(snapshots.some((value) => value.state === "completed")).toBe(false);
});
it("does not persist final completion when ownership changes during cache refresh", async () => {
const snapshots: ImportJob[] = [];
mocks.save.mockImplementation(async (value: ImportJob) => {
snapshots.push(structuredClone(value));
});
mocks.updateCatalog.mockImplementationOnce(async () => {
mocks.eval.mockResolvedValue(0);
});
await drainFurnitureImports();
expect(snapshots.at(-1)?.state).toBe("running");
expect(snapshots.at(-1)?.items[0].state).toBe("done");
expect(snapshots.some((value) => value.state === "completed")).toBe(false);
expect(mocks.updateItems).not.toHaveBeenCalled();
});
it("does not save a translated outcome after ownership is lost during translation", async () => {
const snapshots: ImportJob[] = [];
job.translate = true;
mocks.save.mockImplementation(async (value: ImportJob) => {
snapshots.push(structuredClone(value));
});
mocks.translate.mockImplementationOnce(async () => {
mocks.eval.mockResolvedValue(0);
});
await drainFurnitureImports();
expect(mocks.translate).toHaveBeenCalledOnce();
expect(snapshots).toHaveLength(2);
expect(snapshots.at(-1)?.items[0].state).toBe("running");
expect(mocks.updateCatalog).not.toHaveBeenCalled();
});
+76 -18
View File
@@ -1,4 +1,5 @@
import { randomUUID } from "node:crypto";
import type { ImportJob } from "@/lib/furni/import-job";
import { redis } from "@/lib/redis";
import { logServerError } from "@/lib/server-log";
import { logAudit } from "./audit";
@@ -13,6 +14,7 @@ import { repairFurniture } from "./furniture-repair";
import { rcon } from "./rcon";
const LOCK = "furniture-import-worker:v1";
class ImportLeaseLostError extends Error {}
let running: Promise<void> | undefined;
export function drainFurnitureImports(): Promise<void> {
if (running) return running;
@@ -30,6 +32,28 @@ async function drain() {
const token = randomUUID();
if ((await redis.set(LOCK, token, "EX", 600, "NX")) !== "OK") return;
let lease = true;
async function requireLease() {
if (!lease) throw new ImportLeaseLostError("Import worker lost its lease");
try {
const owned = await 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,
);
if (owned !== 1 || !lease) {
lease = false;
throw new ImportLeaseLostError("Import worker lost its lease");
}
} catch (error) {
lease = false;
if (error instanceof ImportLeaseLostError) throw error;
throw new ImportLeaseLostError(
"Import worker could not verify its lease",
{ cause: error },
);
}
}
const timer = setInterval(() => {
void redis
?.eval(
@@ -47,29 +71,48 @@ async function drain() {
}, 30000);
try {
const store = new ImportJobStore();
// Redis ownership checks and file writes are separate operations. These checks
// prevent known stale writes; they do not provide atomic filesystem fencing.
async function saveOwned(job: ImportJob) {
await requireLease();
await store.save(job);
}
for (const job of await store.list()) {
if (!lease) break;
await requireLease();
if (job.state === "running") {
job.state = "interrupted";
for (const item of job.items)
if (item.state === "running" || item.state === "pending") {
item.state = "failed";
if (item.state === "running") {
item.state = "interrupted";
item.error =
"Server restarted during import. Review local data and retry to complete missing parts.";
"Server restarted during import. Its outcome is uncertain. Inspect local data before starting a new repair.";
}
await store.save(job);
await saveOwned(job);
continue;
}
if (job.state !== "queued") continue;
if (await store.isCancellationRequested(job.id)) {
job.state = "cancelled";
for (const item of job.items)
if (item.state === "pending") item.state = "cancelled";
await saveOwned(job);
continue;
}
job.state = "running";
await store.save(job);
await saveOwned(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");
await requireLease();
if (item.state !== "pending") continue;
if (await store.isCancellationRequested(job.id)) {
for (const remaining of job.items)
if (remaining.state === "pending") remaining.state = "cancelled";
break;
}
item.state = "running";
await store.save(job);
await saveOwned(job);
try {
if (job.sourceId && !source)
throw Error("Import source no longer exists");
@@ -80,6 +123,7 @@ async function drain() {
job.userId,
)
: undefined;
await requireLease();
const result = await (job.mode === "repair"
? repairFurniture
: importSingleFurni)({
@@ -90,6 +134,7 @@ async function drain() {
nitroBaseUrl: source?.nitroBaseUrl,
iconBaseUrl: source?.iconBaseUrl,
});
await requireLease();
item.warnings = result.warnings;
item.itemId = result.itemId;
if (!result.ok) throw Error(result.error || "Import failed");
@@ -114,16 +159,21 @@ async function drain() {
);
}
} catch (error) {
if (error instanceof ImportLeaseLostError) throw error;
await requireLease();
item.state = "failed";
item.error =
error instanceof Error ? error.message : "Import failed";
}
await store.save(job);
await saveOwned(job);
}
await requireLease();
try {
await rcon.updateCatalog();
await requireLease();
await rcon.updateItems();
} catch {
} catch (error) {
if (error instanceof ImportLeaseLostError) throw error;
for (const item of job.items)
if (item.state === "done") {
item.warnings ??= [];
@@ -131,17 +181,25 @@ async function drain() {
}
}
});
job.state = "completed";
await store.save(job);
await requireLease();
job.state = job.items.some((item) => item.state === "cancelled")
? "cancelled"
: "completed";
await saveOwned(job);
}
await requireLease();
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,
);
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,
)
.catch((error) => {
logServerError("furni.worker_release_failed", error);
});
}
}
+6 -2
View File
@@ -12,9 +12,11 @@ import {
UsersSettings,
} from "@/lib/db";
import { resolveHotelName } from "@/lib/hotel-name";
import { loadDashboardEvent } from "./dashboard-event";
import { siteSettings } from "./site-settings";
export async function loadUserDashboard(userId: number) {
const [
dashboardEvent,
userRows,
hotelName,
neededRaw,
@@ -29,6 +31,7 @@ export async function loadUserDashboard(userId: number) {
friends,
currencyRows,
] = await Promise.all([
loadDashboardEvent(),
db
.select({
id: User.id,
@@ -62,7 +65,7 @@ export async function loadUserDashboard(userId: number) {
.from(Rooms)
.orderBy(desc(Rooms.id))
.limit(6)
.catch(() => []),
.catch(() => null),
db
.select({
badgeCode: UsersBadges.badgeCode,
@@ -85,7 +88,7 @@ export async function loadUserDashboard(userId: number) {
.select({ value: count() })
.from(MessengerOffline)
.where(eq(MessengerOffline.userId, userId))
.catch(() => [{ value: 0 }]),
.catch(() => null),
db
.select({ referralsTotal: UserReferrals.referralsTotal })
.from(UserReferrals)
@@ -144,6 +147,7 @@ export async function loadUserDashboard(userId: number) {
];
return {
dashboardEvent,
userRows,
hotelName,
neededRaw,