diff --git a/packages/stats/CHANGELOG.md b/packages/stats/CHANGELOG.md index b0601267d..5e2a32afd 100644 --- a/packages/stats/CHANGELOG.md +++ b/packages/stats/CHANGELOG.md @@ -5,6 +5,7 @@ ### 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)). +- Replaced the native `Bun.JSONL.parseChunk` session parser path with a lenient JS line scanner, avoiding Bun aborts on large stats session files ([#3733](https://github.com/can1357/oh-my-pi/issues/3733)). ## [16.2.3] - 2026-06-28 diff --git a/packages/stats/src/parser.ts b/packages/stats/src/parser.ts index d06c87447..d8453a930 100644 --- a/packages/stats/src/parser.ts +++ b/packages/stats/src/parser.ts @@ -164,54 +164,53 @@ function extractStats( } const LF = 0x0a; +const CR = 0x0d; +const jsonLineDecoder = new TextDecoder(); + +function parseJsonLine(bytes: Uint8Array, start: number, end: number): SessionEntry | null { + while (end > start && bytes[end - 1] === CR) end--; + if (end <= start) return null; + try { + return JSON.parse(jsonLineDecoder.decode(bytes.subarray(start, end))) as SessionEntry; + } catch { + return null; + } +} + +function visitSessionEntriesLenient(bytes: Uint8Array, visit: (entry: SessionEntry) => void): number { + let cursor = 0; + let read = 0; + + while (cursor < bytes.length) { + const newline = bytes.indexOf(LF, cursor); + const hasNewline = newline !== -1; + const lineEnd = hasNewline ? newline : bytes.length; + const entry = parseJsonLine(bytes, cursor, lineEnd); + if (entry) { + visit(entry); + read = hasNewline ? newline + 1 : lineEnd; + } else if (hasNewline) { + read = newline + 1; + } else { + break; + } + cursor = hasNewline ? newline + 1 : lineEnd; + } + + return read; +} function parseSessionEntriesLenient(bytes: Uint8Array): { entries: SessionEntry[]; read: number } { const entries: SessionEntry[] = []; - let cursor = 0; - - while (cursor < bytes.length) { - const { values, error, read, done } = Bun.JSONL.parseChunk(bytes, cursor, bytes.length); - if (values.length > 0) { - entries.push(...(values as SessionEntry[])); - } - - if (error) { - const nextNewline = bytes.indexOf(LF, Math.max(read, cursor)); - if (nextNewline === -1) break; - cursor = nextNewline + 1; - continue; - } - - if (read <= cursor) break; - cursor = read; - if (done) break; - } - - return { entries, read: cursor }; + const read = visitSessionEntriesLenient(bytes, entry => entries.push(entry)); + return { entries, read }; } function scanLastServiceTier(bytes: Uint8Array): ServiceTier | undefined { - let cursor = 0; let currentServiceTier: ServiceTier | undefined; - - while (cursor < bytes.length) { - const { values, error, read, done } = Bun.JSONL.parseChunk(bytes, cursor, bytes.length); - for (const value of values as SessionEntry[]) { - if (isServiceTierChange(value)) currentServiceTier = value.serviceTier ?? undefined; - } - - if (error) { - const nextNewline = bytes.indexOf(LF, Math.max(read, cursor)); - if (nextNewline === -1) break; - cursor = nextNewline + 1; - continue; - } - - if (read <= cursor) break; - cursor = read; - if (done) break; - } - + visitSessionEntriesLenient(bytes, entry => { + if (isServiceTierChange(entry)) currentServiceTier = entry.serviceTier ?? undefined; + }); return currentServiceTier; } /** diff --git a/packages/stats/test/parser-large-session.test.ts b/packages/stats/test/parser-large-session.test.ts new file mode 100644 index 000000000..26de87057 --- /dev/null +++ b/packages/stats/test/parser-large-session.test.ts @@ -0,0 +1,83 @@ +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 { parseSessionFile } from "@oh-my-pi/omp-stats/parser"; +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-large-session-"); + const configDir = path.relative(os.homedir(), tempDir.join("config")); + process.env.PI_CONFIG_DIR = configDir; + setAgentDir(tempDir.join("agent")); +}); + +afterEach(() => { + vi.restoreAllMocks(); + if (originalConfigDir === undefined) { + delete process.env.PI_CONFIG_DIR; + } else { + process.env.PI_CONFIG_DIR = originalConfigDir; + } + setAgentDir(originalAgentDir); + tempDir?.removeSync(); + tempDir = null; +}); + +async function writeLargeSessionFile(): Promise { + const sessionDir = path.join(getSessionsDir(), "--tmp--large-session"); + await fs.mkdir(sessionDir, { recursive: true }); + const sessionFile = path.join(sessionDir, "session.jsonl"); + const timestamp = new Date().toISOString(); + const payload = "x".repeat(16 * 1024); + const lines: string[] = []; + for (let i = 0; i < 256; i++) { + lines.push( + JSON.stringify({ + type: "message", + id: `assistant-${i}`, + parentId: null, + timestamp, + message: { + role: "assistant", + content: [{ type: "text", text: payload }], + 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() + i, + duration: 10, + ttft: 5, + }, + }), + ); + } + await Bun.write(sessionFile, `${lines.join("\n")}\n`); + return sessionFile; +} + +describe("large session parsing", () => { + it("parses multi-megabyte JSONL without entering Bun.JSONL.parseChunk", async () => { + const sessionFile = await writeLargeSessionFile(); + vi.spyOn(Bun.JSONL, "parseChunk").mockImplementation(() => { + throw new Error("native JSONL parser unavailable"); + }); + + const result = await parseSessionFile(sessionFile); + + expect(result.stats).toHaveLength(256); + expect(result.newOffset).toBeGreaterThan(4 * 1024 * 1024); + }); +});