diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index 83f4d96fb..3f3c5439d 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -241,14 +241,16 @@ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: report(sessionFile); }; - const poolSize = Math.max(1, Math.min(files.length, opts?.workers ?? defaultWorkerCount())); - if (poolSize === 1) { + const requestedWorkers = Math.max(1, Math.floor(opts?.workers ?? defaultWorkerCount())); + if (requestedWorkers === 1) { for (const sessionFile of files) { await processFile(sessionFile, parseSessionFile); } return { processed: totalProcessed, files: filesProcessed }; } + const poolSize = Math.min(files.length, requestedWorkers); + const handles: WorkerHandle[] = []; for (let i = 0; i < poolSize; i++) handles.push(spawnWorker()); diff --git a/packages/stats/test/sync-serial.test.ts b/packages/stats/test/sync-serial.test.ts index 2a5daa1cf..fa6d6c0e3 100644 --- a/packages/stats/test/sync-serial.test.ts +++ b/packages/stats/test/sync-serial.test.ts @@ -88,4 +88,15 @@ describe("stats sync serial mode", () => { expect(overall.totalRequests).toBe(1); expect(workerSpy).not.toHaveBeenCalled(); }); + + it("spawns a worker pool when callers explicitly request workers: 2 with a single file", async () => { + await writeSessionFile(); + const workerProbe = new Error("worker probe"); + const workerSpy = vi.spyOn(globalThis, "Worker").mockImplementation(() => { + throw workerProbe; + }); + + await expect(syncAllSessions({ workers: 2 })).rejects.toBe(workerProbe); + expect(workerSpy).toHaveBeenCalled(); + }); });