From ffcd4232d668c8c9d7f6d40d425c933207a8dcce Mon Sep 17 00:00:00 2001 From: openhands Date: Sun, 13 Sep 2026 12:50:17 +0200 Subject: [PATCH] feat(import): durable batch checkpoint and transient auto-retry Mirror interactive batch runs into the import-job store so interrupted imports (restart, time-out, disconnect) can be resumed from Import History. Items are checkpointed as they settle (coalesced, serialized saves) and the mirror starts 'running' so the boot-time worker marks it 'interrupted' instead of double-importing; done items are never re-imported. Add bounded backoff retry for transient download/connection failures before marking an item failed, and point the client's time-out/network toasts at Import History. --- src/app/api/admin/import/furni/batch/route.ts | 167 +++++++++++++++--- src/components/admin/studio/studio-client.tsx | 6 +- .../services/import/transient-retry.test.ts | 91 ++++++++++ src/lib/services/import/transient-retry.ts | 25 +++ 4 files changed, 263 insertions(+), 26 deletions(-) create mode 100644 src/lib/services/import/transient-retry.test.ts create mode 100644 src/lib/services/import/transient-retry.ts diff --git a/src/app/api/admin/import/furni/batch/route.ts b/src/app/api/admin/import/furni/batch/route.ts index 5cbfd36d..e7c719bb 100644 --- a/src/app/api/admin/import/furni/batch/route.ts +++ b/src/app/api/admin/import/furni/batch/route.ts @@ -1,5 +1,8 @@ +import { randomUUID } from "node:crypto"; import { apiError } from "@/lib/api"; import { withAdmin } from "@/lib/api-handler"; +import { getRequestId } from "@/lib/foundation/request-context"; +import type { ImportJob } from "@/lib/furni/import-job"; import { PERMS } from "@/lib/permissions"; import { logAudit } from "@/lib/services/audit"; import { getSource } from "@/lib/services/clone-sources"; @@ -20,6 +23,8 @@ import { verifyAndFixInteractionModesCount, } from "@/lib/services/furni-import"; import { clearFurniImportCache } from "@/lib/services/furni-import-cache"; +import { ImportJobStore } from "@/lib/services/furni-job-store"; +import { withTransientImportRetry } from "@/lib/services/import/transient-retry"; import { rcon } from "@/lib/services/rcon"; import type { ImportSingleResult } from "@/types/furni"; @@ -68,6 +73,49 @@ export const POST = withAdmin( await ensureDirectories(); const encoder = new TextEncoder(); + // Durable checkpoint: mirror this run into the import-job store so an + // interrupted batch can be resumed from Import History after a restart, + // time-out or disconnect. The job starts "running" so the boot-time + // worker drain marks it "interrupted" instead of double-importing. + const store = new ImportJobStore(); + const createdAt = new Date().toISOString(); + const mirrorJob: ImportJob | null = await (async () => { + try { + const job: ImportJob = { + id: randomUUID(), + userId: ctx.session.user.id, + operationId: getRequestId(), + createdAt, + updatedAt: createdAt, + state: "running", + sourceId, + translate: body.translate === true, + langs: + Array.isArray(body.langs) && body.langs.length > 0 + ? body.langs + : undefined, + items: items.map((item) => ({ + id: item.id ?? 0, + classname: item.classname, + name: item.name, + description: item.description ?? "", + type: item.type === "wallitem" ? "wallitem" : "flooritem", + revision: item.revision ?? 0, + category: item.category ?? "unknown", + state: "pending", + })), + }; + await store.save(job); + return job; + } catch (error) { + console.warn( + "[import-furni] Could not create durable batch job", + error, + ); + return null; + } + })(); + // Wire client disconnect to abort controller so we stop processing // when the user navigates away or closes the browser. const abortController = new AbortController(); @@ -90,6 +138,45 @@ export const POST = withAdmin( } }; + // Checkpoint writes are coalesced and chained so concurrent item + // settles can't interleave saves on the same job file; the final + // state is always flushed by finalizeMirror. + const mirrorEntry = (classname: string) => + mirrorJob?.items.find((m) => m.classname === classname); + let saveChain: Promise = Promise.resolve(); + let saveQueued = false; + const checkpoint = () => { + if (!mirrorJob || saveQueued) return; + saveQueued = true; + saveChain = saveChain.then(async () => { + try { + await store.save(mirrorJob); + } catch (error) { + console.warn("[import-furni] Checkpoint write failed", error); + } finally { + saveQueued = false; + } + }); + }; + const finalizeMirror = async (state: "completed" | "interrupted") => { + if (!mirrorJob) return; + if (state === "interrupted") + for (const entry of mirrorJob.items) + if (entry.state === "running") { + entry.state = "interrupted"; + entry.error = + "Import stream was interrupted before finishing. Check imported data before resuming it from history."; + } + mirrorJob.state = state; + mirrorJob.updatedAt = new Date().toISOString(); + await saveChain; + try { + await store.save(mirrorJob); + } catch (error) { + console.warn("[import-furni] Final checkpoint write failed", error); + } + }; + const startTime = Date.now(); send({ type: "batch_start", total: items.length, concurrency }); @@ -116,37 +203,47 @@ export const POST = withAdmin( index, }); + const entry = mirrorEntry(item.classname); + if (entry) { + entry.state = "running"; + checkpoint(); + } + try { // Coalesce micro-step progress events to at most one per // ~120ms per item so large imports don't flood the client // (terminal states are always emitted by importSingleFurni // and sent below). let lastProgressSent = 0; - const result: ImportSingleResult = await importSingleFurni({ - id: item.id ?? 0, - classname: item.classname, - name: item.name, - description: item.description ?? "", - type: item.type ?? "flooritem", - revision: item.revision ?? 0, - category: item.category ?? "unknown", - skipFurniDataWrite: false, - repairExisting: body.repairExisting === true, - sourceSwfBaseUrl: source?.sourceSwfBaseUrl, - nitroBaseUrl: source?.nitroBaseUrl, - iconBaseUrl: source?.iconBaseUrl, - onProgress: (status: string) => { - const now = Date.now(); - if (now - lastProgressSent < 120) return; - lastProgressSent = now; - send({ - type: "item_progress", + const result: ImportSingleResult = await withTransientImportRetry( + () => + importSingleFurni({ + id: item.id ?? 0, classname: item.classname, - status, - index, - }); - }, - }); + name: item.name, + description: item.description ?? "", + type: item.type ?? "flooritem", + revision: item.revision ?? 0, + category: item.category ?? "unknown", + skipFurniDataWrite: false, + repairExisting: body.repairExisting === true, + sourceSwfBaseUrl: source?.sourceSwfBaseUrl, + nitroBaseUrl: source?.nitroBaseUrl, + iconBaseUrl: source?.iconBaseUrl, + onProgress: (status: string) => { + const now = Date.now(); + if (now - lastProgressSent < 120) return; + lastProgressSent = now; + send({ + type: "item_progress", + classname: item.classname, + status, + index, + }); + }, + }), + 2, + ); if (result.ok) { succeeded++; @@ -154,6 +251,14 @@ export const POST = withAdmin( if (result.furniDataEntry) furniDataEntries.push(result.furniDataEntry); + if (entry) { + entry.state = "done"; + entry.itemId = result.itemId; + entry.warnings = + result.warnings.length > 0 ? result.warnings : undefined; + checkpoint(); + } + send({ type: "item_progress", classname: item.classname, @@ -180,6 +285,11 @@ export const POST = withAdmin( }); } else { failed++; + if (entry) { + entry.state = "failed"; + entry.error = result.error; + checkpoint(); + } send({ type: "item_progress", classname: item.classname, @@ -190,6 +300,11 @@ export const POST = withAdmin( } } catch (err) { failed++; + if (entry) { + entry.state = "failed"; + entry.error = (err as Error).message; + checkpoint(); + } send({ type: "item_progress", classname: item.classname, @@ -210,6 +325,7 @@ export const POST = withAdmin( } if (aborted) { + await finalizeMirror("interrupted"); send({ type: "batch_complete", succeeded, @@ -336,6 +452,8 @@ export const POST = withAdmin( const { catalogNameFixed, haveOfferFixed, costCreditsFixed } = await fixDatabaseConsistencyAfterImport(); + await finalizeMirror("completed"); + send({ type: "batch_complete", succeeded, @@ -359,6 +477,7 @@ export const POST = withAdmin( duration: Date.now() - startTime, }); } else { + await finalizeMirror("interrupted"); send({ type: "batch_complete", succeeded, diff --git a/src/components/admin/studio/studio-client.tsx b/src/components/admin/studio/studio-client.tsx index 460d4162..73482336 100644 --- a/src/components/admin/studio/studio-client.tsx +++ b/src/components/admin/studio/studio-client.tsx @@ -933,10 +933,12 @@ export function StudioClient({ if ((err as Error).name === "AbortError") { toast.info("Import cancelled"); } else if ((err as Error).name === "TimeoutError") { - toast.error("Import timed out. Retry or reduce concurrency."); + toast.error( + "Import timed out. Progress was saved — open Import History to resume failed items.", + ); } else if (err instanceof TypeError && err.message?.includes("fetch")) { toast.error( - "Network error during upload. Check your connection and retry.", + "Network error during upload. Progress was saved — open Import History to resume failed items.", ); } else { toast.error( diff --git a/src/lib/services/import/transient-retry.test.ts b/src/lib/services/import/transient-retry.test.ts new file mode 100644 index 00000000..4dc6e262 --- /dev/null +++ b/src/lib/services/import/transient-retry.test.ts @@ -0,0 +1,91 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { + isTransientImportError, + withTransientImportRetry, +} from "./transient-retry"; + +beforeEach(() => { + vi.useFakeTimers(); +}); +afterEach(() => { + vi.useRealTimers(); +}); + +describe("isTransientImportError", () => { + it("flags connection and timeout classes", () => { + for (const message of [ + "fetch failed", + "network error during download", + "service unavailable", + "HTTP 503", + "HTTP 500", + "Empty or truncated response", + "connection timed out", + "ECONNRESET", + ]) { + expect(isTransientImportError(message)).toBe(true); + } + }); + + it("does not flag data errors or aborts", () => { + for (const message of [ + undefined, + "Response is not a SWF file", + "Invalid Nitro bundle", + "This operation was aborted", + "classname missing in assets", + "HTTP 404", + ]) { + expect(isTransientImportError(message)).toBe(false); + } + }); +}); + +describe("withTransientImportRetry", () => { + it("returns success on the first attempt", async () => { + const attempt = vi.fn(async () => ({ ok: true as const, itemId: 1 })); + await expect(withTransientImportRetry(attempt)).resolves.toEqual({ + ok: true, + itemId: 1, + }); + expect(attempt).toHaveBeenCalledTimes(1); + }); + + it("retries transient failures until success", async () => { + const attempt = vi + .fn() + .mockResolvedValueOnce({ ok: false as const, error: "fetch failed" }) + .mockResolvedValueOnce({ ok: false as const, error: "HTTP 503" }) + .mockResolvedValueOnce({ ok: true as const, itemId: 7 }); + const pending = withTransientImportRetry(attempt, 2); + await vi.advanceTimersByTimeAsync(2000); + await expect(pending).resolves.toEqual({ ok: true, itemId: 7 }); + expect(attempt).toHaveBeenCalledTimes(3); + }); + + it("returns the last error after exhausting retries", async () => { + const attempt = vi.fn(async () => ({ + ok: false as const, + error: "fetch failed", + })); + const pending = withTransientImportRetry(attempt, 2); + await vi.advanceTimersByTimeAsync(2000); + await expect(pending).resolves.toEqual({ + ok: false, + error: "fetch failed", + }); + expect(attempt).toHaveBeenCalledTimes(3); + }); + + it("does not retry non-transient failures", async () => { + const attempt = vi.fn(async () => ({ + ok: false as const, + error: "Invalid Nitro bundle", + })); + await expect(withTransientImportRetry(attempt)).resolves.toEqual({ + ok: false, + error: "Invalid Nitro bundle", + }); + expect(attempt).toHaveBeenCalledTimes(1); + }); +}); diff --git a/src/lib/services/import/transient-retry.ts b/src/lib/services/import/transient-retry.ts new file mode 100644 index 00000000..09deef7b --- /dev/null +++ b/src/lib/services/import/transient-retry.ts @@ -0,0 +1,25 @@ +const TRANSIENT_ERROR_PATTERN = + /fetch failed|network|socket|timeout|timed out|ECONNRESET|ETIMEDOUT|ENOTFOUND|EAI_AGAIN|ECONNREFUSED|EHOSTUNREACH|HTTP 5\d\d|Empty or truncated response|service unavailable|temporarily unavailable/i; +const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); + +/** True for download/connection hiccups worth a retry; false for data errors. */ +export function isTransientImportError(message: string | undefined): boolean { + if (!message) return false; + return TRANSIENT_ERROR_PATTERN.test(message); +} + +/** + * Runs `attempt` again (with backoff) when it returns a transient failure. + * Non-transient failures and hard throws are returned/raised immediately. + */ +export async function withTransientImportRetry< + T extends { ok: boolean; error?: string }, +>(attempt: () => Promise, retries = 2): Promise { + for (let n = 0; ; n++) { + const result = await attempt(); + if (result.ok || !isTransientImportError(result.error) || n >= retries) + return result; + const delay = Math.min(400 * 2 ** n, 2000); + await sleep(delay); + } +}