fix(studio): recover uncertain imports with persistent request identities
CI / check (push) Successful in 3m12s
CI / deploy (push) Successful in 19s
CI / publish-container (push) Successful in 1m12s

This commit is contained in:
Simo committed 2026-09-13 18:33:06 +02:00
1 parent afea80708c
commit c977fe95ba
12 files changed
+365 -31

No files matched your search

@@ -1,15 +1,34 @@
import { useState } from "react";
import { FurnitureJobHistory } from "@/components/admin/studio/furniture-jobs";
import { useFurnitureJobs } from "@/components/admin/studio/use-furniture-jobs";
import { furnitureJobFixture } from "../../../src/test/furniture-jobs-fixtures";
export function FurnitureJobsHarness() {
const [completions, setCompletions] = useState(0);
const [submitted, setSubmitted] = useState<string>("");
const history = useFurnitureJobs(
() => setCompletions((value) => value + 1),
true,
7,
);
return (
<section aria-label="Import history hook">
<button
type="button"
disabled={history.busy}
onClick={async () =>
setSubmitted(
String(
await history.submit(furnitureJobFixture.items, {
translate: false,
}),
),
)
}
>
Submit import
</button>
<output aria-label="Submitted">{submitted}</output>
<button
type="button"
onClick={() => {
+43
View File
@@ -143,3 +143,46 @@ test("history shows the current phase and validated recovery source", async ({
),
).toBe(false);
});
test("uncertain import submissions keep their identity after a page reload", async ({
page,
}) => {
const ids: string[] = [];
await page.route("**/api/admin/csrf", (route) =>
route.fulfill({ json: { token: "a".repeat(64) } }),
);
await page.route("**/api/admin/studio/import-jobs**", async (route) => {
if (route.request().method() !== "POST") {
await route.fulfill({ json: { ok: true, jobs: [], nextCursor: null } });
return;
}
const body = route.request().postDataJSON();
ids.push(body.id);
if (ids.length === 1) await route.abort("failed");
else
await route.fulfill({
json: { ok: true, job: { ...furnitureJobFixture, id: body.id } },
});
});
await page.goto("/admin/jobs-harness");
await page
.getByRole("button", { name: "Submit import", exact: true })
.click();
await expect(page.getByLabel("Submitted", { exact: true })).toHaveText(
"false",
);
await page.reload();
await page
.getByRole("button", { name: "Submit import", exact: true })
.click();
await expect(page.getByLabel("Submitted", { exact: true })).toHaveText(
"true",
);
expect(ids).toHaveLength(2);
expect(ids[1]).toBe(ids[0]);
await page
.getByRole("button", { name: "Submit import", exact: true })
.click();
await expect.poll(() => ids.length).toBe(3);
expect(ids[2]).not.toBe(ids[0]);
});
+1
View File
@@ -22,6 +22,7 @@ export default async function StudioFurniPage(_props: {
return (
<StudioClient
actorId={Number(session.user.id)}
source={buildFurniImportSource(hotel)}
initialTree={catalogTree}
defaultTranslate={translateEnabled}
@@ -7,6 +7,7 @@ const mocks = vi.hoisted(() => ({
attachment: vi.fn(),
guard: vi.fn(),
create: vi.fn(),
read: vi.fn(),
list: vi.fn(),
page: vi.fn(),
retry: vi.fn(),
@@ -33,6 +34,7 @@ vi.mock("@/lib/services/furni-job-store", async (original) => ({
...(await original<typeof import("@/lib/services/furni-job-store")>()),
ImportJobStore: class {
create = mocks.create;
read = mocks.read;
list = mocks.list;
page = mocks.page;
retry = mocks.retry;
@@ -64,6 +66,9 @@ const request = (body: unknown) =>
});
beforeEach(() => {
mocks.create.mockClear();
mocks.read.mockRejectedValue(
Object.assign(Error("missing"), { code: "ENOENT" }),
);
mocks.list.mockClear();
mocks.after.mockClear();
mocks.ping.mockResolvedValue("PONG");
@@ -234,3 +239,36 @@ it("rejects invalid recovery attachment paths", async () => {
).toBe(400);
expect(mocks.retry).not.toHaveBeenCalled();
});
it("recovers accepted requests even when Redis and the source are unavailable", async () => {
const body = { id: randomUUID(), sourceId: "source", items: [item] };
mocks.source.mockResolvedValue({ id: "source" });
const accepted = await POST(request(body), ctx);
const saved = (await accepted.json()).job;
mocks.read.mockResolvedValue({ ...saved, state: "completed" });
mocks.ping.mockRejectedValue(Error("offline"));
mocks.source.mockResolvedValue(null);
const replay = await POST(request(body), ctx);
expect(replay.status).toBe(200);
expect((await replay.json()).job).toMatchObject({
id: body.id,
state: "completed",
});
expect(mocks.create).toHaveBeenCalledTimes(1);
});
it("rejects changed content for an accepted request instead of reporting false success", async () => {
const body = { id: randomUUID(), items: [item] };
const accepted = await POST(request(body), ctx);
mocks.read.mockResolvedValue((await accepted.json()).job);
const replay = await POST(request({ ...body, translate: true }), ctx);
expect(replay.status).toBe(409);
expect(mocks.create).toHaveBeenCalledTimes(1);
});
it("does not recover a different owner's request", async () => {
const body = { id: randomUUID(), items: [item] };
const accepted = await POST(request(body), ctx);
mocks.read.mockResolvedValue({ ...(await accepted.json()).job, userId: 99 });
const replay = await POST(request(body), ctx);
expect(replay.status).toBe(409);
expect((await replay.json()).job).toBeUndefined();
});
+48 -16
View File
@@ -1,3 +1,4 @@
import { createHash } from "node:crypto";
import { after } from "next/server";
import { z } from "zod";
import { apiError, apiJson, apiOk } from "@/lib/api";
@@ -8,7 +9,11 @@ import { PERMS } from "@/lib/permission-slugs";
import { redis } from "@/lib/redis";
import { getSource } from "@/lib/services/clone-sources";
import { readFurnitureAttachment } from "@/lib/services/furni-attachment";
import { ImportJobStore, validJobId } from "@/lib/services/furni-job-store";
import {
ImportJobStore,
ImportRequestConflict,
validJobId,
} from "@/lib/services/furni-job-store";
import { drainFurnitureImports } from "@/lib/services/furni-job-worker";
const schema = z.object({
@@ -78,28 +83,55 @@ export const POST = withAdmin(
"Invalid furniture import request (maximum 500 items)",
400,
);
const store = new ImportJobStore(),
body = parsed.data;
const unique = [
...new Map(body.items.map((item) => [item.classname, item])).values(),
];
// Hash validated request fields before worker checkpoints can change the job.
const { id: _id, ...input } = body;
const requestFingerprint = createHash("sha256")
.update(JSON.stringify({ ...input, items: unique }))
.digest("hex");
const existing = await store.read(body.id).catch((error) => {
if (error.code === "ENOENT") return null;
throw error;
});
if (existing) {
// Legacy jobs have no immutable fingerprint: never infer equivalence from progress.
if (
existing.userId !== ctx.session.user.id ||
existing.requestFingerprint !== requestFingerprint
)
return apiError("Request ID already used", 409);
after(drainFurnitureImports);
return apiOk({ job: existing });
}
if (!redis || (await redis.ping().catch(() => null)) !== "PONG")
return apiError(
"Background import queue is temporarily unavailable. Please retry shortly.",
503,
);
if (parsed.data.sourceId && !(await getSource(parsed.data.sourceId)))
if (body.sourceId && !(await getSource(body.sourceId)))
return apiError("Source not found", 404);
const store = new ImportJobStore(),
body = parsed.data;
const createdAt = new Date().toISOString();
const unique = [
...new Map(body.items.map((item) => [item.classname, item])).values(),
];
const job = await store.create({
...body,
userId: ctx.session.user.id,
operationId: getRequestId(),
createdAt,
updatedAt: createdAt,
state: "queued",
items: unique.map((item) => ({ ...item, state: "pending" })),
});
let job: Awaited<ReturnType<ImportJobStore["create"]>>;
try {
job = await store.create({
...body,
requestFingerprint,
userId: ctx.session.user.id,
operationId: getRequestId(),
createdAt,
updatedAt: createdAt,
state: "queued",
items: unique.map((item) => ({ ...item, state: "pending" })),
});
} catch (error) {
if (error instanceof ImportRequestConflict)
return apiError("Request ID already used", 409);
throw error;
}
after(drainFurnitureImports);
return apiOk({ job });
},
@@ -0,0 +1,70 @@
import { expect, it } from "vitest";
import { createFurnitureImportRequests } from "./furniture-import-requests";
function storage() {
const entries = new Map<string, string>();
return {
entries,
getItem: (key: string) => entries.get(key) ?? null,
setItem: (key: string, value: string) => {
entries.set(key, value);
},
removeItem: (key: string) => {
entries.delete(key);
},
};
}
it("reuses an uncertain request across reloads without storing furniture content", async () => {
const saved = storage();
const first = await createFurnitureImportRequests(7, () => saved).get(
'{"name":"private furniture"}',
);
const resumed = await createFurnitureImportRequests(7, () => saved).get(
'{"name":"private furniture"}',
);
expect(resumed.id).toBe(first.id);
expect(JSON.stringify([...saved.entries])).not.toContain("private furniture");
});
it("keeps pending requests distinct across accounts and changed selections", async () => {
const saved = storage();
const client = createFurnitureImportRequests(7, () => saved);
const first = await client.get("chair");
const second = await client.get("lamp");
const other = await createFurnitureImportRequests(8, () => saved).get(
"chair",
);
expect(new Set([first.id, second.id, other.id]).size).toBe(3);
expect((await client.get("chair")).id).toBe(first.id);
});
it("starts a new request only after the previous one was acknowledged", async () => {
const saved = storage();
const client = createFurnitureImportRequests(7, () => saved);
const pending = await client.get("chair");
client.acknowledge(pending);
expect(
(await createFurnitureImportRequests(7, () => saved).get("chair")).id,
).not.toBe(pending.id);
});
it("keeps request identity in memory if session storage is unavailable", async () => {
const client = createFurnitureImportRequests(7, () => {
throw Error("denied");
});
const pending = await client.get("chair");
expect((await client.get("chair")).id).toBe(pending.id);
client.acknowledge(pending);
expect((await client.get("chair")).id).not.toBe(pending.id);
});
it("does not persist an unscoped request and ignores corrupt saved IDs", async () => {
const saved = storage();
await createFurnitureImportRequests(undefined, () => saved).get("chair");
expect(saved.entries.size).toBe(0);
const first = await createFurnitureImportRequests(7, () => saved).get(
"chair",
);
saved.setItem(first.key, "../invalid");
const repaired = await createFurnitureImportRequests(7, () => saved).get(
"chair",
);
expect(repaired.id).not.toBe("../invalid");
expect(saved.getItem(first.key)).toBe(repaired.id);
});
@@ -0,0 +1,59 @@
type RequestIdentity = { key: string; id: string };
type RequestStorage = Pick<Storage, "getItem" | "setItem" | "removeItem">;
/** Store only request identities; furniture payloads and files stay out of browser storage. */
export function createFurnitureImportRequests(
actorId?: number,
storage: () => RequestStorage = () => sessionStorage,
) {
const pending = new Map<string, string>();
const scoped = Number.isSafeInteger(actorId) && Number(actorId) > 0;
return {
async get(payload: string): Promise<RequestIdentity> {
const digest = await crypto.subtle.digest(
"SHA-256",
new TextEncoder().encode(payload),
);
const fingerprint = Array.from(new Uint8Array(digest), (byte) =>
byte.toString(16).padStart(2, "0"),
).join("");
const key = `furniture-import-request:studio:${actorId ?? "memory"}:${fingerprint}`;
let id = pending.get(key);
if (!id && scoped) {
try {
const stored = storage().getItem(key);
if (
stored &&
/^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test(
stored,
)
)
id = stored;
} catch {
/* Keep in-memory recovery when browser storage is unavailable. */
}
}
id ||= crypto.randomUUID();
pending.set(key, id);
if (scoped) {
try {
storage().setItem(key, id);
} catch {
/* The current page can still retry safely. */
}
}
return { key, id };
},
acknowledge(request: RequestIdentity) {
if (pending.get(request.key) === request.id) pending.delete(request.key);
if (scoped) {
try {
if (storage().getItem(request.key) === request.id)
storage().removeItem(request.key);
} catch {
/* A later replay remains safe on the server. */
}
}
},
};
}
+10 -4
View File
@@ -167,10 +167,12 @@ const TABLE_COLS =
const TABLE_ROW_H = 40;
export function StudioClient({
actorId,
source,
initialTree,
defaultTranslate,
}: {
actorId?: number;
source: FurniImportSource;
initialTree: TreeNode[];
defaultTranslate: boolean;
@@ -180,10 +182,14 @@ export function StudioClient({
single: boolean;
cloneAll?: boolean;
} | null>(null);
const jobs = useFurnitureJobs(() => {
void fetchItems(activeSearch, 1, activeSource);
void fetchStats();
});
const jobs = useFurnitureJobs(
() => {
void fetchItems(activeSearch, 1, activeSource);
void fetchStats();
},
false,
actorId,
);
const [preparingReview, setPreparingReview] = useState(false);
// Catalog tree
const [tree] = useState<TreeNode[]>(initialTree);
@@ -1,11 +1,12 @@
"use client";
import { useQuery } from "@tanstack/react-query";
import { useTranslations } from "next-intl";
import { useCallback, useEffect, useRef, useState } from "react";
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { toast } from "sonner";
import { useVisiblePolling } from "@/hooks/use-visible-polling";
import { adminFetch } from "@/lib/admin-fetch";
import type { ImportJobItem } from "@/lib/furni/import-job";
import { createFurnitureImportRequests } from "./furniture-import-requests";
import {
createFurnitureJobsClient,
FurnitureJobsHttpError,
@@ -14,7 +15,11 @@ import {
readFurnitureJobResponse,
} from "./furniture-jobs-query";
export function useFurnitureJobs(onComplete: () => void, paginated = false) {
export function useFurnitureJobs(
onComplete: () => void,
paginated = false,
actorId?: number,
) {
const t = useTranslations("pages.admin.importHistory");
const [client] = useState(createFurnitureJobsClient);
const cursor = useRef<string | null>(null);
@@ -38,7 +43,10 @@ export function useFurnitureJobs(onComplete: () => void, paginated = false) {
const completed = useRef(new Set<string>());
const callback = useRef(onComplete);
callback.current = onComplete;
const request = useRef<{ payload: string; id: string } | null>(null);
const requests = useMemo(
() => createFurnitureImportRequests(actorId),
[actorId],
);
const submitting = useRef(false);
useEffect(
() => () => {
@@ -154,16 +162,21 @@ export function useFurnitureJobs(onComplete: () => void, paginated = false) {
setBusy(true);
const body = { items, ...options },
payload = JSON.stringify(body);
if (request.current?.payload !== payload)
request.current = { payload, id: crypto.randomUUID() };
try {
const request = await requests.get(payload);
const response = await adminFetch("/api/admin/studio/import-jobs", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ ...body, id: request.current.id }),
body: JSON.stringify({ ...body, id: request.id }),
});
await readFurnitureJobResponse(response, "Could not queue import");
request.current = null;
const accepted = await readFurnitureJobResponse(
response,
"Could not queue import",
);
if (accepted.id !== request.id)
throw Error("Could not confirm the queued import");
requests.acknowledge(request);
await client.cancelQueries();
resetPage();
await refresh();
+1
View File
@@ -25,6 +25,7 @@ export interface ImportJob {
mode?: "repair";
syncKind?: "official" | "clone";
operationId?: string;
requestFingerprint?: string;
retryOf?: string;
id: string;
userId: number;
+44 -1
View File
@@ -243,11 +243,12 @@ it("attaches recovery files only to safe failed items and preserves retry histor
).rejects.toThrow("latest retry");
});
it("clears old source provenance and phase when creating a retry", async () => {
it("clears old source provenance, request identity and phase when creating a retry", async () => {
const s = await store();
const original: ImportJob = {
...job(),
state: "completed",
requestFingerprint: "original-request",
items: [
{
id: 1,
@@ -271,7 +272,49 @@ it("clears old source provenance and phase when creating a retry", async () => {
};
await s.create(original);
const child = await s.retry(original.id, 4);
expect(child?.requestFingerprint).toBeUndefined();
expect(child?.items[0].phase).toBeUndefined();
expect(child?.items[0].recoveredSource).toBeUndefined();
expect(child?.items[0].sourceAttempt).toBeUndefined();
});
it("rejects reuse of a request ID with a different immutable fingerprint", async () => {
const s = await store();
const original = { ...job(), requestFingerprint: "original" };
await s.create(original);
await expect(
s.create({ ...original, requestFingerprint: "changed" }),
).rejects.toThrow("Request ID already used");
expect((await s.read(original.id)).requestFingerprint).toBe("original");
});
it("preserves the request identity while checkpoint data changes", async () => {
const s = await store();
const original = { ...job(), requestFingerprint: "original" };
await s.create(original);
await s.save({ ...original, state: "completed" });
expect((await s.create(original)).state).toBe("completed");
});
it("does not guess request identity for legacy jobs without a fingerprint", async () => {
const s = await store();
const original = job();
await s.create(original);
await expect(
s.create({ ...original, requestFingerprint: "new" }),
).rejects.toThrow("Request ID already used");
});
it("rejects concurrent conflicting submissions without overwriting the winner", async () => {
const s = await store();
const original = { ...job(), requestFingerprint: "one" };
const results = await Promise.allSettled([
s.create(original),
s.create({ ...original, requestFingerprint: "two" }),
]);
expect(
results.filter((result) => result.status === "fulfilled"),
).toHaveLength(1);
expect(results.filter((result) => result.status === "rejected")).toHaveLength(
1,
);
expect(await s.list()).toHaveLength(1);
});
+11 -2
View File
@@ -10,6 +10,11 @@ export const validJobId = (value: unknown): value is 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 ImportRequestConflict extends Error {
constructor() {
super("Request ID already used");
}
}
export class ImportJobStore {
constructor(readonly root = importRoot()) {}
async save(job: ImportJob) {
@@ -31,8 +36,11 @@ export class ImportJobStore {
} 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");
if (
existing.userId !== job.userId ||
existing.requestFingerprint !== job.requestFingerprint
)
throw new ImportRequestConflict();
return existing;
} finally {
await fs.unlink(tmp);
@@ -106,6 +114,7 @@ export class ImportJobStore {
...original,
id: retryId,
retryOf: id,
requestFingerprint: undefined,
createdAt: now,
updatedAt: now,
state: "queued",