perf: prevent import hanging with timeouts and batching
- verifyAndFixInteractionModesCount: paginated DB queries (500/batch) instead of loading all items into memory at once - withFurniDataLock: add 60s chain timeout to prevent deadlocks when a lock holder stalls or crashes - withGamedataLock: same timeout protection for gamedata locks - conversion-pool: add 60s per-job timeout, fall back to main thread - furni/batch: abort signal + post_import progress events + skip post-import steps when client disconnects - furni/batch-regen: abort signal + early exit when disconnected - clone/sync-all: pre-fetch furnidata once per source instead of per item (eliminates 2000+ redundant fetches) - sse-client: add 60s idle timeout to prevent infinite hangs
This commit is contained in:
1 parent
f70d96b81c
commit
7556b1f3fa
5 files changed
+222
-93
No files matched your search
@@ -47,9 +47,16 @@ export const POST = withAdmin(
|
||||
const { swfDir, nitroDir } = await getFurniAssetDirs();
|
||||
const encoder = new TextEncoder();
|
||||
|
||||
// Wire client disconnect to abort controller
|
||||
let aborted = false;
|
||||
request.signal.addEventListener("abort", () => {
|
||||
aborted = true;
|
||||
});
|
||||
|
||||
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`),
|
||||
@@ -240,6 +247,18 @@ export const POST = withAdmin(
|
||||
await new Promise((r) => setImmediate(r));
|
||||
}
|
||||
|
||||
if (aborted) {
|
||||
send({
|
||||
type: "regen_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
skipped,
|
||||
aborted: true,
|
||||
});
|
||||
controller.close();
|
||||
return;
|
||||
}
|
||||
|
||||
// Batch write FurnitureData.json
|
||||
if (furniDataEntries.length > 0) {
|
||||
try {
|
||||
@@ -291,25 +310,27 @@ export const POST = withAdmin(
|
||||
);
|
||||
}
|
||||
|
||||
send({
|
||||
type: "regen_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
skipped,
|
||||
furniDataAdded: furniDataEntries.length,
|
||||
ownershipFixed: ownership.fixed,
|
||||
nitrosSynced: assetsSynced.copiedNitros,
|
||||
iconsSynced: assetsSynced.copiedIcons,
|
||||
furniDataFixedIds: furniReconcile.fixedIds,
|
||||
furniDataMissing: furniReconcile.missing,
|
||||
furniDataConflicts: furniReconcile.conflicts,
|
||||
offerRebuildChecked: offerRebuild.checked,
|
||||
offerRebuildFixed: offerRebuild.fixed,
|
||||
spriteChecked: spriteVerify.checked,
|
||||
spriteFixed: spriteVerify.fixed,
|
||||
languages,
|
||||
duration: Date.now() - startTime,
|
||||
});
|
||||
if (!aborted) {
|
||||
send({
|
||||
type: "regen_complete",
|
||||
succeeded,
|
||||
failed,
|
||||
skipped,
|
||||
furniDataAdded: furniDataEntries.length,
|
||||
ownershipFixed: ownership.fixed,
|
||||
nitrosSynced: assetsSynced.copiedNitros,
|
||||
iconsSynced: assetsSynced.copiedIcons,
|
||||
furniDataFixedIds: furniReconcile.fixedIds,
|
||||
furniDataMissing: furniReconcile.missing,
|
||||
furniDataConflicts: furniReconcile.conflicts,
|
||||
offerRebuildChecked: offerRebuild.checked,
|
||||
offerRebuildFixed: offerRebuild.fixed,
|
||||
spriteChecked: spriteVerify.checked,
|
||||
spriteFixed: spriteVerify.fixed,
|
||||
languages,
|
||||
duration: Date.now() - startTime,
|
||||
});
|
||||
}
|
||||
|
||||
controller.close();
|
||||
},
|
||||
|
||||
@@ -184,14 +184,42 @@ export async function findFurniDataIdConflict(
|
||||
}
|
||||
}
|
||||
|
||||
const CHAIN_TIMEOUT_MS = 60_000;
|
||||
|
||||
/**
|
||||
* Wait for the previous lock holder with a timeout.
|
||||
* If the holder stalled (e.g. process crash, unhandled rejection),
|
||||
* we break the chain instead of waiting forever.
|
||||
*/
|
||||
async function waitForFurniDataLock(): Promise<void> {
|
||||
const prev = furniDataLock;
|
||||
if (prev === Promise.resolve()) return;
|
||||
|
||||
try {
|
||||
await Promise.race([
|
||||
prev,
|
||||
new Promise((_, reject) =>
|
||||
setTimeout(
|
||||
() => reject(new Error("furni-data chain timeout")),
|
||||
CHAIN_TIMEOUT_MS,
|
||||
),
|
||||
),
|
||||
]);
|
||||
} catch {
|
||||
// Chain holder timed out or crashed — break the chain
|
||||
furniDataLock = Promise.resolve();
|
||||
}
|
||||
}
|
||||
|
||||
export async function withFurniDataLock<T>(fn: () => Promise<T>): Promise<T> {
|
||||
let release!: () => void;
|
||||
const acquired = new Promise<void>((r) => {
|
||||
release = r;
|
||||
});
|
||||
const prev = furniDataLock;
|
||||
|
||||
await waitForFurniDataLock();
|
||||
furniDataLock = acquired;
|
||||
await prev;
|
||||
|
||||
const lockBase = await getFurnitureDataReadPath();
|
||||
await fs.mkdir(path.dirname(lockBase), { recursive: true });
|
||||
const lockPath = `${lockBase}.lock`;
|
||||
@@ -200,6 +228,7 @@ export async function withFurniDataLock<T>(fn: () => Promise<T>): Promise<T> {
|
||||
return await fn();
|
||||
} finally {
|
||||
await fs.unlink(lockPath).catch(() => {});
|
||||
furniDataLock = Promise.resolve();
|
||||
release?.();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1128,83 +1128,84 @@ export async function verifyAndFixInteractionModesCount(): Promise<{
|
||||
let fixed = 0;
|
||||
let checked = 0;
|
||||
let skipped = 0;
|
||||
const BATCH_SIZE = 500;
|
||||
|
||||
try {
|
||||
const allItems = await db
|
||||
.select({
|
||||
id: ItemsBase.id,
|
||||
itemName: ItemsBase.itemName,
|
||||
publicName: ItemsBase.publicName,
|
||||
allowSit: ItemsBase.allowSit,
|
||||
allowLay: ItemsBase.allowLay,
|
||||
allowWalk: ItemsBase.allowWalk,
|
||||
interactionType: ItemsBase.interactionType,
|
||||
interactionModesCount: ItemsBase.interactionModesCount,
|
||||
})
|
||||
.from(ItemsBase);
|
||||
let offset = 0;
|
||||
while (true) {
|
||||
const batch = await db
|
||||
.select({
|
||||
id: ItemsBase.id,
|
||||
itemName: ItemsBase.itemName,
|
||||
publicName: ItemsBase.publicName,
|
||||
allowSit: ItemsBase.allowSit,
|
||||
allowLay: ItemsBase.allowLay,
|
||||
allowWalk: ItemsBase.allowWalk,
|
||||
interactionType: ItemsBase.interactionType,
|
||||
interactionModesCount: ItemsBase.interactionModesCount,
|
||||
})
|
||||
.from(ItemsBase)
|
||||
.limit(BATCH_SIZE)
|
||||
.offset(offset);
|
||||
|
||||
for (const item of allItems) {
|
||||
checked++;
|
||||
const realFlags = await lookupRealFurniFlags(item.itemName);
|
||||
if (batch.length === 0) break;
|
||||
offset += batch.length;
|
||||
checked += batch.length;
|
||||
|
||||
// Real data absent → do not guess, leave the item untouched.
|
||||
if (realFlags.source === "none") {
|
||||
skipped++;
|
||||
continue;
|
||||
}
|
||||
for (const item of batch) {
|
||||
const realFlags = await lookupRealFurniFlags(item.itemName);
|
||||
|
||||
const animationModes = await realNitroModes(item.itemName);
|
||||
const auto = autoDetectInteraction(item.itemName, item.publicName, {
|
||||
cansiton: realFlags.cansiton,
|
||||
canlayon: realFlags.canlayon,
|
||||
canstandon: realFlags.canstandon,
|
||||
// Furnidata flags are curated → authoritative over keywords.
|
||||
hasActionData: true,
|
||||
animationStates: animationModes,
|
||||
});
|
||||
if (realFlags.source === "none") {
|
||||
skipped++;
|
||||
continue;
|
||||
}
|
||||
|
||||
const changes: Partial<Record<string, unknown>> = {};
|
||||
// allow_sit/lay/walk of 2 is emulator-special (e.g. tents lay=2,
|
||||
// walk-through walk=2) — never overwrite those with a 0/1.
|
||||
if (
|
||||
item.allowSit !== 2 &&
|
||||
(item.allowSit ? 1 : 0) !== (auto.canSit ? 1 : 0)
|
||||
)
|
||||
changes.allowSit = auto.canSit ? 1 : 0;
|
||||
if (
|
||||
item.allowLay !== 2 &&
|
||||
(item.allowLay ? 1 : 0) !== (auto.canLay ? 1 : 0)
|
||||
)
|
||||
changes.allowLay = auto.canLay ? 1 : 0;
|
||||
if (
|
||||
item.allowWalk !== 2 &&
|
||||
(item.allowWalk ? 1 : 0) !== (auto.canStand ? 1 : 0)
|
||||
)
|
||||
changes.allowWalk = auto.canStand ? 1 : 0;
|
||||
if ((item.interactionModesCount ?? 0) !== auto.interactionModesCount)
|
||||
changes.interactionModesCount = auto.interactionModesCount;
|
||||
// Only fill interaction_type when the DB still has the generic
|
||||
// "default" (nothing emulator-specific configured). Never overwrite
|
||||
// an existing handler type (e.g. intelligence_bookcase, nest, wf_*)
|
||||
// — those are emulator semantics that real furnidata cannot replace.
|
||||
if (
|
||||
(item.interactionType === "default" || !item.interactionType) &&
|
||||
auto.interactionType !== "default"
|
||||
) {
|
||||
changes.interactionType = auto.interactionType;
|
||||
}
|
||||
const animationModes = await realNitroModes(item.itemName);
|
||||
const auto = autoDetectInteraction(item.itemName, item.publicName, {
|
||||
cansiton: realFlags.cansiton,
|
||||
canlayon: realFlags.canlayon,
|
||||
canstandon: realFlags.canstandon,
|
||||
hasActionData: true,
|
||||
animationStates: animationModes,
|
||||
});
|
||||
|
||||
if (Object.keys(changes).length > 0) {
|
||||
try {
|
||||
await db
|
||||
.update(ItemsBase)
|
||||
.set(changes)
|
||||
.where(eq(ItemsBase.id, item.id));
|
||||
fixed++;
|
||||
} catch (err) {
|
||||
errors.push(
|
||||
`Failed to fix ${item.itemName}: ${(err as Error).message}`,
|
||||
);
|
||||
const changes: Partial<Record<string, unknown>> = {};
|
||||
if (
|
||||
item.allowSit !== 2 &&
|
||||
(item.allowSit ? 1 : 0) !== (auto.canSit ? 1 : 0)
|
||||
)
|
||||
changes.allowSit = auto.canSit ? 1 : 0;
|
||||
if (
|
||||
item.allowLay !== 2 &&
|
||||
(item.allowLay ? 1 : 0) !== (auto.canLay ? 1 : 0)
|
||||
)
|
||||
changes.allowLay = auto.canLay ? 1 : 0;
|
||||
if (
|
||||
item.allowWalk !== 2 &&
|
||||
(item.allowWalk ? 1 : 0) !== (auto.canStand ? 1 : 0)
|
||||
)
|
||||
changes.allowWalk = auto.canStand ? 1 : 0;
|
||||
if ((item.interactionModesCount ?? 0) !== auto.interactionModesCount)
|
||||
changes.interactionModesCount = auto.interactionModesCount;
|
||||
if (
|
||||
(item.interactionType === "default" || !item.interactionType) &&
|
||||
auto.interactionType !== "default"
|
||||
) {
|
||||
changes.interactionType = auto.interactionType;
|
||||
}
|
||||
|
||||
if (Object.keys(changes).length > 0) {
|
||||
try {
|
||||
await db
|
||||
.update(ItemsBase)
|
||||
.set(changes)
|
||||
.where(eq(ItemsBase.id, item.id));
|
||||
fixed++;
|
||||
} catch (err) {
|
||||
errors.push(
|
||||
`Failed to fix ${item.itemName}: ${(err as Error).message}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,10 +2,12 @@ import { promises as fs } from "node:fs";
|
||||
|
||||
const LOCK_STALE_MS = 60_000;
|
||||
const LOCK_ACQUIRE_TIMEOUT_MS = 30_000;
|
||||
const CHAIN_TIMEOUT_MS = 60_000;
|
||||
|
||||
// In-process serialization, keyed by absolute file path, layered on top of the
|
||||
// on-disk lock so multiple awaiters in the same Node process queue cleanly.
|
||||
const chains = new Map<string, Promise<void>>();
|
||||
const chainTimestamps = new Map<string, number>();
|
||||
|
||||
async function acquireDiskLock(lockPath: string): Promise<void> {
|
||||
const deadline = Date.now() + LOCK_ACQUIRE_TIMEOUT_MS;
|
||||
@@ -34,6 +36,48 @@ async function acquireDiskLock(lockPath: string): Promise<void> {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait for the previous chain holder with a timeout.
|
||||
* If the chain holder stalled (e.g. process crash, unhandled rejection),
|
||||
* we break the chain instead of waiting forever.
|
||||
*/
|
||||
async function waitForChain(filePath: string): Promise<void> {
|
||||
const prev = chains.get(filePath);
|
||||
if (!prev) return;
|
||||
|
||||
const timestamp = chainTimestamps.get(filePath) ?? 0;
|
||||
const elapsed = Date.now() - timestamp;
|
||||
|
||||
if (elapsed > CHAIN_TIMEOUT_MS) {
|
||||
// Previous holder took too long — break the chain
|
||||
console.warn(
|
||||
`[gamedata-lock] Breaking stale chain for ${filePath} (held for ${elapsed}ms)`,
|
||||
);
|
||||
chains.delete(filePath);
|
||||
chainTimestamps.delete(filePath);
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
await Promise.race([
|
||||
prev,
|
||||
new Promise((_, reject) =>
|
||||
setTimeout(
|
||||
() => reject(new Error("chain timeout")),
|
||||
Math.max(1, CHAIN_TIMEOUT_MS - elapsed),
|
||||
),
|
||||
),
|
||||
]);
|
||||
} catch {
|
||||
// Chain holder timed out or crashed — break the chain
|
||||
console.warn(
|
||||
`[gamedata-lock] Chain timeout for ${filePath}, breaking chain`,
|
||||
);
|
||||
chains.delete(filePath);
|
||||
chainTimestamps.delete(filePath);
|
||||
}
|
||||
}
|
||||
|
||||
export async function withGamedataLock<T>(
|
||||
filePath: string,
|
||||
fn: () => Promise<T>,
|
||||
@@ -42,15 +86,20 @@ export async function withGamedataLock<T>(
|
||||
const acquired = new Promise<void>((r) => {
|
||||
release = r;
|
||||
});
|
||||
const prev = chains.get(filePath) ?? Promise.resolve();
|
||||
|
||||
await waitForChain(filePath);
|
||||
|
||||
chains.set(filePath, acquired);
|
||||
await prev;
|
||||
chainTimestamps.set(filePath, Date.now());
|
||||
|
||||
const lockPath = `${filePath}.lock`;
|
||||
await acquireDiskLock(lockPath);
|
||||
try {
|
||||
return await fn();
|
||||
} finally {
|
||||
await fs.unlink(/*turbopackIgnore: true*/ lockPath).catch(() => {});
|
||||
chains.delete(filePath);
|
||||
chainTimestamps.delete(filePath);
|
||||
release();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -87,6 +87,8 @@ function pump() {
|
||||
}
|
||||
}
|
||||
|
||||
const CONVERSION_TIMEOUT_MS = 60_000;
|
||||
|
||||
function run(req: Omit<WorkerRequest, "id">): Promise<WorkerResponse> {
|
||||
if (!ensurePool()) {
|
||||
try {
|
||||
@@ -118,6 +120,33 @@ function run(req: Omit<WorkerRequest, "id">): Promise<WorkerResponse> {
|
||||
};
|
||||
queue.push(job);
|
||||
pump();
|
||||
|
||||
// Timeout: if a job takes too long, remove it from the queue
|
||||
// and reject so the caller can fall back to main-thread conversion.
|
||||
const timer = setTimeout(() => {
|
||||
const idx = queue.findIndex((j) => j.req.id === job.req.id);
|
||||
if (idx !== -1) {
|
||||
queue.splice(idx, 1);
|
||||
pending.delete(job.req.id);
|
||||
reject(new Error("Conversion timed out"));
|
||||
} else {
|
||||
// Already dispatched to a worker — the worker is
|
||||
// responsible for timing out its own task.
|
||||
reject(new Error("Conversion timed out"));
|
||||
}
|
||||
}, CONVERSION_TIMEOUT_MS);
|
||||
|
||||
// Clear the timer once the job resolves/rejects
|
||||
const origResolve = resolve;
|
||||
const origReject = reject;
|
||||
job.resolve = (r) => {
|
||||
clearTimeout(timer);
|
||||
origResolve(r);
|
||||
};
|
||||
job.reject = (e) => {
|
||||
clearTimeout(timer);
|
||||
origReject(e);
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user