From aac5054c3f5993a313f2b37b8eefa3b3e0dee180 Mon Sep 17 00:00:00 2001 From: metaphorics <152830360+metaphorics@users.noreply.github.com> Date: Thu, 2 Jul 2026 17:25:50 +0900 Subject: [PATCH] perf(stats): stream session lookup and cache snapcompact JSONL reads getSessionEntry in packages/stats/src/parser.ts previously read the entire session file into memory and parsed every entry just to find one by ID. It now streams the file line-by-line via readLines and returns on the first matching entry, stopping I/O and parsing early. In packages/stats/src/gain-aggregator.ts, readSnapcompactRecords now stats snapcompact-savings.jsonl and caches parsed records keyed by mtimeMs and size, reusing them when the append-only journal has not changed. Cutoff filtering, per-request session:toolCallId deduplication, and project filtering are still applied per call. Verified with bun test test/gain-aggregator.test.ts in packages/stats (3 pass, 0 fail, 16 expect() calls, 1 file, 757.00ms) and bun run check:types. Closes #4251 --- packages/stats/CHANGELOG.md | 5 +++ packages/stats/src/gain-aggregator.ts | 64 ++++++++++++++++++++------- packages/stats/src/parser.ts | 17 +++---- 3 files changed, 59 insertions(+), 27 deletions(-) diff --git a/packages/stats/CHANGELOG.md b/packages/stats/CHANGELOG.md index 1cbaed0ca..088dc7404 100644 --- a/packages/stats/CHANGELOG.md +++ b/packages/stats/CHANGELOG.md @@ -2,6 +2,11 @@ ## [Unreleased] +### Changed + +- Stream session-entry lookup and cache snapcompact-savings.jsonl reads by mtime/size to avoid repeated full-file scans ([#4251](https://github.com/can1357/oh-my-pi/issues/4251)) + + ## [16.2.7] - 2026-06-30 ### Fixed diff --git a/packages/stats/src/gain-aggregator.ts b/packages/stats/src/gain-aggregator.ts index ba2c1ee4a..d874e9bfa 100644 --- a/packages/stats/src/gain-aggregator.ts +++ b/packages/stats/src/gain-aggregator.ts @@ -7,6 +7,8 @@ * Missing files are treated as zero records — never an error. */ +import type { Stats } from "node:fs"; +import * as fs from "node:fs/promises"; import * as path from "node:path"; import { getStatsDbPath, isEnoent, logger } from "@oh-my-pi/pi-utils"; import { getTimeRangeConfig } from "./aggregator"; @@ -149,37 +151,65 @@ async function readProjectsBySession(sessions: readonly string[]): Promise { const filePath = path.join(path.dirname(getStatsDbPath()), "snapcompact-savings.jsonl"); - let text: string; + + let stat: Stats; try { - text = await Bun.file(filePath).text(); + stat = await fs.stat(filePath); } catch (err) { if (isEnoent(err)) return { records: [], projects: new Set() }; - logger.debug("gain-aggregator: failed to read snapcompact-savings.jsonl", { err: String(err) }); + logger.debug("gain-aggregator: failed to stat snapcompact-savings.jsonl", { err: String(err) }); return { records: [], projects: new Set() }; } - const seen = new Set(); - const parsed: SnapcompactRecord[] = []; - for (const line of text.split("\n")) { - if (!line.trim()) continue; + const cacheKey = `${filePath}:${stat.mtimeMs}:${stat.size}`; + let parsed: SnapcompactRecord[]; + if (snapcompactCache?.key === cacheKey) { + parsed = snapcompactCache.records; + } else { + let text: string; try { - const rec = JSON.parse(line) as SnapcompactRecord; - if (cutoff !== null && rec.ts < cutoff) continue; - const key = `${rec.session}:${rec.toolCallId}`; - if (seen.has(key)) continue; - seen.add(key); - parsed.push(rec); - } catch { - /* skip malformed line */ + text = await Bun.file(filePath).text(); + } catch (readErr) { + if (isEnoent(readErr)) return { records: [], projects: new Set() }; + logger.debug("gain-aggregator: failed to read snapcompact-savings.jsonl", { err: String(readErr) }); + return { records: [], projects: new Set() }; } + + parsed = []; + for (const line of text.split("\n")) { + if (!line.trim()) continue; + try { + const rec = JSON.parse(line) as SnapcompactRecord; + parsed.push(rec); + } catch { + /* skip malformed line */ + } + } + snapcompactCache = { key: cacheKey, records: parsed }; } - const projectsBySession = await readProjectsBySession(parsed.map(rec => rec.session)); + const filtered = cutoff === null ? parsed : parsed.filter(rec => rec.ts >= cutoff); + const seen = new Set(); + const deduped: SnapcompactRecord[] = []; + for (const rec of filtered) { + const key = `${rec.session}:${rec.toolCallId}`; + if (seen.has(key)) continue; + seen.add(key); + deduped.push(rec); + } + const projectsBySession = await readProjectsBySession(deduped.map(rec => rec.session)); const projects = new Set(); const records: SnapcompactRecord[] = []; - for (const rec of parsed) { + for (const rec of deduped) { const sessionProjects = projectsBySession.get(rec.session); if (sessionProjects) { for (const sessionProject of sessionProjects) projects.add(sessionProject); diff --git a/packages/stats/src/parser.ts b/packages/stats/src/parser.ts index b77f9e549..2f837a649 100644 --- a/packages/stats/src/parser.ts +++ b/packages/stats/src/parser.ts @@ -7,7 +7,7 @@ import { resolveModelServiceTier, type ServiceTierByFamily, } from "@oh-my-pi/pi-ai"; -import { getSessionsDir, isEnoent } from "@oh-my-pi/pi-utils"; +import { getSessionsDir, isEnoent, readLines } from "@oh-my-pi/pi-utils"; import type { AgentType, MessageStats, @@ -350,19 +350,16 @@ export async function listAllSessionFiles(): Promise { * Find a specific entry in a session file. */ export async function getSessionEntry(sessionPath: string, entryId: string): Promise { - let bytes: Uint8Array; try { - bytes = await Bun.file(sessionPath).bytes(); + for await (const line of readLines(Bun.file(sessionPath).stream())) { + const entry = parseJsonLine(line, 0, line.length); + if (entry && "id" in entry && entry.id === entryId) { + return entry; + } + } } catch (err) { if (isEnoent(err)) return null; throw err; } - - const { entries } = parseSessionEntriesLenient(bytes); - for (const entry of entries) { - if ("id" in entry && entry.id === entryId) { - return entry; - } - } return null; }