revert fb8e77bb68
Local Build and Deploy / deploy (push) Successful in 1m11s

revert style: clean up code with prettier and eslint
This commit is contained in:
remco committed 2026-07-12 21:02:03 +02:00
1 parent 27de078f54
commit e85e4d74ea
378 files changed
+20710 -22428

No files matched your search

@@ -1,21 +1,21 @@
import path from "node:path";
import { describe, expect, it } from "vitest";
import { resolveGamedataFile } from "./asset-paths";
import path from 'node:path'
import { describe, expect, it } from 'vitest'
import { resolveGamedataFile } from './asset-paths'
describe("resolveGamedataFile", () => {
it("uses absolute external paths as-is", () => {
describe('resolveGamedataFile', () => {
it('uses absolute external paths as-is', () => {
const absolutePath = path.win32.join(
"E:\\",
"Users",
"simol",
"Desktop",
"DEV",
"Nitro-Files",
"nitro-assets",
"gamedata",
"FigureMap.json",
);
'E:\\',
'Users',
'simol',
'Desktop',
'DEV',
'Nitro-Files',
'nitro-assets',
'gamedata',
'FigureMap.json',
)
expect(resolveGamedataFile(absolutePath, "FigureMap.json")).toBe(absolutePath);
});
});
expect(resolveGamedataFile(absolutePath, 'FigureMap.json')).toBe(absolutePath)
})
})
+9 -9
View File
@@ -1,5 +1,5 @@
import { existsSync } from "node:fs";
import path from "node:path";
import { existsSync } from 'node:fs'
import path from 'node:path'
// The CMS default gamedata JSON paths point under public/nitro-assets/gamedata/,
// but this deployment keeps the gamedata JSON (FigureMap/FigureData/EffectMap…)
@@ -9,7 +9,7 @@ import path from "node:path";
// the client's served asset dir — see *_NITRO_DIR consts in the import services.)
function pub(...segs: string[]): string {
return path.join(/*turbopackIgnore: true*/ process.cwd(), "public", ...segs);
return path.join(/*turbopackIgnore: true*/ process.cwd(), 'public', ...segs)
}
/**
@@ -19,10 +19,10 @@ function pub(...segs: string[]): string {
* Returns the configured path when neither exists (write target for fresh setups).
*/
export function resolveGamedataFile(url: string, fallbackName: string): string {
if (path.win32.isAbsolute(url) || (path.isAbsolute(url) && !url.startsWith("/"))) return url;
const configured = pub(url.replace(/^\/+/, ""));
if (existsSync(/*turbopackIgnore: true*/ configured)) return configured;
const gamedata = pub("Gamedata", "config", fallbackName);
if (existsSync(/*turbopackIgnore: true*/ gamedata)) return gamedata;
return configured;
if (path.win32.isAbsolute(url) || (path.isAbsolute(url) && !url.startsWith('/'))) return url
const configured = pub(url.replace(/^\/+/, ''))
if (existsSync(/*turbopackIgnore: true*/ configured)) return configured
const gamedata = pub('Gamedata', 'config', fallbackName)
if (existsSync(/*turbopackIgnore: true*/ gamedata)) return gamedata
return configured
}
+17 -17
View File
@@ -1,20 +1,20 @@
import { describe, expect, it } from "vitest";
import { validatePngBytes, validateSwfBytes } from "./download";
import { describe, expect, it } from 'vitest'
import { validatePngBytes, validateSwfBytes } from './download'
describe("import/core/download validators", () => {
it("accepts FWS/CWS/ZWS swf magic", () => {
expect(validateSwfBytes(Buffer.from("FWS\x06\x00\x00\x00\x00"))).toBe(true);
expect(validateSwfBytes(Buffer.from("CWS\x06\x00\x00\x00\x00"))).toBe(true);
expect(validateSwfBytes(Buffer.from("ZWS\x06\x00\x00\x00\x00"))).toBe(true);
});
describe('import/core/download validators', () => {
it('accepts FWS/CWS/ZWS swf magic', () => {
expect(validateSwfBytes(Buffer.from('FWS\x06\x00\x00\x00\x00'))).toBe(true)
expect(validateSwfBytes(Buffer.from('CWS\x06\x00\x00\x00\x00'))).toBe(true)
expect(validateSwfBytes(Buffer.from('ZWS\x06\x00\x00\x00\x00'))).toBe(true)
})
it("rejects non-swf and too-short buffers", () => {
expect(validateSwfBytes(Buffer.from("PNG\x00\x00\x00\x00\x00"))).toBe(false);
expect(validateSwfBytes(Buffer.from("FW"))).toBe(false);
});
it('rejects non-swf and too-short buffers', () => {
expect(validateSwfBytes(Buffer.from('PNG\x00\x00\x00\x00\x00'))).toBe(false)
expect(validateSwfBytes(Buffer.from('FW'))).toBe(false)
})
it("accepts png magic and rejects others", () => {
expect(validatePngBytes(Buffer.from([137, 80, 78, 71, 13, 10, 26, 10]))).toBe(true);
expect(validatePngBytes(Buffer.from([1, 2, 3, 4, 5, 6, 7, 8]))).toBe(false);
});
});
it('accepts png magic and rejects others', () => {
expect(validatePngBytes(Buffer.from([137, 80, 78, 71, 13, 10, 26, 10]))).toBe(true)
expect(validatePngBytes(Buffer.from([1, 2, 3, 4, 5, 6, 7, 8]))).toBe(false)
})
})
+32 -32
View File
@@ -1,67 +1,67 @@
import { promises as fs } from "node:fs";
import { promises as fs } from 'node:fs'
export function validateSwfBytes(buffer: Buffer): boolean {
if (buffer.length < 8) return false;
const sig = buffer.toString("ascii", 0, 3);
return sig === "FWS" || sig === "CWS" || sig === "ZWS";
if (buffer.length < 8) return false
const sig = buffer.toString('ascii', 0, 3)
return sig === 'FWS' || sig === 'CWS' || sig === 'ZWS'
}
export function validatePngBytes(buffer: Buffer): boolean {
if (buffer.length < 8) return false;
return buffer[0] === 137 && buffer[1] === 80 && buffer[2] === 78 && buffer[3] === 71;
if (buffer.length < 8) return false
return buffer[0] === 137 && buffer[1] === 80 && buffer[2] === 78 && buffer[3] === 71
}
export async function downloadFile(
url: string,
destPath: string,
options?: { maxRetries?: number; validate?: "swf" | "png" },
options?: { maxRetries?: number; validate?: 'swf' | 'png' },
): Promise<{ ok: boolean; size: number }> {
const maxRetries = options?.maxRetries ?? 3;
const baseDelay = 1000;
const maxRetries = options?.maxRetries ?? 3
const baseDelay = 1000
for (let attempt = 0; attempt <= maxRetries; attempt++) {
try {
if (attempt > 0) {
await new Promise((r) => setTimeout(r, baseDelay * 2 ** (attempt - 1)));
await new Promise((r) => setTimeout(r, baseDelay * 2 ** (attempt - 1)))
}
const res = await fetch(url, { signal: AbortSignal.timeout(15000) });
const res = await fetch(url, { signal: AbortSignal.timeout(15000) })
if (!res.ok) {
const deterministic = res.status === 404 || res.status === 410 || res.status === 403;
const deterministic = res.status === 404 || res.status === 410 || res.status === 403
if (deterministic || attempt === maxRetries) {
console.warn(
`[import-download] Download failed ${url}: ${res.status}${deterministic ? " (deterministic, not retrying)" : ` after ${maxRetries + 1} attempts`}`,
);
return { ok: false, size: 0 };
`[import-download] Download failed ${url}: ${res.status}${deterministic ? ' (deterministic, not retrying)' : ` after ${maxRetries + 1} attempts`}`,
)
return { ok: false, size: 0 }
}
continue;
continue
}
const buffer = Buffer.from(await res.arrayBuffer());
const buffer = Buffer.from(await res.arrayBuffer())
if (buffer.length < 8) {
if (attempt === maxRetries) return { ok: false, size: 0 };
continue;
if (attempt === maxRetries) return { ok: false, size: 0 }
continue
}
if (options?.validate === "swf" && !validateSwfBytes(buffer)) {
if (options?.validate === 'swf' && !validateSwfBytes(buffer)) {
if (attempt === maxRetries) {
console.warn(`[import-download] Invalid SWF magic bytes from ${url}`);
return { ok: false, size: 0 };
console.warn(`[import-download] Invalid SWF magic bytes from ${url}`)
return { ok: false, size: 0 }
}
continue;
continue
}
if (options?.validate === "png" && !validatePngBytes(buffer)) {
if (options?.validate === 'png' && !validatePngBytes(buffer)) {
if (attempt === maxRetries) {
console.warn(`[import-download] Invalid PNG magic bytes from ${url}`);
return { ok: false, size: 0 };
console.warn(`[import-download] Invalid PNG magic bytes from ${url}`)
return { ok: false, size: 0 }
}
continue;
continue
}
await fs.writeFile(/*turbopackIgnore: true*/ destPath, buffer);
return { ok: true, size: buffer.length };
await fs.writeFile(/*turbopackIgnore: true*/ destPath, buffer)
return { ok: true, size: buffer.length }
} catch (err) {
if (attempt === maxRetries) {
console.warn(`[import-download] Download error ${url}:`, (err as Error).message);
return { ok: false, size: 0 };
console.warn(`[import-download] Download error ${url}:`, (err as Error).message)
return { ok: false, size: 0 }
}
}
}
return { ok: false, size: 0 };
return { ok: false, size: 0 }
}
@@ -1,43 +1,43 @@
import { promises as fs } from "node:fs";
import os from "node:os";
import path from "node:path";
import { afterEach, beforeEach, describe, expect, it } from "vitest";
import { readGamedataJson, withGamedataLock, writeGamedataJsonAtomic } from "./gamedata-json";
import { promises as fs } from 'node:fs'
import os from 'node:os'
import path from 'node:path'
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
import { readGamedataJson, withGamedataLock, writeGamedataJsonAtomic } from './gamedata-json'
let dir: string;
let file: string;
let dir: string
let file: string
beforeEach(async () => {
dir = await fs.mkdtemp(path.join(os.tmpdir(), "gamedata-"));
file = path.join(dir, "EffectMap.json");
await fs.writeFile(file, JSON.stringify({ effects: [{ id: "1", lib: "A" }] }));
});
dir = await fs.mkdtemp(path.join(os.tmpdir(), 'gamedata-'))
file = path.join(dir, 'EffectMap.json')
await fs.writeFile(file, JSON.stringify({ effects: [{ id: '1', lib: 'A' }] }))
})
afterEach(async () => {
await fs.rm(dir, { recursive: true, force: true });
});
await fs.rm(dir, { recursive: true, force: true })
})
describe("import/core/gamedata-json", () => {
it("reads parsed JSON", async () => {
const data = await readGamedataJson<{ effects: unknown[] }>(file);
expect(data.effects).toHaveLength(1);
});
describe('import/core/gamedata-json', () => {
it('reads parsed JSON', async () => {
const data = await readGamedataJson<{ effects: unknown[] }>(file)
expect(data.effects).toHaveLength(1)
})
it("writes atomically (round-trips)", async () => {
await writeGamedataJsonAtomic(file, { effects: [{ id: "9", lib: "Z" }] });
const data = await readGamedataJson<{ effects: Array<{ id: string }> }>(file);
expect(data.effects[0].id).toBe("9");
});
it('writes atomically (round-trips)', async () => {
await writeGamedataJsonAtomic(file, { effects: [{ id: '9', lib: 'Z' }] })
const data = await readGamedataJson<{ effects: Array<{ id: string }> }>(file)
expect(data.effects[0].id).toBe('9')
})
it("serializes concurrent writers (no lost update)", async () => {
it('serializes concurrent writers (no lost update)', async () => {
const bump = (n: number) =>
withGamedataLock(file, async () => {
const d = await readGamedataJson<{ effects: unknown[] }>(file);
d.effects.push({ id: String(n), lib: `L${n}` });
await writeGamedataJsonAtomic(file, d);
});
await Promise.all([bump(2), bump(3), bump(4)]);
const d = await readGamedataJson<{ effects: unknown[] }>(file);
expect(d.effects).toHaveLength(4); // 1 seed + 3 appended, none lost
});
});
const d = await readGamedataJson<{ effects: unknown[] }>(file)
d.effects.push({ id: String(n), lib: `L${n}` })
await writeGamedataJsonAtomic(file, d)
})
await Promise.all([bump(2), bump(3), bump(4)])
const d = await readGamedataJson<{ effects: unknown[] }>(file)
expect(d.effects).toHaveLength(4) // 1 seed + 3 appended, none lost
})
})
+32 -32
View File
@@ -1,63 +1,63 @@
import { promises as fs } from "node:fs";
import { promises as fs } from 'node:fs'
const LOCK_STALE_MS = 60_000;
const LOCK_ACQUIRE_TIMEOUT_MS = 30_000;
const LOCK_STALE_MS = 60_000
const LOCK_ACQUIRE_TIMEOUT_MS = 30_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 chains = new Map<string, Promise<void>>()
async function acquireDiskLock(lockPath: string): Promise<void> {
const deadline = Date.now() + LOCK_ACQUIRE_TIMEOUT_MS;
const deadline = Date.now() + LOCK_ACQUIRE_TIMEOUT_MS
while (true) {
try {
const fh = await fs.open(/*turbopackIgnore: true*/ lockPath, "wx");
await fh.writeFile(String(process.pid));
await fh.close();
return;
const fh = await fs.open(/*turbopackIgnore: true*/ lockPath, 'wx')
await fh.writeFile(String(process.pid))
await fh.close()
return
} catch (err) {
const e = err as NodeJS.ErrnoException;
if (e.code !== "EEXIST") throw err;
const e = err as NodeJS.ErrnoException
if (e.code !== 'EEXIST') throw err
try {
const stat = await fs.stat(/*turbopackIgnore: true*/ lockPath);
const stat = await fs.stat(/*turbopackIgnore: true*/ lockPath)
if (Date.now() - stat.mtimeMs > LOCK_STALE_MS) {
await fs.unlink(/*turbopackIgnore: true*/ lockPath).catch(() => {});
continue;
await fs.unlink(/*turbopackIgnore: true*/ lockPath).catch(() => {})
continue
}
} catch {
/* lock vanished between checks */
}
if (Date.now() > deadline) throw new Error(`gamedata lock timeout: ${lockPath}`);
await new Promise((r) => setTimeout(r, 100));
if (Date.now() > deadline) throw new Error(`gamedata lock timeout: ${lockPath}`)
await new Promise((r) => setTimeout(r, 100))
}
}
}
export async function withGamedataLock<T>(filePath: string, fn: () => Promise<T>): Promise<T> {
let release!: () => void;
let release!: () => void
const acquired = new Promise<void>((r) => {
release = r;
});
const prev = chains.get(filePath) ?? Promise.resolve();
chains.set(filePath, acquired);
await prev;
const lockPath = `${filePath}.lock`;
await acquireDiskLock(lockPath);
release = r
})
const prev = chains.get(filePath) ?? Promise.resolve()
chains.set(filePath, acquired)
await prev
const lockPath = `${filePath}.lock`
await acquireDiskLock(lockPath)
try {
return await fn();
return await fn()
} finally {
await fs.unlink(/*turbopackIgnore: true*/ lockPath).catch(() => {});
release();
await fs.unlink(/*turbopackIgnore: true*/ lockPath).catch(() => {})
release()
}
}
export async function readGamedataJson<T>(filePath: string): Promise<T> {
const raw = await fs.readFile(/*turbopackIgnore: true*/ filePath, "utf-8");
return JSON.parse(raw) as T;
const raw = await fs.readFile(/*turbopackIgnore: true*/ filePath, 'utf-8')
return JSON.parse(raw) as T
}
export async function writeGamedataJsonAtomic(filePath: string, data: unknown): Promise<void> {
const tmp = `${filePath}.tmp`;
await fs.writeFile(/*turbopackIgnore: true*/ tmp, JSON.stringify(data, null, 2), "utf-8");
await fs.rename(/*turbopackIgnore: true*/ tmp, filePath);
const tmp = `${filePath}.tmp`
await fs.writeFile(/*turbopackIgnore: true*/ tmp, JSON.stringify(data, null, 2), 'utf-8')
await fs.rename(/*turbopackIgnore: true*/ tmp, filePath)
}
+25 -25
View File
@@ -1,31 +1,31 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { __clearGordonCache, resolveGordonBuildUrl } from "./gordon";
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { __clearGordonCache, resolveGordonBuildUrl } from './gordon'
beforeEach(() => {
__clearGordonCache();
vi.restoreAllMocks();
});
afterEach(() => vi.restoreAllMocks());
__clearGordonCache()
vi.restoreAllMocks()
})
afterEach(() => vi.restoreAllMocks())
describe("import/core/gordon", () => {
it("parses flash.client.url from external_variables and normalizes", async () => {
vi.spyOn(globalThis, "fetch").mockResolvedValue(
new Response("flash.client.url=//images.habbo.com/gordon/PRODUCTION-XYZ/\n", { status: 200 }),
);
expect(await resolveGordonBuildUrl()).toBe("https://images.habbo.com/gordon/PRODUCTION-XYZ/");
});
describe('import/core/gordon', () => {
it('parses flash.client.url from external_variables and normalizes', async () => {
vi.spyOn(globalThis, 'fetch').mockResolvedValue(
new Response('flash.client.url=//images.habbo.com/gordon/PRODUCTION-XYZ/\n', { status: 200 }),
)
expect(await resolveGordonBuildUrl()).toBe('https://images.habbo.com/gordon/PRODUCTION-XYZ/')
})
it("caches the resolved build (second call does not refetch)", async () => {
it('caches the resolved build (second call does not refetch)', async () => {
const spy = vi
.spyOn(globalThis, "fetch")
.mockResolvedValue(new Response("flash.client.url=//x/gordon/B/\n", { status: 200 }));
await resolveGordonBuildUrl();
await resolveGordonBuildUrl();
expect(spy).toHaveBeenCalledTimes(1);
});
.spyOn(globalThis, 'fetch')
.mockResolvedValue(new Response('flash.client.url=//x/gordon/B/\n', { status: 200 }))
await resolveGordonBuildUrl()
await resolveGordonBuildUrl()
expect(spy).toHaveBeenCalledTimes(1)
})
it("throws when flash.client.url is absent", async () => {
vi.spyOn(globalThis, "fetch").mockResolvedValue(new Response("other=1\n", { status: 200 }));
await expect(resolveGordonBuildUrl()).rejects.toThrow();
});
});
it('throws when flash.client.url is absent', async () => {
vi.spyOn(globalThis, 'fetch').mockResolvedValue(new Response('other=1\n', { status: 200 }))
await expect(resolveGordonBuildUrl()).rejects.toThrow()
})
})
+14 -14
View File
@@ -1,22 +1,22 @@
const EXTERNAL_VARIABLES_URL = "https://www.habbo.com/gamedata/external_variables/1";
const CACHE_TTL = 30 * 60 * 1000;
const EXTERNAL_VARIABLES_URL = 'https://www.habbo.com/gamedata/external_variables/1'
const CACHE_TTL = 30 * 60 * 1000
let gordonCache: { url: string; ts: number } | null = null;
let gordonCache: { url: string; ts: number } | null = null
export function __clearGordonCache(): void {
gordonCache = null;
gordonCache = null
}
/** Resolve the current Habbo "gordon" asset build URL (trailing slash) from external_variables. */
export async function resolveGordonBuildUrl(): Promise<string> {
const now = Date.now();
if (gordonCache && now - gordonCache.ts < CACHE_TTL) return gordonCache.url;
const res = await fetch(EXTERNAL_VARIABLES_URL, { signal: AbortSignal.timeout(20000) });
if (!res.ok) throw new Error(`external_variables fetch failed: ${res.status}`);
const text = await res.text();
const line = /^flash\.client\.url=(.+)$/m.exec(text)?.[1]?.trim();
if (!line) throw new Error("flash.client.url not found in external_variables");
const normalized = (line.startsWith("//") ? `https:${line}` : line).replace(/\/?$/, "/");
gordonCache = { url: normalized, ts: now };
return normalized;
const now = Date.now()
if (gordonCache && now - gordonCache.ts < CACHE_TTL) return gordonCache.url
const res = await fetch(EXTERNAL_VARIABLES_URL, { signal: AbortSignal.timeout(20000) })
if (!res.ok) throw new Error(`external_variables fetch failed: ${res.status}`)
const text = await res.text()
const line = /^flash\.client\.url=(.+)$/m.exec(text)?.[1]?.trim()
if (!line) throw new Error('flash.client.url not found in external_variables')
const normalized = (line.startsWith('//') ? `https:${line}` : line).replace(/\/?$/, '/')
gordonCache = { url: normalized, ts: now }
return normalized
}
+20 -20
View File
@@ -1,28 +1,28 @@
import { describe, expect, it } from "vitest";
import { runSseBatch } from "./sse-batch";
import { describe, expect, it } from 'vitest'
import { runSseBatch } from './sse-batch'
async function collect(res: Response): Promise<Array<Record<string, unknown>>> {
const text = await res.text();
const text = await res.text()
return text
.split("\n\n")
.filter((l) => l.startsWith("data: "))
.map((l) => JSON.parse(l.slice(6)));
.split('\n\n')
.filter((l) => l.startsWith('data: '))
.map((l) => JSON.parse(l.slice(6)))
}
describe("import/core/sse-batch", () => {
it("emits batch_start, per-item, and batch_complete with correct counts", async () => {
describe('import/core/sse-batch', () => {
it('emits batch_start, per-item, and batch_complete with correct counts', async () => {
const res = runSseBatch({
items: ["a", "b", "c"],
items: ['a', 'b', 'c'],
concurrency: 2,
worker: async (item) => ({ ok: item !== "b", error: item === "b" ? "boom" : undefined }),
worker: async (item) => ({ ok: item !== 'b', error: item === 'b' ? 'boom' : undefined }),
labelOf: (item) => item,
});
const events = await collect(res);
expect(events[0]).toMatchObject({ type: "batch_start", total: 3, concurrency: 2 });
const done = events.filter((e) => e.type === "item_progress" && e.status === "done");
const failed = events.filter((e) => e.type === "item_progress" && e.status === "failed");
expect(done).toHaveLength(2);
expect(failed).toHaveLength(1);
expect(events.at(-1)).toMatchObject({ type: "batch_complete", succeeded: 2, failed: 1 });
});
});
})
const events = await collect(res)
expect(events[0]).toMatchObject({ type: 'batch_start', total: 3, concurrency: 2 })
const done = events.filter((e) => e.type === 'item_progress' && e.status === 'done')
const failed = events.filter((e) => e.type === 'item_progress' && e.status === 'failed')
expect(done).toHaveLength(2)
expect(failed).toHaveLength(1)
expect(events.at(-1)).toMatchObject({ type: 'batch_complete', succeeded: 2, failed: 1 })
})
})
+45 -45
View File
@@ -1,16 +1,16 @@
export interface SseWorkerResult {
ok: boolean;
warnings?: string[];
error?: string;
ok: boolean
warnings?: string[]
error?: string
}
export interface RunSseBatchOptions<T> {
items: T[];
concurrency: number;
items: T[]
concurrency: number
/** Label used as the `classname` field on item_progress (kept for the existing client parser). */
labelOf: (item: T) => string;
labelOf: (item: T) => string
/** Per-item worker. `report(status)` streams intermediate progress (e.g. 'downloading'). */
worker: (item: T, index: number, report: (status: string) => void) => Promise<SseWorkerResult>;
worker: (item: T, index: number, report: (status: string) => void) => Promise<SseWorkerResult>
}
/**
@@ -18,88 +18,88 @@ export interface RunSseBatchOptions<T> {
* parser consumes: batch_start / item_progress / batch_complete.
*/
export function runSseBatch<T>(opts: RunSseBatchOptions<T>): Response {
const { items, labelOf, worker } = opts;
const concurrency = Math.min(Math.max(opts.concurrency || 3, 1), 5);
const encoder = new TextEncoder();
const { items, labelOf, worker } = opts
const concurrency = Math.min(Math.max(opts.concurrency || 3, 1), 5)
const encoder = new TextEncoder()
const stream = new ReadableStream({
async start(controller) {
const send = (data: unknown) => {
try {
controller.enqueue(encoder.encode(`data: ${JSON.stringify(data)}\n\n`));
controller.enqueue(encoder.encode(`data: ${JSON.stringify(data)}\n\n`))
} catch {
/* stream closed by client */
}
};
}
const startTime = Date.now();
send({ type: "batch_start", total: items.length, concurrency });
const startTime = Date.now()
send({ type: 'batch_start', total: items.length, concurrency })
let succeeded = 0;
let failed = 0;
let withWarnings = 0;
let succeeded = 0
let failed = 0
let withWarnings = 0
for (let i = 0; i < items.length; i += concurrency) {
const chunk = items.slice(i, i + concurrency);
const chunk = items.slice(i, i + concurrency)
await Promise.allSettled(
chunk.map(async (item, chunkIdx) => {
const index = i + chunkIdx;
const classname = labelOf(item);
send({ type: "item_progress", classname, status: "started", index });
const index = i + chunkIdx
const classname = labelOf(item)
send({ type: 'item_progress', classname, status: 'started', index })
try {
const result = await worker(item, index, (status) =>
send({ type: "item_progress", classname, status, index }),
);
send({ type: 'item_progress', classname, status, index }),
)
if (result.ok) {
succeeded++;
if (result.warnings?.length) withWarnings++;
succeeded++
if (result.warnings?.length) withWarnings++
send({
type: "item_progress",
type: 'item_progress',
classname,
status: "done",
status: 'done',
index,
warnings: result.warnings?.length ? result.warnings : undefined,
});
})
} else {
failed++;
failed++
send({
type: "item_progress",
type: 'item_progress',
classname,
status: "failed",
status: 'failed',
index,
message: result.error,
});
})
}
} catch (err) {
failed++;
failed++
send({
type: "item_progress",
type: 'item_progress',
classname,
status: "failed",
status: 'failed',
index,
message: (err as Error).message,
});
})
}
}),
);
)
}
send({
type: "batch_complete",
type: 'batch_complete',
succeeded,
failed,
warnings: withWarnings,
duration: Date.now() - startTime,
});
controller.close();
})
controller.close()
},
});
})
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
Connection: "keep-alive",
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
Connection: 'keep-alive',
},
});
})
}