fix(coding-agent): fixed session JSONL persistence races for immediate and mid-close writes
- Fixed initial assistant persistence by synchronously materializing in-memory entries and keeping a writer open. - Fixed mid-close append handling to write entries through a one-shot sync writer instead of queuing a rewrite. - Fixed persist-task gating by tracking pending writes before starting immediate persistence.
This commit is contained in:
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed session JSONL persistence so the first assistant turn materializes the file synchronously, leaves the append writer open, and writes later entries with a sync append writer even during writer-close races instead of waiting on a queued rewrite.
|
||||
|
||||
## [15.12.5] - 2026-06-13
|
||||
### Changed
|
||||
|
||||
|
||||
@@ -2028,6 +2028,7 @@ export class SessionManager {
|
||||
#persistWriter: NdjsonFileWriter | undefined;
|
||||
#persistWriterPath: string | undefined;
|
||||
#persistChain: Promise<void> = Promise.resolve();
|
||||
#pendingPersistTasks = 0;
|
||||
#persistError: Error | undefined;
|
||||
#persistErrorReported = false;
|
||||
#artifactManager: ArtifactManager | null = null;
|
||||
@@ -2430,14 +2431,18 @@ export class SessionManager {
|
||||
}
|
||||
|
||||
#queuePersistTask(task: () => Promise<void>, options?: { ignoreError?: boolean }): Promise<void> {
|
||||
this.#pendingPersistTasks++;
|
||||
const next = this.#persistChain.then(async () => {
|
||||
if (this.#persistError && !options?.ignoreError) throw this.#persistError;
|
||||
await task();
|
||||
});
|
||||
this.#persistChain = next.catch(err => {
|
||||
const tracked = next.finally(() => {
|
||||
this.#pendingPersistTasks--;
|
||||
});
|
||||
this.#persistChain = tracked.catch(err => {
|
||||
this.#recordPersistError(err);
|
||||
});
|
||||
return next;
|
||||
return tracked;
|
||||
}
|
||||
|
||||
#ensurePersistWriter(): NdjsonFileWriter | undefined {
|
||||
@@ -2564,6 +2569,85 @@ export class SessionManager {
|
||||
}
|
||||
}
|
||||
|
||||
#writeEntriesToWriterSync(writer: NdjsonFileWriter): void {
|
||||
for (const entry of this.#fileEntries) {
|
||||
writer.writeSync(prepareEntryForPersistenceSync(entry, this.#blobStore));
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* First durable-session write: synchronously materialize the current in-memory
|
||||
* entries into the real session file, then keep that writer open so every
|
||||
* later append goes through the hot `writeSync` path. This intentionally does
|
||||
* not fsync; it only guarantees the JSONL bytes are handed to the kernel
|
||||
* before `appendMessage()` returns.
|
||||
*/
|
||||
#tryStartImmediatePersistWriter(): boolean {
|
||||
if (!this.#sessionFile) return false;
|
||||
if (this.#flushed || this.#needsFullRewriteOnNextPersist) return false;
|
||||
if (this.#persistWriter || this.#pendingPersistTasks > 0) return false;
|
||||
if (this.storage.existsSync(this.#sessionFile)) return false;
|
||||
|
||||
const writer = new NdjsonFileWriter(this.storage, this.#sessionFile, {
|
||||
flags: "w",
|
||||
onError: err => this.#recordPersistError(err),
|
||||
});
|
||||
try {
|
||||
this.#writeEntriesToWriterSync(writer);
|
||||
} catch (err) {
|
||||
void writer.close().catch(() => {});
|
||||
throw toError(err);
|
||||
}
|
||||
|
||||
this.#persistWriter = writer;
|
||||
this.#persistWriterPath = this.#sessionFile;
|
||||
this.#needsFullRewriteOnNextPersist = false;
|
||||
this.#flushed = true;
|
||||
return true;
|
||||
}
|
||||
|
||||
#rewriteFileInPlaceSync(): void {
|
||||
if (!this.#sessionFile) return;
|
||||
const previousWriter = this.#persistWriter;
|
||||
this.#persistWriter = undefined;
|
||||
this.#persistWriterPath = undefined;
|
||||
if (previousWriter) {
|
||||
void previousWriter.close().catch(err => {
|
||||
this.#recordPersistError(err);
|
||||
});
|
||||
}
|
||||
|
||||
const writer = new NdjsonFileWriter(this.storage, this.#sessionFile, {
|
||||
flags: "w",
|
||||
onError: err => this.#recordPersistError(err),
|
||||
});
|
||||
try {
|
||||
this.#writeEntriesToWriterSync(writer);
|
||||
} catch (err) {
|
||||
void writer.close().catch(() => {});
|
||||
throw toError(err);
|
||||
}
|
||||
|
||||
this.#persistWriter = writer;
|
||||
this.#persistWriterPath = this.#sessionFile;
|
||||
this.#needsFullRewriteOnNextPersist = false;
|
||||
this.#flushed = true;
|
||||
}
|
||||
|
||||
#appendEntryWithOneShotWriter(entry: SessionEntry): void {
|
||||
if (!this.#sessionFile) return;
|
||||
const writer = new NdjsonFileWriter(this.storage, this.#sessionFile, {
|
||||
onError: err => this.#recordPersistError(err),
|
||||
});
|
||||
try {
|
||||
writer.writeSync(prepareEntryForPersistenceSync(entry, this.#blobStore));
|
||||
} finally {
|
||||
void writer.close().catch(err => {
|
||||
this.#recordPersistError(err);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
async #rewriteFile(): Promise<void> {
|
||||
if (!this.persist || !this.#sessionFile) return;
|
||||
await this.#queuePersistTask(async () => {
|
||||
@@ -2992,12 +3076,15 @@ export class SessionManager {
|
||||
}
|
||||
|
||||
if (this.#needsFullRewriteOnNextPersist || !this.#flushed) {
|
||||
// Cold path: rewrite the whole file atomically. Async — the writer is
|
||||
// closed/reopened and every entry is re-prepared. Errors flow through
|
||||
// `#persistChain` → `#recordPersistError`; we swallow the rejection
|
||||
// here to avoid an unhandled rejection when the persist dir races with
|
||||
// test-level tempDir cleanup.
|
||||
this.#rewriteFile().catch(() => {});
|
||||
// First assistant turn after lazy session creation: write the full file
|
||||
// synchronously and leave the append writer open. After this returns,
|
||||
// subsequent entries are never pending only in memory.
|
||||
try {
|
||||
if (this.#tryStartImmediatePersistWriter()) return;
|
||||
this.#rewriteFileInPlaceSync();
|
||||
} catch (err) {
|
||||
this.#recordPersistError(err);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -3009,11 +3096,10 @@ export class SessionManager {
|
||||
try {
|
||||
const writer = this.#ensurePersistWriter();
|
||||
if (!writer) {
|
||||
// `#ensurePersistWriter` returns undefined here only when the cached
|
||||
// writer is mid-close (the `!persist`/`!sessionFile` cases are
|
||||
// rejected above). Route through `#rewriteFile` so the entry — which
|
||||
// is already in `#fileEntries` — persists once the close drains.
|
||||
this.#rewriteFile().catch(() => {});
|
||||
// The cached writer is mid-close. Write the new entry through a
|
||||
// short-lived append writer instead of queueing a rewrite, so the
|
||||
// JSONL line is on disk before this append returns.
|
||||
this.#appendEntryWithOneShotWriter(entry);
|
||||
return;
|
||||
}
|
||||
const persistedEntry = prepareEntryForPersistenceSync(entry, this.#blobStore);
|
||||
|
||||
@@ -207,6 +207,11 @@ describe("SessionManager close/appendMessage race", () => {
|
||||
});
|
||||
}).not.toThrow();
|
||||
|
||||
const sessionFile = sm.getSessionFile();
|
||||
if (!sessionFile) throw new Error("Expected session file");
|
||||
const duringCloseContent = await storage.readText(sessionFile);
|
||||
expect(duringCloseContent).toContain('"content":"during-close"');
|
||||
|
||||
// Drain everything.
|
||||
await settle(closePromise, storage);
|
||||
// Pre-fix `flush()` rejects with the stashed Error("Writer closed").
|
||||
|
||||
@@ -0,0 +1,90 @@
|
||||
import { afterEach, describe, expect, it } from "bun:test";
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
|
||||
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
|
||||
import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
|
||||
const tempDirs: TempDir[] = [];
|
||||
|
||||
function makeTempDir(prefix: string): string {
|
||||
const dir = TempDir.createSync(prefix);
|
||||
tempDirs.push(dir);
|
||||
return dir.path();
|
||||
}
|
||||
|
||||
afterEach(async () => {
|
||||
await Promise.all(tempDirs.splice(0).map(dir => dir.remove()));
|
||||
});
|
||||
|
||||
function assistantMessage(text: string) {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
|
||||
if (!model) throw new Error("Expected built-in anthropic model to exist");
|
||||
return {
|
||||
role: "assistant" as const,
|
||||
content: [{ type: "text" as const, text }],
|
||||
api: model.api,
|
||||
provider: model.provider,
|
||||
model: model.id,
|
||||
usage: {
|
||||
input: 0,
|
||||
output: 0,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 0,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
},
|
||||
stopReason: "stop" as const,
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
}
|
||||
|
||||
function readJsonl(file: string): Array<Record<string, unknown>> {
|
||||
return fs
|
||||
.readFileSync(file, "utf8")
|
||||
.trimEnd()
|
||||
.split("\n")
|
||||
.filter(Boolean)
|
||||
.map(line => JSON.parse(line) as Record<string, unknown>);
|
||||
}
|
||||
|
||||
function messageRole(entry: Record<string, unknown>): string | undefined {
|
||||
const message = entry.message;
|
||||
if (!message || typeof message !== "object") return undefined;
|
||||
const role = (message as { role?: unknown }).role;
|
||||
return typeof role === "string" ? role : undefined;
|
||||
}
|
||||
|
||||
function messageContent(entry: Record<string, unknown>): unknown {
|
||||
const message = entry.message;
|
||||
if (!message || typeof message !== "object") return undefined;
|
||||
return (message as { content?: unknown }).content;
|
||||
}
|
||||
|
||||
describe("SessionManager immediate JSONL persistence", () => {
|
||||
it("writes the first assistant turn and later entries before appendMessage returns", () => {
|
||||
const cwd = makeTempDir("@pi-immediate-cwd-");
|
||||
const sessionDir = path.join(cwd, "sessions");
|
||||
const manager = SessionManager.create(cwd, sessionDir);
|
||||
const sessionFile = manager.getSessionFile();
|
||||
if (!sessionFile) throw new Error("Expected a persisted session file path");
|
||||
|
||||
manager.appendMessage({ role: "user", content: "queued before assistant", timestamp: Date.now() });
|
||||
expect(fs.existsSync(sessionFile)).toBe(false);
|
||||
|
||||
manager.appendMessage(assistantMessage("hello"));
|
||||
expect(fs.existsSync(sessionFile)).toBe(true);
|
||||
|
||||
let entries = readJsonl(sessionFile);
|
||||
expect(entries).toHaveLength(3);
|
||||
expect(messageRole(entries[1] ?? {})).toBe("user");
|
||||
expect(messageRole(entries[2] ?? {})).toBe("assistant");
|
||||
|
||||
manager.appendMessage({ role: "user", content: "written immediately", timestamp: Date.now() });
|
||||
|
||||
entries = readJsonl(sessionFile);
|
||||
expect(entries).toHaveLength(4);
|
||||
expect(messageRole(entries[3] ?? {})).toBe("user");
|
||||
expect(messageContent(entries[3] ?? {})).toBe("written immediately");
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user