fix(session): looped title fallback + drained backend on close

SessionManager.#persistTitleChangeEntry's catch fallback previously did a single-shot atomic rewrite: any prompt/tool appended while it awaited was fenced with #atomicRewriteDirty=true but never re-serialized. Extracted the fenced-rewrite do-while loop into #runFencedAtomicRewrite and used it from both #rewriteAtomically and #persistTitleChangeEntry, so fenced entries during either path are captured before the task resolves.

Added SessionStorage.drain(): for FileSessionStorage and MemorySessionStorage it is a no-op; IndexedSessionStorage already had one and now conforms to the interface. SessionManager.flush() and close() await it so a graceful shutdown does not exit while a fire-and-forget writeTextSync publish (queued by flushSync on an indexed backend) is still on the wire — reducing the residual publish-window race for Redis/SQL where the backend cannot be aborted mid-flight.

Regression covers the title fallback loop: TitleFallbackPausingStorage forces updateSessionTitle to throw, pauses the fallback's writeTextAtomic, appends a message and a custom entry during the pause, and asserts (a) both fenced entries land on the current JSONL, (b) the final title is applied, and (c) writeTextAtomicCalls >= 2 proving the loop iterated.

Fixes #4338
This commit is contained in:
roboomp
2026-07-02 20:32:27 +00:00
parent f614ec1537
commit babb731cf5
4 changed files with 148 additions and 25 deletions
@@ -573,30 +573,45 @@ export class SessionManager {
const startEpoch = this.#diskEpoch;
await this.#scheduleDiskWork(
async () => {
this.#atomicRewriteActive = true;
try {
do {
this.#atomicRewriteDirty = false;
await this.#closeWriterHandle();
const sessionFile = this.#sessionFile;
if (!sessionFile) return;
if (this.#diskEpoch !== startEpoch) return;
await this.#storage.writeTextAtomic(sessionFile, this.#fileBody(), {
commitGuard: () => this.#diskEpoch === startEpoch,
});
if (this.#diskEpoch !== startEpoch) return;
} while (this.#atomicRewriteDirty);
if (await this.#runFencedAtomicRewrite(startEpoch)) {
this.#fileIsCurrent = true;
this.#rewriteRequired = false;
this.#hasTitleSlot = true;
} finally {
this.#atomicRewriteActive = false;
}
},
{ epoch: startEpoch },
);
}
/**
* Shared fenced atomic-rewrite loop used by `#rewriteAtomically` and the
* `#persistTitleChangeEntry` fallback. Holds `#atomicRewriteActive` across
* the writer close and the full-file replace, and loops on
* `#atomicRewriteDirty` so any fenced append that lands during the rewrite
* is captured before the task resolves. Returns `false` when the disk epoch
* moved (a superseding synchronous rewrite has taken over) so callers skip
* their post-publish state updates.
*/
async #runFencedAtomicRewrite(epoch: number): Promise<boolean> {
this.#atomicRewriteActive = true;
try {
do {
this.#atomicRewriteDirty = false;
await this.#closeWriterHandle();
const sessionFile = this.#sessionFile;
if (!sessionFile) return false;
if (this.#diskEpoch !== epoch) return false;
await this.#storage.writeTextAtomic(sessionFile, this.#fileBody(), {
commitGuard: () => this.#diskEpoch === epoch,
});
if (this.#diskEpoch !== epoch) return false;
} while (this.#atomicRewriteDirty);
return true;
} finally {
this.#atomicRewriteActive = false;
}
}
#appendToSessionFile(entry: SessionEntry): void {
if (!this.#persist || !this.#sessionFile) return;
if (this.#diskFailure) throw this.#diskFailure;
@@ -669,16 +684,7 @@ export class SessionManager {
await this.#storage.updateSessionTitle(sessionFile, update);
if (this.#diskEpoch === epoch) this.#fileIsCurrent = true;
} catch {
this.#atomicRewriteActive = true;
try {
await this.#closeWriterHandle();
await this.#storage.writeTextAtomic(sessionFile, this.#fileBody(), {
commitGuard: () => this.#diskEpoch === epoch,
});
} finally {
this.#atomicRewriteActive = false;
}
if (this.#diskEpoch !== epoch) return;
if (!(await this.#runFencedAtomicRewrite(epoch))) return;
this.#clearDiskError();
this.#fileIsCurrent = true;
this.#rewriteRequired = false;
@@ -1087,6 +1093,10 @@ export class SessionManager {
await this.#scheduleDiskWork(async () => {
if (this.#writer?.isOpen()) await this.#writer.flush();
});
// Drain any fire-and-forget backing writes (e.g. `writeTextSync` queued
// on IndexedSessionStorage during `flushSync`) so callers relying on
// flush() see the write durably visible to readers.
await this.#storage.drain();
if (this.#diskFailure) throw this.#diskFailure;
}
@@ -1116,6 +1126,10 @@ export class SessionManager {
if (hadWriter || (this.#sessionFile && this.#storage.existsSync(this.#sessionFile)))
this.#fileIsCurrent = true;
});
// Wait for any queued backing writes (IndexedSessionStorage per-path
// tail) to become durable so a graceful shutdown does not exit while
// a fire-and-forget publish is still on the wire.
await this.#storage.drain();
if (this.#diskFailure) throw this.#diskFailure;
}
@@ -65,6 +65,15 @@ export interface SessionStorage {
unlink(path: string): Promise<void>;
deleteSessionWithArtifacts(sessionPath: string): Promise<void>;
openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter;
/**
* Wait for every backing write scheduled by this storage to become durably
* visible. Sync backends (file, memory) return immediately because their
* writes complete in-body; async backends (Redis/SQL via
* {@link IndexedSessionStorage}) await their per-path queues so a caller
* driving a graceful shutdown does not exit while a fire-and-forget
* `writeTextSync` publish is still on the wire.
*/
drain(): Promise<void>;
}
// FinalizationRegistry to clean up leaked file descriptors
@@ -378,6 +387,12 @@ export class FileSessionStorage implements SessionStorage {
return fs.promises.unlink(path);
}
drain(): Promise<void> {
// File writes complete synchronously in-body via fs.writeFileSync /
// fs.renameSync, so there is no queued work to await.
return Promise.resolve();
}
openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter {
return new FileSessionStorageWriter(path, options);
}
@@ -709,6 +724,10 @@ export class MemorySessionStorage implements SessionStorage {
return Promise.resolve();
}
drain(): Promise<void> {
return Promise.resolve();
}
openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter {
return new MemorySessionStorageWriter(this, path, options);
}
@@ -6,6 +6,7 @@ import {
type SessionStorageWriter,
type WriteTextAtomicOptions,
} from "@oh-my-pi/pi-coding-agent/session/session-storage";
import type { SessionTitleUpdate } from "@oh-my-pi/pi-coding-agent/session/session-title-slot";
interface DetachableWriter extends SessionStorageWriter {
detach(): void;
@@ -371,3 +372,89 @@ describe("SessionManager atomic rewrite fence spans writer.close()", () => {
expect(storage.detachedLines).toEqual([]);
});
});
class TitleFallbackPausingStorage extends MemorySessionStorage {
readonly writeStarted = Promise.withResolvers<void>();
readonly allowWrite = Promise.withResolvers<void>();
writeTextAtomicCalls = 0;
failNextUpdateTitle = false;
override async updateSessionTitle(path: string, update: SessionTitleUpdate): Promise<void> {
if (this.failNextUpdateTitle) {
this.failNextUpdateTitle = false;
throw new Error("updateSessionTitle forced failure");
}
return super.updateSessionTitle(path, update);
}
override async writeTextAtomic(path: string, content: string, options?: WriteTextAtomicOptions): Promise<void> {
this.writeTextAtomicCalls += 1;
this.writeStarted.resolve();
await this.allowWrite.promise;
if (options?.commitGuard && !options.commitGuard()) return;
this.writeTextSync(path, content);
}
}
describe("SessionManager title-change fallback fenced-append durability", () => {
it("loops on the dirty flag so fenced appends during the fallback rewrite persist", async () => {
const storage = new TitleFallbackPausingStorage();
const sessionManager = SessionManager.create("/cwd", "/sessions", storage);
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
if (!model) throw new Error("Expected built-in anthropic model");
// Materialize the session on disk with a title slot present so a later
// setSessionName takes the append-then-updateSessionTitle try branch
// instead of the up-front #rewriteAtomically fallback.
sessionManager.appendMessage({
role: "assistant",
content: [{ type: "text", text: "seed response" }],
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",
timestamp: Date.now(),
});
await sessionManager.flush();
await sessionManager.setSessionName("initial title", "user", "seed");
await sessionManager.flush();
expect(storage.writeTextAtomicCalls).toBe(0);
// Force the try branch to fail so the catch runs the atomic-rewrite loop.
storage.failNextUpdateTitle = true;
const rename = sessionManager.setSessionName("second title", "user", "test");
await storage.writeStarted.promise;
// Fenced appends during the paused fallback rewrite: pre-fix these
// would be marked dirty and dropped from the serialized body because
// the fallback never looped on that flag.
sessionManager.appendMessage({
role: "user",
content: "during title fallback",
timestamp: Date.now(),
});
sessionManager.appendCustomEntry("during_title_fallback_custom", { reason: "test" });
storage.allowWrite.resolve();
await rename;
await sessionManager.flush();
const sessionFile = sessionManager.getSessionFile();
if (!sessionFile) throw new Error("Expected session file");
const content = await storage.readText(sessionFile);
expect(content).toContain('"title":"second title"');
expect(content).toContain("during title fallback");
expect(content).toContain('"customType":"during_title_fallback_custom"');
// Loop must have executed at least twice: first pass paused, dirty from
// the fenced appends triggers a second pass that includes them.
expect(storage.writeTextAtomicCalls).toBeGreaterThanOrEqual(2);
});
});
@@ -117,6 +117,9 @@ class CloseHoldingStorage implements SessionStorage {
deleteSessionWithArtifacts(sessionPath: string): Promise<void> {
return this.#inner.deleteSessionWithArtifacts(sessionPath);
}
drain(): Promise<void> {
return this.#inner.drain();
}
}
/** Drive microtasks while releasing every parked close until `promise` settles. */