refactor(session): simplify streamed entry visits

This commit is contained in:
Aleksandr Khaustov
2026-08-15 23:54:30 +04:00
parent cbbaceac36
commit 9c695cf98a
3 changed files with 39 additions and 88 deletions
@@ -46,27 +46,20 @@ function isValidSessionHeader(entry: FileEntry | undefined): entry is SessionHea
return entry?.type === "session" && typeof entry.id === "string";
}
function applyTitleSlot(header: SessionHeader, slot: SessionTitleUpdate): void {
function applyTitleSlot(entry: FileEntry | undefined, slot: SessionTitleUpdate | undefined): void {
if (!slot || !isValidSessionHeader(entry)) return;
if (slot.title && slot.title.length > 0) {
header.title = slot.title;
entry.title = slot.title;
} else {
delete header.title;
delete entry.title;
}
if (slot.source) {
header.titleSource = slot.source;
entry.titleSource = slot.source;
} else {
delete header.titleSource;
delete entry.titleSource;
}
}
function foldTitleSlot(entries: FileEntry[], slot: SessionTitleUpdate | undefined): FileEntry[] {
if (!slot || entries.length === 0) return entries;
const header = entries[0];
if (!isValidSessionHeader(header)) return entries;
applyTitleSlot(header, slot);
return entries;
}
/** Parse session JSONL while stripping and folding the optional fixed title slot. */
export function parseSessionContent(content: string): {
entries: FileEntry[];
@@ -74,17 +67,19 @@ export function parseSessionContent(content: string): {
} {
const { body, slot } = splitTitleSlot(content);
const entries = parseJsonlLenient<RawFileEntry>(body) as FileEntry[];
return { entries: foldTitleSlot(entries, slot), titleSlot: slot };
applyTitleSlot(entries[0], slot);
return { entries, titleSlot: slot };
}
/** Parse session JSONL and visit each entry without retaining prior entries. */
export async function visitEntriesFromFileStream(
filePath: string,
visit: (entry: FileEntry, titleSlot: SessionTitleUpdate | undefined) => void | boolean,
visit: (entry: FileEntry) => void | boolean,
options: VisitEntriesFromFileStreamOptions = {},
): Promise<SessionTitleUpdate | undefined> {
let titleSlot: SessionTitleUpdate | undefined;
let sawFirstLine = false;
let sawFirstEntry = false;
let bytesSinceYield = 0;
let entriesSinceYield = 0;
let recordsSeen = 0;
@@ -129,8 +124,13 @@ export async function visitEntriesFromFileStream(
stopped = true;
break;
}
const entry = value as FileEntry;
if (!sawFirstEntry) {
sawFirstEntry = true;
applyTitleSlot(entry, titleSlot);
}
try {
if (visit(value as FileEntry, titleSlot) === false) {
if (visit(entry) === false) {
stopped = true;
break;
}
@@ -219,7 +219,7 @@ export async function loadEntriesFromFileStream(filePath: string): Promise<{
const titleSlot = await visitEntriesFromFileStream(filePath, entry => {
entries.push(entry);
});
return { entries: foldTitleSlot(entries, titleSlot), titleSlot };
return { entries, titleSlot };
}
/** Read only the fixed-size head window to detect a physical title slot. */
@@ -277,37 +277,21 @@ export async function visitEntriesFromFile(
visit: (entry: FileEntry) => void | boolean,
storage: SessionStorage = new FileSessionStorage(),
): Promise<void> {
let visitorThrew = false;
const callVisitor = (entry: FileEntry): void | boolean => {
try {
const size = storage.statSync(filePath).size;
if (shouldStreamEntries(storage, size)) {
let sawFirstEntry = false;
await visitEntriesFromFileStream(filePath, entry => {
if (!sawFirstEntry) {
sawFirstEntry = true;
if (!isValidSessionHeader(entry)) return false;
}
return visit(entry);
} catch (err) {
visitorThrew = true;
throw err;
}
};
try {
const size = storage.statSync(filePath).size;
if (shouldStreamEntries(storage, size)) {
let sawFirstEntry = false;
await visitEntriesFromFileStream(filePath, (entry, titleSlot) => {
if (!sawFirstEntry) {
sawFirstEntry = true;
if (!isValidSessionHeader(entry)) return false;
if (titleSlot) applyTitleSlot(entry, titleSlot);
}
return callVisitor(entry);
});
return;
}
});
return;
}
for (const entry of await loadEntriesWithKnownSize(filePath, storage, size)) {
if (callVisitor(entry) === false) return;
}
} catch (err) {
if (visitorThrew) throw err;
if (isEnoent(err)) return;
throw err;
for (const entry of await loadEntriesWithKnownSize(filePath, storage, size)) {
if (visit(entry) === false) return;
}
}
@@ -4,7 +4,6 @@ import * as os from "node:os";
import * as path from "node:path";
import type { FileEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import * as sessionLoader from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { FileSessionStorage } from "@oh-my-pi/pi-coding-agent/session/session-storage";
import { serializeTitleSlot } from "@oh-my-pi/pi-coding-agent/session/session-title-slot";
// Parity contract for the ≥8MiB streaming loader (now Bun.JSONL-based): it must
@@ -17,13 +16,6 @@ import { serializeTitleSlot } from "@oh-my-pi/pi-coding-agent/session/session-ti
const ISO = "2026-06-29T12:00:00.000Z";
const HEADER = { type: "session", version: 3, id: "s1", timestamp: ISO, cwd: "/tmp" };
const LARGE_SESSION_BYTES = 9 * 1024 * 1024;
class LargeFileSessionStorage extends FileSessionStorage {
override statSync(filePath: string) {
return { ...super.statSync(filePath), size: LARGE_SESSION_BYTES };
}
}
const msg = (id: string, parentId: string, text: string) => ({
type: "message",
@@ -101,6 +93,7 @@ describe("loadEntriesFromFileStream (Bun.JSONL parity)", () => {
expect(titleSlot?.title).toBe("Visitor");
expect(entryIds(visited)).toEqual(["s1", "m1", "m2"]);
expect(visited[0]).toMatchObject({ title: "Visitor", titleSource: "user" });
});
it("visits a large journal before reading its tail", async () => {
@@ -182,26 +175,6 @@ describe("loadEntriesFromFileStream (Bun.JSONL parity)", () => {
).rejects.toBe(failure);
});
it("propagates ENOENT errors through the routed visitor", async () => {
const file = await writeTemp(`${JSON.stringify(HEADER)}\n`);
const failure = Object.assign(new Error("visitor failed"), { code: "ENOENT" });
await expect(
sessionLoader.visitEntriesFromFile(file, () => {
throw failure;
}),
).rejects.toBe(failure);
await expect(
sessionLoader.visitEntriesFromFile(
file,
() => {
throw failure;
},
new LargeFileSessionStorage(),
),
).rejects.toBe(failure);
});
it("matches parseSessionContent on title slot + valid + malformed + blank lines", async () => {
const slotLine = serializeTitleSlot({ title: "Hello world", source: "user", updatedAt: ISO });
// title slot | header | valid | blank | malformed | valid | malformed-no-newline-at-EOF
@@ -1,7 +1,6 @@
import { afterEach, describe, expect, it, spyOn } from "bun:test";
import { afterEach, describe, expect, it } from "bun:test";
import * as path from "node:path";
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
import * as sessionLoader from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { FileSessionStorage, MemorySessionStorage } from "@oh-my-pi/pi-coding-agent/session/session-storage";
import { TempDir } from "@oh-my-pi/pi-utils";
@@ -23,6 +22,9 @@ class LargeFileSessionStorage extends FileSessionStorage {
override statSync(filePath: string) {
return { ...super.statSync(filePath), size: LARGE_SESSION_BYTES };
}
override async readText(): Promise<string> {
throw new Error("Large sessions must stream");
}
}
function assistantMessage(text: string) {
@@ -76,7 +78,7 @@ describe("SessionManager.peekSessionInit", () => {
expect(peek?.init?.restrictToolNames).toBe(true);
});
it("routes the latest large-file metadata through the entry visitor", async () => {
it("streams large file-backed sessions without a full read", async () => {
const cwd = makeTempDir("@pi-peek-stream-");
const manager = SessionManager.create(cwd, path.join(cwd, "sessions"));
const sessionFile = manager.getSessionFile();
@@ -86,17 +88,9 @@ describe("SessionManager.peekSessionInit", () => {
manager.appendSessionInit({ systemPrompt: "second", task: "task", tools: ["read"], spawns: "" });
manager.appendMessage(assistantMessage("journal tail"));
const storage = new LargeFileSessionStorage();
const visitEntries = spyOn(sessionLoader, "visitEntriesFromFile");
try {
const peek = await SessionManager.peekSessionInit(sessionFile, storage);
expect(peek?.cwd).toBe(manager.getCwd());
expect(peek?.init?.systemPrompt).toBe("second");
expect(visitEntries).toHaveBeenCalledTimes(1);
} finally {
visitEntries.mockRestore();
}
const peek = await SessionManager.peekSessionInit(sessionFile, new LargeFileSessionStorage());
expect(peek?.cwd).toBe(manager.getCwd());
expect(peek?.init?.systemPrompt).toBe("second");
});
it("preserves non-file storage behavior", async () => {