style: format code biome
Local Build and Deploy / deploy (push) Failing after 46s

This commit is contained in:
openhands committed 2026-07-13 21:57:41 +02:00
1 parent 8efd032cc6
commit df38dccbf1
735 files changed
+128321 -120870

No files matched your search

@@ -3,19 +3,21 @@ import { describe, expect, it } from "vitest";
import { resolveGamedataFile } from "./asset-paths";
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",
);
it("uses absolute external paths as-is", () => {
const absolutePath = path.win32.join(
"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,
);
});
});
+11 -7
View File
@@ -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,14 @@ 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 -13
View File
@@ -2,19 +2,23 @@ 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);
});
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);
});
});
+64 -55
View File
@@ -1,67 +1,76 @@
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" },
url: string,
destPath: string,
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)));
}
const res = await fetch(url, { signal: AbortSignal.timeout(15000) });
if (!res.ok) {
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 };
}
continue;
}
const buffer = Buffer.from(await res.arrayBuffer());
if (buffer.length < 8) {
if (attempt === maxRetries) return { ok: false, size: 0 };
continue;
}
if (options?.validate === "swf" && !validateSwfBytes(buffer)) {
if (attempt === maxRetries) {
console.warn(`[import-download] Invalid SWF magic bytes from ${url}`);
return { ok: false, size: 0 };
}
continue;
}
if (options?.validate === "png" && !validatePngBytes(buffer)) {
if (attempt === maxRetries) {
console.warn(`[import-download] Invalid PNG magic bytes from ${url}`);
return { ok: false, size: 0 };
}
continue;
}
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 };
}
}
}
return { ok: false, size: 0 };
for (let attempt = 0; attempt <= maxRetries; attempt++) {
try {
if (attempt > 0) {
await new Promise((r) => setTimeout(r, baseDelay * 2 ** (attempt - 1)));
}
const res = await fetch(url, { signal: AbortSignal.timeout(15000) });
if (!res.ok) {
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 };
}
continue;
}
const buffer = Buffer.from(await res.arrayBuffer());
if (buffer.length < 8) {
if (attempt === maxRetries) return { ok: false, size: 0 };
continue;
}
if (options?.validate === "swf" && !validateSwfBytes(buffer)) {
if (attempt === maxRetries) {
console.warn(`[import-download] Invalid SWF magic bytes from ${url}`);
return { ok: false, size: 0 };
}
continue;
}
if (options?.validate === "png" && !validatePngBytes(buffer)) {
if (attempt === maxRetries) {
console.warn(`[import-download] Invalid PNG magic bytes from ${url}`);
return { ok: false, size: 0 };
}
continue;
}
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 };
}
}
}
return { ok: false, size: 0 };
}
@@ -2,42 +2,51 @@ 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 {
readGamedataJson,
withGamedataLock,
writeGamedataJsonAtomic,
} from "./gamedata-json";
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);
});
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 () => {
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
});
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
});
});
+56 -45
View File
@@ -8,56 +8,67 @@ const LOCK_ACQUIRE_TIMEOUT_MS = 30_000;
const chains = new Map<string, Promise<void>>();
async function acquireDiskLock(lockPath: string): Promise<void> {
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;
} catch (err) {
const e = err as NodeJS.ErrnoException;
if (e.code !== "EEXIST") throw err;
try {
const stat = await fs.stat(/*turbopackIgnore: true*/ lockPath);
if (Date.now() - stat.mtimeMs > LOCK_STALE_MS) {
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));
}
}
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;
} catch (err) {
const e = err as NodeJS.ErrnoException;
if (e.code !== "EEXIST") throw err;
try {
const stat = await fs.stat(/*turbopackIgnore: true*/ lockPath);
if (Date.now() - stat.mtimeMs > LOCK_STALE_MS) {
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));
}
}
}
export async function withGamedataLock<T>(filePath: string, fn: () => Promise<T>): Promise<T> {
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);
try {
return await fn();
} finally {
await fs.unlink(/*turbopackIgnore: true*/ lockPath).catch(() => {});
release();
}
export async function withGamedataLock<T>(
filePath: string,
fn: () => Promise<T>,
): Promise<T> {
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);
try {
return await fn();
} finally {
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);
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);
}
+29 -20
View File
@@ -2,30 +2,39 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { __clearGordonCache, resolveGordonBuildUrl } from "./gordon";
beforeEach(() => {
__clearGordonCache();
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/");
});
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 () => {
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);
});
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);
});
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();
});
});
+20 -12
View File
@@ -1,22 +1,30 @@
const EXTERNAL_VARIABLES_URL = "https://www.habbo.com/gamedata/external_variables/1";
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;
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;
}
+35 -20
View File
@@ -2,27 +2,42 @@ 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();
return text
.split("\n\n")
.filter((l) => l.startsWith("data: "))
.map((l) => JSON.parse(l.slice(6)));
const text = await res.text();
return text
.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 () => {
const res = runSseBatch({
items: ["a", "b", "c"],
concurrency: 2,
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 });
});
it("emits batch_start, per-item, and batch_complete with correct counts", async () => {
const res = runSseBatch({
items: ["a", "b", "c"],
concurrency: 2,
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,
});
});
});
+100 -87
View File
@@ -1,16 +1,20 @@
export interface SseWorkerResult {
ok: boolean;
warnings?: string[];
error?: string;
ok: boolean;
warnings?: string[];
error?: string;
}
export interface RunSseBatchOptions<T> {
items: T[];
concurrency: number;
/** Label used as the `classname` field on item_progress (kept for the existing client parser). */
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>;
items: T[];
concurrency: number;
/** Label used as the `classname` field on item_progress (kept for the existing client parser). */
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>;
}
/**
@@ -18,88 +22,97 @@ 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`));
} catch {
/* stream closed by client */
}
};
const stream = new ReadableStream({
async start(controller) {
const send = (data: unknown) => {
try {
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);
await Promise.allSettled(
chunk.map(async (item, chunkIdx) => {
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 }),
);
if (result.ok) {
succeeded++;
if (result.warnings?.length) withWarnings++;
send({
type: "item_progress",
classname,
status: "done",
index,
warnings: result.warnings?.length ? result.warnings : undefined,
});
} else {
failed++;
send({
type: "item_progress",
classname,
status: "failed",
index,
message: result.error,
});
}
} catch (err) {
failed++;
send({
type: "item_progress",
classname,
status: "failed",
index,
message: (err as Error).message,
});
}
}),
);
}
for (let i = 0; i < items.length; 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,
});
try {
const result = await worker(item, index, (status) =>
send({ type: "item_progress", classname, status, index }),
);
if (result.ok) {
succeeded++;
if (result.warnings?.length) withWarnings++;
send({
type: "item_progress",
classname,
status: "done",
index,
warnings: result.warnings?.length
? result.warnings
: undefined,
});
} else {
failed++;
send({
type: "item_progress",
classname,
status: "failed",
index,
message: result.error,
});
}
} catch (err) {
failed++;
send({
type: "item_progress",
classname,
status: "failed",
index,
message: (err as Error).message,
});
}
}),
);
}
send({
type: "batch_complete",
succeeded,
failed,
warnings: withWarnings,
duration: Date.now() - startTime,
});
controller.close();
},
});
send({
type: "batch_complete",
succeeded,
failed,
warnings: withWarnings,
duration: Date.now() - startTime,
});
controller.close();
},
});
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
Connection: "keep-alive",
},
});
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
Connection: "keep-alive",
},
});
}