feat(import): durable batch checkpoint and transient auto-retry
CI / check (push) Successful in 2m24s
CI / deploy (push) Successful in 1m32s
CI / publish-container (push) Successful in 47s

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.
This commit is contained in:
openhands committed 2026-09-13 12:50:17 +02:00
1 parent af1a06b0a3
commit ffcd4232d6
4 files changed
+263 -26

No files matched your search

+143 -24
View File
@@ -1,5 +1,8 @@
import { randomUUID } from "node:crypto";
import { apiError } from "@/lib/api"; import { apiError } from "@/lib/api";
import { withAdmin } from "@/lib/api-handler"; 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 { PERMS } from "@/lib/permissions";
import { logAudit } from "@/lib/services/audit"; import { logAudit } from "@/lib/services/audit";
import { getSource } from "@/lib/services/clone-sources"; import { getSource } from "@/lib/services/clone-sources";
@@ -20,6 +23,8 @@ import {
verifyAndFixInteractionModesCount, verifyAndFixInteractionModesCount,
} from "@/lib/services/furni-import"; } from "@/lib/services/furni-import";
import { clearFurniImportCache } from "@/lib/services/furni-import-cache"; 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 { rcon } from "@/lib/services/rcon";
import type { ImportSingleResult } from "@/types/furni"; import type { ImportSingleResult } from "@/types/furni";
@@ -68,6 +73,49 @@ export const POST = withAdmin(
await ensureDirectories(); await ensureDirectories();
const encoder = new TextEncoder(); 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 // Wire client disconnect to abort controller so we stop processing
// when the user navigates away or closes the browser. // when the user navigates away or closes the browser.
const abortController = new AbortController(); 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<void> = 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(); const startTime = Date.now();
send({ type: "batch_start", total: items.length, concurrency }); send({ type: "batch_start", total: items.length, concurrency });
@@ -116,37 +203,47 @@ export const POST = withAdmin(
index, index,
}); });
const entry = mirrorEntry(item.classname);
if (entry) {
entry.state = "running";
checkpoint();
}
try { try {
// Coalesce micro-step progress events to at most one per // Coalesce micro-step progress events to at most one per
// ~120ms per item so large imports don't flood the client // ~120ms per item so large imports don't flood the client
// (terminal states are always emitted by importSingleFurni // (terminal states are always emitted by importSingleFurni
// and sent below). // and sent below).
let lastProgressSent = 0; let lastProgressSent = 0;
const result: ImportSingleResult = await importSingleFurni({ const result: ImportSingleResult = await withTransientImportRetry(
id: item.id ?? 0, () =>
classname: item.classname, importSingleFurni({
name: item.name, id: item.id ?? 0,
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, classname: item.classname,
status, name: item.name,
index, 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) { if (result.ok) {
succeeded++; succeeded++;
@@ -154,6 +251,14 @@ export const POST = withAdmin(
if (result.furniDataEntry) if (result.furniDataEntry)
furniDataEntries.push(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({ send({
type: "item_progress", type: "item_progress",
classname: item.classname, classname: item.classname,
@@ -180,6 +285,11 @@ export const POST = withAdmin(
}); });
} else { } else {
failed++; failed++;
if (entry) {
entry.state = "failed";
entry.error = result.error;
checkpoint();
}
send({ send({
type: "item_progress", type: "item_progress",
classname: item.classname, classname: item.classname,
@@ -190,6 +300,11 @@ export const POST = withAdmin(
} }
} catch (err) { } catch (err) {
failed++; failed++;
if (entry) {
entry.state = "failed";
entry.error = (err as Error).message;
checkpoint();
}
send({ send({
type: "item_progress", type: "item_progress",
classname: item.classname, classname: item.classname,
@@ -210,6 +325,7 @@ export const POST = withAdmin(
} }
if (aborted) { if (aborted) {
await finalizeMirror("interrupted");
send({ send({
type: "batch_complete", type: "batch_complete",
succeeded, succeeded,
@@ -336,6 +452,8 @@ export const POST = withAdmin(
const { catalogNameFixed, haveOfferFixed, costCreditsFixed } = const { catalogNameFixed, haveOfferFixed, costCreditsFixed } =
await fixDatabaseConsistencyAfterImport(); await fixDatabaseConsistencyAfterImport();
await finalizeMirror("completed");
send({ send({
type: "batch_complete", type: "batch_complete",
succeeded, succeeded,
@@ -359,6 +477,7 @@ export const POST = withAdmin(
duration: Date.now() - startTime, duration: Date.now() - startTime,
}); });
} else { } else {
await finalizeMirror("interrupted");
send({ send({
type: "batch_complete", type: "batch_complete",
succeeded, succeeded,
@@ -933,10 +933,12 @@ export function StudioClient({
if ((err as Error).name === "AbortError") { if ((err as Error).name === "AbortError") {
toast.info("Import cancelled"); toast.info("Import cancelled");
} else if ((err as Error).name === "TimeoutError") { } 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")) { } else if (err instanceof TypeError && err.message?.includes("fetch")) {
toast.error( 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 { } else {
toast.error( toast.error(
@@ -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);
});
});
@@ -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<T>, retries = 2): Promise<T> {
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);
}
}