diff --git a/packages/stats/CHANGELOG.md b/packages/stats/CHANGELOG.md index 8e5db7136..b0601267d 100644 --- a/packages/stats/CHANGELOG.md +++ b/packages/stats/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Kept stats session sync on the serial parser path for `workers: 1` and macOS defaults, avoiding Bun worker re-entry aborts when launching `/stats` ([#3733](https://github.com/can1357/oh-my-pi/issues/3733)). + ## [16.2.3] - 2026-06-28 ### Added diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index 7c1b337fa..83f4d96fb 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -23,7 +23,7 @@ import { setFileOffset, updateUserMessageLinks, } from "./db"; -import { getSessionEntry, listAllSessionFiles, type ParseSessionResult } from "./parser"; +import { getSessionEntry, listAllSessionFiles, type ParseSessionResult, parseSessionFile } from "./parser"; import type { SyncWorkerRequest, SyncWorkerResponse } from "./sync-worker"; // Coding-agent binary/bundle workers route through the CLI entrypoint with a // hidden argv mode, so the compiled binary and npm bundle only need one @@ -68,6 +68,9 @@ export interface SyncOptions { } function defaultWorkerCount(): number { + // Bun 1.3.x can abort the macOS process when stats sync workers re-enter + // the compiled `omp` binary. Keep macOS on the documented serial path. + if (process.platform === "darwin") return 1; // `navigator.hardwareConcurrency` is the portable answer in Bun; fall // back to a small fixed pool if it's somehow unavailable. const hw = typeof navigator !== "undefined" ? (navigator.hardwareConcurrency ?? 0) : 0; @@ -183,11 +186,11 @@ export async function smokeTestSyncWorker({ timeoutMs = 5_000 }: { timeoutMs?: n /** * Sync all session files to the database. * - * Parsing fans out across a worker pool (one in-flight job per worker) - * while DB writes and offset bookkeeping stay on the calling thread so the - * single SQLite handle stays uncontended. `onProgress` fires once per - * completed file (skipped files included so the bar walks at a steady - * rate). + * `workers: 1` parses inline. Larger pools fan parsing out across workers + * (one in-flight job per worker) while DB writes and offset bookkeeping stay on + * the calling thread so the single SQLite handle stays uncontended. + * `onProgress` fires once per completed file (skipped files included so the + * bar walks at a steady rate). */ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: number; files: number }> { await initDb(); @@ -200,10 +203,6 @@ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: let completed = 0; let cursor = 0; - const poolSize = Math.max(1, Math.min(files.length, opts?.workers ?? defaultWorkerCount())); - const handles: WorkerHandle[] = []; - for (let i = 0; i < poolSize; i++) handles.push(spawnWorker()); - const report = (sessionFile: string) => { completed++; opts?.onProgress?.({ @@ -214,34 +213,51 @@ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: }); }; + const processFile = async ( + sessionFile: string, + parse: (sessionFile: string, fromOffset: number) => Promise, + ): Promise => { + let fileStats: fs.Stats; + try { + fileStats = await fs.promises.stat(sessionFile); + } catch { + report(sessionFile); + return; + } + const lastModified = fileStats.mtimeMs; + const stored = getFileOffset(sessionFile); + if (stored && stored.lastModified >= lastModified) { + report(sessionFile); + return; + } + + const fromOffset = stored?.offset ?? 0; + const result = await parse(sessionFile, fromOffset); + const inserted = applyParseResult(sessionFile, lastModified, result); + if (inserted > 0) { + totalProcessed += inserted; + filesProcessed++; + } + report(sessionFile); + }; + + const poolSize = Math.max(1, Math.min(files.length, opts?.workers ?? defaultWorkerCount())); + if (poolSize === 1) { + for (const sessionFile of files) { + await processFile(sessionFile, parseSessionFile); + } + return { processed: totalProcessed, files: filesProcessed }; + } + + const handles: WorkerHandle[] = []; + for (let i = 0; i < poolSize; i++) handles.push(spawnWorker()); + async function drain(handle: WorkerHandle): Promise { while (true) { const idx = cursor++; if (idx >= files.length) return; const sessionFile = files[idx]; - - let fileStats: fs.Stats; - try { - fileStats = await fs.promises.stat(sessionFile); - } catch { - report(sessionFile); - continue; - } - const lastModified = fileStats.mtimeMs; - const stored = getFileOffset(sessionFile); - if (stored && stored.lastModified >= lastModified) { - report(sessionFile); - continue; - } - - const fromOffset = stored?.offset ?? 0; - const result = await dispatch(handle, { sessionFile, fromOffset }); - const inserted = applyParseResult(sessionFile, lastModified, result); - if (inserted > 0) { - totalProcessed += inserted; - filesProcessed++; - } - report(sessionFile); + await processFile(sessionFile, (file, fromOffset) => dispatch(handle, { sessionFile: file, fromOffset })); } } diff --git a/packages/stats/test/sync-serial.test.ts b/packages/stats/test/sync-serial.test.ts new file mode 100644 index 000000000..2a5daa1cf --- /dev/null +++ b/packages/stats/test/sync-serial.test.ts @@ -0,0 +1,91 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; +import * as fs from "node:fs/promises"; +import * as os from "node:os"; +import * as path from "node:path"; +import { syncAllSessions } from "@oh-my-pi/omp-stats/aggregator"; +import { closeDb, getOverallStats } from "@oh-my-pi/omp-stats/db"; +import { getAgentDir, getSessionsDir, setAgentDir, TempDir } from "@oh-my-pi/pi-utils"; + +const originalConfigDir = process.env.PI_CONFIG_DIR; +const originalAgentDir = getAgentDir(); +let tempDir: TempDir | null = null; + +beforeEach(() => { + tempDir = TempDir.createSync("@pi-stats-sync-serial-"); + const configDir = path.relative(os.homedir(), tempDir.join("config")); + process.env.PI_CONFIG_DIR = configDir; + setAgentDir(tempDir.join("agent")); +}); + +afterEach(() => { + vi.restoreAllMocks(); + closeDb(); + if (originalConfigDir === undefined) { + delete process.env.PI_CONFIG_DIR; + } else { + process.env.PI_CONFIG_DIR = originalConfigDir; + } + setAgentDir(originalAgentDir); + tempDir?.removeSync(); + tempDir = null; +}); + +async function writeSessionFile(): Promise { + const sessionDir = path.join(getSessionsDir(), "--tmp--sync-serial"); + await fs.mkdir(sessionDir, { recursive: true }); + const timestamp = new Date().toISOString(); + const sessionFile = path.join(sessionDir, "session.jsonl"); + const assistant = { + type: "message", + id: "assistant-1", + parentId: null, + timestamp, + message: { + role: "assistant", + content: [{ type: "text", text: "ok" }], + api: "openai-responses", + provider: "openai", + model: "gpt-5.4", + usage: { + input: 1, + output: 2, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 3, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop", + timestamp: Date.now(), + duration: 10, + ttft: 5, + }, + }; + await Bun.write(sessionFile, `${JSON.stringify(assistant)}\n`); +} + +describe("stats sync serial mode", () => { + it("honors workers: 1 without spawning a worker", async () => { + await writeSessionFile(); + const workerSpy = vi.spyOn(globalThis, "Worker"); + + const synced = await syncAllSessions({ workers: 1 }); + const overall = getOverallStats(); + + expect(synced.files).toBe(1); + expect(overall.totalRequests).toBe(1); + expect(workerSpy).not.toHaveBeenCalled(); + }); + + it("uses the serial parser by default on macOS", async () => { + await writeSessionFile(); + vi.spyOn(process, "platform", "get").mockReturnValue("darwin"); + const workerSpy = vi.spyOn(globalThis, "Worker"); + + const synced = await syncAllSessions(); + const overall = getOverallStats(); + + expect(synced.files).toBe(1); + expect(overall.totalRequests).toBe(1); + expect(workerSpy).not.toHaveBeenCalled(); + }); +});