Merge PR #8598: fix(session): repair torn JSONL appends (@roboomp)

# Conflicts:
#	packages/coding-agent/src/session/session-loader.ts
#	packages/coding-agent/src/session/session-manager.ts
#	packages/coding-agent/test/session-loader-stream.test.ts
This commit is contained in:
can1357
2026-08-16 02:48:12 +02:00
11 changed files with 232 additions and 104 deletions
@@ -50,6 +50,7 @@ import {
logger,
postmortem,
prompt,
sanitizeText,
setProjectDir,
} from "@oh-my-pi/pi-utils";
import chalk from "@oh-my-pi/pi-utils/chalk";
@@ -1112,6 +1113,15 @@ export class InteractiveMode implements InteractiveModeContext {
// before initHooksAndCustomTools/#reconcileModeFromSession/#enterPlanMode —
// all of which can reach setSessionName during init.
this.#eventBusUnsubscribers.push(
this.sessionManager.onPersistenceError(error => {
const detail = truncateToWidth(
replaceTabs(sanitizeText(error.message)).replace(/[\r\n]+/g, " "),
TRUNCATE_LENGTHS.LINE,
);
this.showWarning(
`Session persistence failed: ${detail}. Unsaved entries remain in memory; persistence will retry on the next entry.`,
);
}),
this.sessionManager.onSessionNameChanged(() => {
setSessionTerminalTitle(this.sessionManager.getSessionName(), this.sessionManager.getCwd());
this.#handleSessionAccentInputsChanged();
@@ -2,14 +2,7 @@ import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { getBlobsDir, isEnoent, parseJsonlLenient } from "@oh-my-pi/pi-utils";
import { BlobStore, isBlobRef, resolveImageData, resolveImageDataUrl } from "./blob-store";
import { buildSessionContext } from "./session-context";
import {
type FileEntry,
type RawFileEntry,
SESSION_TITLE_SLOT_BYTES,
type SessionEntry,
type SessionHeader,
type SessionTitleSlotEntry,
} from "./session-entries";
import type { FileEntry, RawFileEntry, SessionEntry, SessionHeader } from "./session-entries";
import { migrateToCurrentVersion } from "./session-migrations";
import { isImageBlock, isImageDataPayload } from "./session-persistence";
import { FileSessionStorage, type SessionStorage } from "./session-storage";
@@ -33,6 +26,15 @@ export interface VisitEntriesFromFileStreamOptions {
yieldEveryBytes?: number;
/** Yield to the macrotask queue after this many entries have been visited. */
yieldEveryEntries?: number;
/** Called once for every malformed JSONL record skipped by the stream. */
onMalformedRecord?: () => void;
}
/** Parsed session entries plus corruption metadata needed by writable loaders. */
export interface SessionLoadResult {
entries: FileEntry[];
titleSlot: SessionTitleUpdate | undefined;
malformedRecords: number;
}
function splitTitleSlot(content: string): { body: string; slot: SessionTitleUpdate | undefined } {
@@ -61,14 +63,16 @@ function applyTitleSlot(entry: FileEntry | undefined, slot: SessionTitleUpdate |
}
/** Parse session JSONL while stripping and folding the optional fixed title slot. */
export function parseSessionContent(content: string): {
entries: FileEntry[];
titleSlot: SessionTitleUpdate | undefined;
} {
export function parseSessionContent(content: string): SessionLoadResult {
const { body, slot } = splitTitleSlot(content);
const entries = parseJsonlLenient<RawFileEntry>(body) as FileEntry[];
let malformedRecords = 0;
const entries = parseJsonlLenient<RawFileEntry>(body, {
onMalformedRecord: () => {
malformedRecords++;
},
}) as FileEntry[];
applyTitleSlot(entries[0], slot);
return { entries, titleSlot: slot };
return { entries, titleSlot: slot, malformedRecords };
}
/** Parse session JSONL and visit each entry without retaining prior entries. */
@@ -151,6 +155,15 @@ export async function visitEntriesFromFileStream(
// Malformed record: skip past the next newline and continue.
const nextNewline = buffer.indexOf(0x0a, read);
if (nextNewline === -1) break; // rest of the bad line not yet received
let nonWhitespace = false;
for (let index = read; index < nextNewline; index++) {
const byte = buffer[index];
if (byte !== 0x09 && byte !== 0x0d && byte !== 0x20) {
nonWhitespace = true;
break;
}
}
if (nonWhitespace) options.onMalformedRecord?.();
recordsSeen++;
buffer = buffer.subarray(nextNewline + 1);
if (recordsSeen >= maxRecords) {
@@ -211,33 +224,23 @@ export async function visitEntriesFromFileStream(
}
/** Exported for testing — the ≥8MiB streaming path (works on any file size). */
export async function loadEntriesFromFileStream(filePath: string): Promise<{
entries: FileEntry[];
titleSlot: SessionTitleUpdate | undefined;
}> {
export async function loadEntriesFromFileStream(filePath: string): Promise<SessionLoadResult> {
const entries: FileEntry[] = [];
const titleSlot = await visitEntriesFromFileStream(filePath, entry => {
entries.push(entry);
});
return { entries, titleSlot };
let malformedRecords = 0;
const titleSlot = await visitEntriesFromFileStream(
filePath,
entry => {
entries.push(entry);
},
{
onMalformedRecord: () => {
malformedRecords++;
},
},
);
return { entries, titleSlot, malformedRecords };
}
/** Read only the fixed-size head window to detect a physical title slot. */
export async function readTitleSlotFromFile(
filePath: string,
storage: SessionStorage = new FileSessionStorage(),
): Promise<SessionTitleSlotEntry | undefined> {
let head: string;
try {
[head] = await storage.readTextSlices(filePath, SESSION_TITLE_SLOT_BYTES, 0);
} catch (err) {
if (isEnoent(err)) return undefined;
throw err;
}
const newlineIndex = head.indexOf("\n");
if (newlineIndex < 0) return undefined;
return parseTitleSlotLine(head.slice(0, newlineIndex));
}
/** Exported for compaction.test.ts */
export function parseSessionEntries(content: string): FileEntry[] {
return parseSessionContent(content).entries;
@@ -247,25 +250,32 @@ function shouldStreamEntries(storage: SessionStorage, size: number): boolean {
return storage instanceof FileSessionStorage && size >= STREAM_LOAD_THRESHOLD_BYTES;
}
async function loadEntriesWithKnownSize(filePath: string, storage: SessionStorage, size: number): Promise<FileEntry[]> {
async function loadWithKnownSize(filePath: string, storage: SessionStorage, size: number): Promise<SessionLoadResult> {
const loaded = shouldStreamEntries(storage, size)
? await loadEntriesFromFileStream(filePath)
: parseSessionContent(await storage.readText(filePath));
const { entries } = loaded;
return isValidSessionHeader(entries[0]) ? entries : [];
return isValidSessionHeader(loaded.entries[0]) ? loaded : { ...loaded, entries: [] };
}
/** Exported for testing */
/** Load and validate a session while retaining malformed-record diagnostics. */
export async function loadSessionFile(
filePath: string,
storage: SessionStorage = new FileSessionStorage(),
): Promise<SessionLoadResult> {
try {
return await loadWithKnownSize(filePath, storage, storage.statSync(filePath).size);
} catch (err) {
if (isEnoent(err)) return { entries: [], titleSlot: undefined, malformedRecords: 0 };
throw err;
}
}
/** Load the valid entries from a session file, skipping malformed records. */
export async function loadEntriesFromFile(
filePath: string,
storage: SessionStorage = new FileSessionStorage(),
): Promise<FileEntry[]> {
try {
return await loadEntriesWithKnownSize(filePath, storage, storage.statSync(filePath).size);
} catch (err) {
if (isEnoent(err)) return [];
throw err;
}
return (await loadSessionFile(filePath, storage)).entries;
}
/**
@@ -290,7 +300,7 @@ export async function visitEntriesFromFile(
return;
}
for (const entry of await loadEntriesWithKnownSize(filePath, storage, size)) {
for (const entry of (await loadWithKnownSize(filePath, storage, size)).entries) {
if (visit(entry) === false) return;
}
}
@@ -61,8 +61,10 @@ import {
import { findMostRecentSession, listAllSessions, listSessions, type SessionInfo } from "./session-listing";
import {
loadEntriesFromFile,
loadSessionFile,
readTitleSlotFromFile,
resolveBlobRefsInEntries,
type SessionLoadResult,
visitEntriesFromFile,
} from "./session-loader";
import { generateId, migrateToCurrentVersion } from "./session-migrations";
@@ -545,6 +547,7 @@ export class SessionManager {
*/
#breadcrumbFresh = false;
#sessionNameChangedCallbacks = new Set<() => void>();
#persistenceErrorCallbacks = new Set<(error: Error) => void>();
private constructor(cwd: string, sessionDir: string, persist: boolean, storage: SessionStorage) {
this.#cwd = cwd;
@@ -586,6 +589,15 @@ export class SessionManager {
error: error.message,
stack: error.stack,
});
for (const callback of this.#persistenceErrorCallbacks) {
try {
callback(error);
} catch (callbackError) {
logger.warn("Session persistence error observer failed", {
error: toError(callbackError).message,
});
}
}
}
return this.#diskFailure;
@@ -843,6 +855,7 @@ export class SessionManager {
this.#diskTail = Promise.resolve();
this.#closeWriterEventually();
this.#storage.writeTextSync(targetPath, body);
this.#clearDiskError();
// Only mark the manager current when writing the active session path.
// Mid-move writes update the live relocation path; `#sessionFile` is
// still the pre-repoint source until moveTo repoints it.
@@ -935,7 +948,13 @@ export class SessionManager {
this.#atomicRewriteDirty = true;
return;
}
if (this.#diskFailure) throw this.#diskFailure;
if (this.#diskFailure) {
// The failed entry and any later entries remain in memory. A full
// replacement is the writability probe and restores all of them once
// transient storage pressure clears.
this.#fileIsCurrent = false;
this.#rewriteRequired = true;
}
// Lazy gate: a brand-new session is not written until it has an assistant
// message (or someone forced creation), so sessions that never produce
@@ -973,9 +992,9 @@ export class SessionManager {
// chain. Prefer appendSync so write failures latch `#diskFailure` before
// this call returns (not via a discarded rejected Promise after a later
// microtask). Callers stay non-throwing here — the core turn loop invokes
// appendMessage/appendCustomEntry without try/catch; flushSync/close and
// subsequent appends still throw the latched error. File writers apply
// each line to the OS page cache before return.
// appendMessage/appendCustomEntry without try/catch. A later entry retries
// all in-memory state through a full rewrite. File writers apply each line
// to the OS page cache before return.
// A mid-close writer leaves `#writer` undefined, so `#appendWriter` simply
// opens a fresh append handle and the entry still lands.
try {
@@ -984,16 +1003,28 @@ export class SessionManager {
if (writer.appendSync) {
writer.appendSync(line);
} else {
void writer.append(line).catch(err => this.#noteDiskFailure(err));
void writer.append(line).catch(err => {
this.#fileIsCurrent = false;
this.#rewriteRequired = true;
this.#noteDiskFailure(err);
});
}
} catch (err) {
this.#fileIsCurrent = false;
this.#rewriteRequired = true;
this.#noteDiskFailure(err);
}
}
async #persistTitleChangeEntry(entry: TitleChangeEntry, update: SessionTitleUpdate): Promise<void> {
if (!this.#persist || !this.#sessionFile) return;
if (this.#diskFailure) throw this.#diskFailure;
if (this.#diskFailure) {
this.#fileIsCurrent = false;
this.#rewriteRequired = true;
this.#rewriteSynchronously();
if (this.#diskFailure) throw this.#diskFailure;
return;
}
if (!this.#shouldHaveSessionFile()) {
this.#fileIsCurrent = false;
@@ -1317,7 +1348,7 @@ export class SessionManager {
await this.#setSessionFile(sessionFile);
}
async #setSessionFile(sessionFile: string, loadedEntries?: FileEntry[]): Promise<void> {
async #setSessionFile(sessionFile: string, loadedSession?: SessionLoadResult): Promise<void> {
await this.#drainAndCloseWriter();
this.#clearDiskError();
this.#draftOnlySessionCleanupArmed = false;
@@ -1326,8 +1357,8 @@ export class SessionManager {
this.#sessionFile = resolvedSessionFile;
this.#rememberBreadcrumb(this.#cwd, resolvedSessionFile);
const titleSlot = await readTitleSlotFromFile(resolvedSessionFile, this.#storage);
const fileEntries = loadedEntries ?? (await loadEntriesFromFile(resolvedSessionFile, this.#storage));
const loaded = loadedSession ?? (await loadSessionFile(resolvedSessionFile, this.#storage));
const { entries: fileEntries, titleSlot } = loaded;
if (fileEntries.length === 0) {
// Explicit but empty/missing path (e.g. --session flag): start fresh but
// keep the requested path and materialize the header immediately.
@@ -1361,7 +1392,7 @@ export class SessionManager {
this.#titleUpdatedAt = titleSlot?.updatedAt ?? header.timestamp;
this.#hasTitleSlot = titleSlot !== undefined;
this.#fileIsCurrent = true;
this.#rewriteRequired = migrated;
this.#rewriteRequired = migrated || loaded.malformedRecords > 0;
this.#forceFileCreation = true;
this.#artifactManager = null;
this.#artifactManagerSessionFile = null;
@@ -2017,6 +2048,14 @@ export class SessionManager {
};
}
/** Subscribe to persistence failures so hosts can surface lost-durability state. */
onPersistenceError(cb: (error: Error) => void): () => void {
this.#persistenceErrorCallbacks.add(cb);
return () => {
this.#persistenceErrorCallbacks.delete(cb);
};
}
/**
* Set the session display name.
* @param source "user" for explicit renames; "auto" for generated titles.
@@ -2624,8 +2663,8 @@ export class SessionManager {
storage: SessionStorage = new FileSessionStorage(),
options?: { initialCwd?: string; suppressBreadcrumb?: boolean },
): Promise<SessionManager> {
const loaded = await loadEntriesFromFile(filePath, storage);
const header = loaded.find(entry => entry.type === "session") as SessionHeader | undefined;
const loaded = await loadSessionFile(filePath, storage);
const header = loaded.entries.find(entry => entry.type === "session") as SessionHeader | undefined;
// Resume into the session's recorded cwd only when that directory still
// exists. A deleted project dir would make the constructor's #cwd — and the
// `setProjectDir` chdir interactive mode runs next — point at (and fail on)
@@ -128,14 +128,27 @@ class FileSessionStorageWriter implements SessionStorageWriter {
}
#writeNow(line: string): void {
const originalSize = fs.fstatSync(this.#fd).size;
const buf = Buffer.from(line, "utf-8");
let offset = 0;
while (offset < buf.length) {
const written = fs.writeSync(this.#fd, buf, offset, buf.length - offset);
if (written === 0) {
throw new Error("Short write");
try {
while (offset < buf.length) {
const written = fs.writeSync(this.#fd, buf, offset, buf.length - offset);
if (written === 0) {
throw new Error("Short write");
}
offset += written;
}
offset += written;
} catch (writeError) {
try {
fs.ftruncateSync(this.#fd, originalSize);
} catch (rollbackError) {
throw new AggregateError(
[toError(writeError), toError(rollbackError)],
"Session append failed and its partial bytes could not be rolled back",
);
}
throw writeError;
}
}