fix(stats): avoided native jsonl parser trap
Replaced the stats session parser's Bun.JSONL.parseChunk path with a lenient JS line scanner so large session files do not enter the native parser before /stats launches. Fixes #3733
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
/**
|
||||
|
||||
@@ -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<string> {
|
||||
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);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user