fix: prevent import hanging by adding abort signals and idle timeouts
- Furni batch: wire request.signal to abort controller, send post_import progress events, skip post-import steps when aborted - Clone sync-all: pre-fetch furnidata once per source instead of per item (eliminates 2000+ redundant DB reads + HTTP requests) - SSE client: add 60s idle timeout to prevent infinite hangs when server stops responding
This commit is contained in:
1 parent
96c33efced
commit
f70d96b81c
3 files changed
+179
-85
No files matched your search
@@ -6,6 +6,7 @@ import { logAudit } from "@/lib/services/audit";
|
||||
import {
|
||||
cloneSingleFurni,
|
||||
fetchSourceFurnidata,
|
||||
type SourceFurni,
|
||||
} from "@/lib/services/clone-import";
|
||||
import { listSources } from "@/lib/services/clone-sources";
|
||||
import {
|
||||
@@ -39,6 +40,20 @@ export const POST = withAdmin(
|
||||
);
|
||||
}
|
||||
|
||||
// Pre-fetch furnidata ONCE per source and build a lookup map.
|
||||
// The old code re-fetched inside the worker loop for every item,
|
||||
// causing 2000+ redundant DB reads + HTTP requests.
|
||||
const sourceFurniDataMap = new Map<string, Map<string, SourceFurni>>();
|
||||
|
||||
for (const source of sources) {
|
||||
try {
|
||||
const entries = await fetchSourceFurnidata(source.furnidataUrl);
|
||||
const byClassname = new Map(entries.map((e) => [e.classname, e]));
|
||||
sourceFurniDataMap.set(source.id, byClassname);
|
||||
} catch {}
|
||||
}
|
||||
|
||||
// Collect all missing items using pre-fetched data
|
||||
const allItems: Array<{
|
||||
classname: string;
|
||||
sourceName: string;
|
||||
@@ -46,14 +61,10 @@ export const POST = withAdmin(
|
||||
}> = [];
|
||||
|
||||
for (const source of sources) {
|
||||
let entries: Awaited<ReturnType<typeof fetchSourceFurnidata>>;
|
||||
try {
|
||||
entries = await fetchSourceFurnidata(source.furnidataUrl);
|
||||
} catch {
|
||||
continue;
|
||||
}
|
||||
const byClassname = sourceFurniDataMap.get(source.id);
|
||||
if (!byClassname) continue;
|
||||
|
||||
const classnames = entries.map((e) => e.classname);
|
||||
const classnames = [...byClassname.keys()];
|
||||
if (classnames.length === 0) continue;
|
||||
|
||||
const rows = await db
|
||||
@@ -62,10 +73,10 @@ export const POST = withAdmin(
|
||||
.where(inArray(ItemsBase.itemName, classnames));
|
||||
const have = new Set(rows.map((r) => r.item_name));
|
||||
|
||||
for (const entry of entries) {
|
||||
if (!have.has(entry.classname)) {
|
||||
for (const classname of classnames) {
|
||||
if (!have.has(classname)) {
|
||||
allItems.push({
|
||||
classname: entry.classname,
|
||||
classname,
|
||||
sourceName: source.name,
|
||||
sourceId: source.id,
|
||||
});
|
||||
@@ -106,29 +117,30 @@ export const POST = withAdmin(
|
||||
itemType: string;
|
||||
}> = [];
|
||||
|
||||
// Source lookup cache — pre-fetched above, no need to re-fetch per item
|
||||
const sourceMap = new Map(sources.map((s) => [s.id, s]));
|
||||
|
||||
return runSseBatch({
|
||||
items: allItems,
|
||||
concurrency: 10,
|
||||
signal: request.signal,
|
||||
labelOf: (it) => `${it.sourceName}/${it.classname}`,
|
||||
worker: async (it, _index, report) => {
|
||||
const allSrc = await listSources();
|
||||
const source = allSrc.find((s) => s.id === it.sourceId);
|
||||
const source = sourceMap.get(it.sourceId);
|
||||
if (!source) {
|
||||
return { ok: false, error: "source not found" };
|
||||
}
|
||||
|
||||
let allEntries: Awaited<ReturnType<typeof fetchSourceFurnidata>>;
|
||||
try {
|
||||
allEntries = await fetchSourceFurnidata(source.furnidataUrl);
|
||||
} catch (err) {
|
||||
// Use pre-fetched furnidata instead of re-fetching
|
||||
const byClassname = sourceFurniDataMap.get(it.sourceId);
|
||||
if (!byClassname) {
|
||||
return {
|
||||
ok: false,
|
||||
error: `furnidata fetch failed: ${(err as Error).message}`,
|
||||
error: "furnidata not available for this source",
|
||||
};
|
||||
}
|
||||
|
||||
const entry = allEntries.find((e) => e.classname === it.classname);
|
||||
const entry = byClassname.get(it.classname);
|
||||
if (!entry) {
|
||||
return {
|
||||
ok: false,
|
||||
|
||||
@@ -67,9 +67,19 @@ export const POST = withAdmin(
|
||||
await ensureDirectories();
|
||||
const encoder = new TextEncoder();
|
||||
|
||||
// Wire client disconnect to abort controller so we stop processing
|
||||
// when the user navigates away or closes the browser.
|
||||
const abortController = new AbortController();
|
||||
let aborted = false;
|
||||
request.signal.addEventListener("abort", () => {
|
||||
aborted = true;
|
||||
abortController.abort();
|
||||
});
|
||||
|
||||
const stream = new ReadableStream({
|
||||
async start(controller) {
|
||||
const send = (data: unknown) => {
|
||||
if (aborted) return;
|
||||
try {
|
||||
controller.enqueue(
|
||||
encoder.encode(`data: ${JSON.stringify(data)}\n\n`),
|
||||
@@ -92,6 +102,8 @@ export const POST = withAdmin(
|
||||
|
||||
// Process in chunks of `concurrency`
|
||||
for (let i = 0; i < items.length; i += concurrency) {
|
||||
if (aborted) break;
|
||||
|
||||
const chunk = items.slice(i, i + concurrency);
|
||||
|
||||
const promises = chunk.map(async (item, chunkIdx) => {
|
||||
@@ -182,6 +194,22 @@ export const POST = withAdmin(
|
||||
await new Promise((r) => setImmediate(r));
|
||||
}
|
||||
|
||||
if (aborted) {
|
||||
send({
|
||||
type: "batch_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
warnings: withWarnings,
|
||||
aborted: true,
|
||||
duration: Date.now() - startTime,
|
||||
});
|
||||
controller.close();
|
||||
return;
|
||||
}
|
||||
|
||||
// Post-import steps: send progress so the user knows what's happening
|
||||
send({ type: "post_import", status: "Writing FurnitureData.json..." });
|
||||
|
||||
// Batch write FurnitureData.json once at the end — BEFORE RCON so the
|
||||
// emulator never sees items that the Nitro client can't render.
|
||||
let furniDataWriteOk = true;
|
||||
@@ -201,8 +229,11 @@ export const POST = withAdmin(
|
||||
}
|
||||
|
||||
// Refresh game server caches once — skip if FurnitureData write failed
|
||||
// to avoid an inconsistent state (emulator sees items, client doesn't).
|
||||
if (furniDataWriteOk) {
|
||||
if (furniDataWriteOk && !aborted) {
|
||||
send({
|
||||
type: "post_import",
|
||||
status: "Refreshing emulator caches...",
|
||||
});
|
||||
try {
|
||||
const okCat = await rcon.updateCatalog();
|
||||
const okItems = await rcon.updateItems();
|
||||
@@ -223,77 +254,106 @@ export const POST = withAdmin(
|
||||
}
|
||||
}
|
||||
|
||||
// Force offer_id to its row id and fix asset ownership across the tree
|
||||
let offerIdsFixed = 0;
|
||||
let ownershipFixed: string[] = [];
|
||||
let nitrosSynced: string[] = [];
|
||||
let iconsSynced: string[] = [];
|
||||
let furniDataFixedIds = 0;
|
||||
let furniDataFixedOfferIds = 0;
|
||||
let furniDataMissing = 0;
|
||||
let furniDataConflicts = 0;
|
||||
let offerRebuildFixed = 0;
|
||||
let spriteFixed = 0;
|
||||
let interactionFixed = 0;
|
||||
let languages: Awaited<
|
||||
ReturnType<typeof buildLocalizedFurniDataFiles>
|
||||
> = [];
|
||||
try {
|
||||
offerIdsFixed = (await reconcileImportedOfferIds()).fixed;
|
||||
ownershipFixed = (await ensureFurniOwnership()).fixed;
|
||||
const assetsSynced = await syncAssetsToGamedataBundle();
|
||||
nitrosSynced = assetsSynced.copiedNitros;
|
||||
iconsSynced = assetsSynced.copiedIcons;
|
||||
const furniReconcile = await reconcileFurniDataWithItemsBase();
|
||||
furniDataFixedIds = furniReconcile.fixedIds;
|
||||
furniDataFixedOfferIds = furniReconcile.fixedOfferIds;
|
||||
furniDataMissing = furniReconcile.missing;
|
||||
furniDataConflicts = furniReconcile.conflicts;
|
||||
offerRebuildFixed = (await rebuildCatalogOfferIds()).fixed;
|
||||
spriteFixed = (await verifyAndFixSpriteIds()).fixed;
|
||||
interactionFixed = (await verifyAndFixInteractionModesCount()).fixed;
|
||||
if (!aborted) {
|
||||
send({
|
||||
type: "post_import",
|
||||
status: "Reconciling offer IDs & ownership...",
|
||||
});
|
||||
|
||||
// Post-import reconciliation — wrapped in try/catch so one
|
||||
// failure doesn't block the rest.
|
||||
let offerIdsFixed = 0;
|
||||
let ownershipFixed: string[] = [];
|
||||
let nitrosSynced: string[] = [];
|
||||
let iconsSynced: string[] = [];
|
||||
let furniDataFixedIds = 0;
|
||||
let furniDataFixedOfferIds = 0;
|
||||
let furniDataMissing = 0;
|
||||
let furniDataConflicts = 0;
|
||||
let offerRebuildFixed = 0;
|
||||
let spriteFixed = 0;
|
||||
let interactionFixed = 0;
|
||||
let languages: Awaited<
|
||||
ReturnType<typeof buildLocalizedFurniDataFiles>
|
||||
> = [];
|
||||
try {
|
||||
languages = await buildLocalizedFurniDataFiles();
|
||||
offerIdsFixed = (await reconcileImportedOfferIds()).fixed;
|
||||
ownershipFixed = (await ensureFurniOwnership()).fixed;
|
||||
const assetsSynced = await syncAssetsToGamedataBundle();
|
||||
nitrosSynced = assetsSynced.copiedNitros;
|
||||
iconsSynced = assetsSynced.copiedIcons;
|
||||
const furniReconcile = await reconcileFurniDataWithItemsBase();
|
||||
furniDataFixedIds = furniReconcile.fixedIds;
|
||||
furniDataFixedOfferIds = furniReconcile.fixedOfferIds;
|
||||
furniDataMissing = furniReconcile.missing;
|
||||
furniDataConflicts = furniReconcile.conflicts;
|
||||
offerRebuildFixed = (await rebuildCatalogOfferIds()).fixed;
|
||||
spriteFixed = (await verifyAndFixSpriteIds()).fixed;
|
||||
|
||||
if (!aborted) {
|
||||
send({
|
||||
type: "post_import",
|
||||
status: "Verifying interaction modes...",
|
||||
});
|
||||
interactionFixed = (await verifyAndFixInteractionModesCount())
|
||||
.fixed;
|
||||
}
|
||||
|
||||
try {
|
||||
languages = await buildLocalizedFurniDataFiles();
|
||||
} catch (err) {
|
||||
console.warn(
|
||||
"[import-furni] Localized furnidata build failed:",
|
||||
(err as Error).message,
|
||||
);
|
||||
}
|
||||
} catch (err) {
|
||||
console.warn(
|
||||
"[import-furni] Localized furnidata build failed:",
|
||||
"[import-furni] Post-import reconcile/ownership failed:",
|
||||
(err as Error).message,
|
||||
);
|
||||
}
|
||||
} catch (err) {
|
||||
console.warn(
|
||||
"[import-furni] Post-import reconcile/ownership failed:",
|
||||
(err as Error).message,
|
||||
);
|
||||
|
||||
send({
|
||||
type: "post_import",
|
||||
status: "Fixing database consistency...",
|
||||
});
|
||||
const { catalogNameFixed, haveOfferFixed, costCreditsFixed } =
|
||||
await fixDatabaseConsistencyAfterImport();
|
||||
|
||||
send({
|
||||
type: "batch_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
warnings: withWarnings,
|
||||
offerIdsFixed,
|
||||
ownershipFixed,
|
||||
nitrosSynced,
|
||||
iconsSynced,
|
||||
furniDataFixedIds,
|
||||
furniDataFixedOfferIds,
|
||||
furniDataMissing,
|
||||
furniDataConflicts,
|
||||
offerRebuildFixed,
|
||||
spriteFixed,
|
||||
interactionFixed,
|
||||
catalogNameFixed,
|
||||
haveOfferFixed,
|
||||
costCreditsFixed,
|
||||
languages,
|
||||
duration: Date.now() - startTime,
|
||||
});
|
||||
} else {
|
||||
send({
|
||||
type: "batch_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
warnings: withWarnings,
|
||||
aborted: true,
|
||||
duration: Date.now() - startTime,
|
||||
});
|
||||
}
|
||||
|
||||
// Fix database consistency (catalog_name + have_offer + cost_credits) after import
|
||||
const { catalogNameFixed, haveOfferFixed, costCreditsFixed } =
|
||||
await fixDatabaseConsistencyAfterImport();
|
||||
|
||||
send({
|
||||
type: "batch_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
warnings: withWarnings,
|
||||
offerIdsFixed,
|
||||
ownershipFixed,
|
||||
nitrosSynced,
|
||||
iconsSynced,
|
||||
furniDataFixedIds,
|
||||
furniDataFixedOfferIds,
|
||||
furniDataMissing,
|
||||
furniDataConflicts,
|
||||
offerRebuildFixed,
|
||||
spriteFixed,
|
||||
interactionFixed,
|
||||
catalogNameFixed,
|
||||
haveOfferFixed,
|
||||
costCreditsFixed,
|
||||
languages,
|
||||
duration: Date.now() - startTime,
|
||||
});
|
||||
|
||||
controller.close();
|
||||
},
|
||||
});
|
||||
|
||||
+24
-2
@@ -3,26 +3,46 @@ import { adminFetch } from "@/lib/admin-fetch";
|
||||
|
||||
export type SseEvent = Record<string, unknown>;
|
||||
|
||||
/** Default idle timeout: if no SSE event arrives for this long, abort. */
|
||||
const DEFAULT_IDLE_TIMEOUT_MS = 60_000;
|
||||
|
||||
/**
|
||||
* Read an SSE response body and invoke `onEvent` for each `data:` JSON payload.
|
||||
* Aborts if no event arrives within `idleTimeoutMs` (default 60s).
|
||||
*/
|
||||
export async function readSseStream(
|
||||
body: ReadableStream<Uint8Array>,
|
||||
onEvent: (event: SseEvent) => void,
|
||||
signal?: AbortSignal,
|
||||
idleTimeoutMs: number = DEFAULT_IDLE_TIMEOUT_MS,
|
||||
): Promise<void> {
|
||||
const reader = body.getReader();
|
||||
const decoder = new TextDecoder();
|
||||
let buf = "";
|
||||
let idleTimer: ReturnType<typeof setTimeout> | undefined;
|
||||
let done = false;
|
||||
|
||||
const resetIdleTimer = () => {
|
||||
if (idleTimer) clearTimeout(idleTimer);
|
||||
if (idleTimeoutMs > 0) {
|
||||
idleTimer = setTimeout(() => {
|
||||
if (!done) {
|
||||
reader.cancel("SSE idle timeout").catch(() => {});
|
||||
}
|
||||
}, idleTimeoutMs);
|
||||
}
|
||||
};
|
||||
|
||||
try {
|
||||
resetIdleTimer();
|
||||
while (true) {
|
||||
if (signal?.aborted) {
|
||||
await reader.cancel();
|
||||
break;
|
||||
}
|
||||
const { value, done } = await reader.read();
|
||||
if (done) break;
|
||||
const { value, done: streamDone } = await reader.read();
|
||||
if (streamDone) break;
|
||||
resetIdleTimer();
|
||||
buf += decoder.decode(value, { stream: true });
|
||||
const parts = buf.split("\n\n");
|
||||
buf = parts.pop() ?? "";
|
||||
@@ -36,6 +56,8 @@ export async function readSseStream(
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
done = true;
|
||||
if (idleTimer) clearTimeout(idleTimer);
|
||||
reader.releaseLock();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user