From fd6700302837ec9474e6ec20e7ecebd4dccf9359 Mon Sep 17 00:00:00 2001 From: can1357 Date: Mon, 5 Jan 2026 02:23:36 +0100 Subject: [PATCH] feat(coding-agent): added streaming persistence with backpressure control - Added streaming NDJSON writer with backpressure handling for session persistence. - Added SessionManager.flush() method for explicit write synchronization before session operations. - Added content truncation for oversized blocks (500K char limit) to prevent memory exhaustion. - Replaced synchronous file appends with streaming NDJSON writer in session persistence. --- bun.lock | 15 +- packages/coding-agent/CHANGELOG.md | 5 +- packages/coding-agent/package.json | 1 + .../coding-agent/src/core/agent-session.ts | 7 + .../coding-agent/src/core/session-manager.ts | 407 +++++++++++++++++- .../src/modes/interactive/interactive-mode.ts | 3 + 6 files changed, 431 insertions(+), 7 deletions(-) diff --git a/bun.lock b/bun.lock index fdecb2e3f..ef58928f3 100644 --- a/bun.lock +++ b/bun.lock @@ -75,6 +75,7 @@ "marked": "^15.0.12", "minimatch": "^10.1.1", "nanoid": "^5.1.6", + "ndjson": "^2.0.0", "node-html-parser": "^6.1.13", "smol-toml": "^1.6.0", "strip-ansi": "^7.1.2", @@ -782,6 +783,8 @@ "json-schema-traverse": ["json-schema-traverse@1.0.0", "", {}, "sha512-NM8/P9n3XjXhIZn1lLhkFaACTOURQXjWhV4BA/RnOv8xvgqtqpAX9IO4mRQxSx1Rlo4tqzeqb0sOlruaOy3dug=="], + "json-stringify-safe": ["json-stringify-safe@5.0.1", "", {}, "sha512-ZClg6AaYvamvYEE82d3Iyd3vSSIjQ+odgjaTzRuO3s7toCdFKczob2i0zCh7JE8kWn17yvAWhUVxvqGwUalsRA=="], + "jsonschema": ["jsonschema@1.5.0", "", {}, "sha512-K+A9hhqbn0f3pJX17Q/7H6yQfD/5OXgdrR5UE12gMXCiN9D5Xq2o5mddV2QEcX/bjla99ASsAAQUyMCCRWAEhw=="], "jszip": ["jszip@3.10.1", "", { "dependencies": { "lie": "~3.3.0", "pako": "~1.0.2", "readable-stream": "~2.3.6", "setimmediate": "^1.0.5" } }, "sha512-xXDvecyTpGLrqFrvkrUSoxxfJI5AH7U8zxxtVclpsUtMCq4JQ290LY8AW5c7Ggnr/Y/oK+bQMbqK2qmtk3pN4g=="], @@ -880,6 +883,8 @@ "napi-build-utils": ["napi-build-utils@2.0.0", "", {}, "sha512-GEbrYkbfF7MoNaoh2iGG84Mnf/WZfB0GdGEsM8wz7Expx/LlWf5U8t9nvJKXSp3qr5IsEbK04cBGhol/KwOsWA=="], + "ndjson": ["ndjson@2.0.0", "", { "dependencies": { "json-stringify-safe": "^5.0.1", "minimist": "^1.2.5", "readable-stream": "^3.6.0", "split2": "^3.0.0", "through2": "^4.0.0" }, "bin": { "ndjson": "cli.js" } }, "sha512-nGl7LRGrzugTtaFcJMhLbpzJM6XdivmbkdlaGcrk/LXg2KL/YBC6z1g70xh0/al+oFuVFP8N8kiWRucmeEH/qQ=="], + "node-abi": ["node-abi@3.85.0", "", { "dependencies": { "semver": "^7.3.5" } }, "sha512-zsFhmbkAzwhTft6nd3VxcG0cvJsT70rL+BIGHWVq5fi6MwGrHwzqKaxXE+Hl2GmnGItnDKPPkO5/LQqjVkIdFg=="], "node-addon-api": ["node-addon-api@7.1.1", "", {}, "sha512-5m3bsyrjFWE1xf7nz7YXdN4udnVtXK6/Yfgn5qnahL6bCkf2yKt4k3nuTKAtT4r3IG8JNR2ncsIMdZuAzJjHQQ=="], @@ -1000,6 +1005,8 @@ "source-map-js": ["source-map-js@1.2.1", "", {}, "sha512-UXWMKhLOwVKb728IUtQPXxfYU+usdybtUrK/8uGE8CQMvrhOpwvzDBwj0QhSL7MQc7vIsISBG8VQ8+IDQxpfQA=="], + "split2": ["split2@3.2.2", "", { "dependencies": { "readable-stream": "^3.0.0" } }, "sha512-9NThjpgZnifTkJpzTZ7Eue85S49QwpNhZTq6GRJwObb6jnLFNGB7Qm73V5HewTROPyxD0C29xqmaI68bQtV+hg=="], + "stack-trace": ["stack-trace@0.0.10", "", {}, "sha512-KGzahc7puUKkzyMt+IqAep+TVNbKP+k2Lmwhub39m1AsTSkaDutx56aDCo+HLDzf/D26BIHTJWNiTG1KAJiQCg=="], "stackback": ["stackback@0.0.2", "", {}, "sha512-1XMJE5fQo1jGH6Y/7ebnwPOBEkIEnT4QF32d5R1+VXdXveM0IBMJt8zfaxX1P3QhVwrYe+576+jkANtSS2mBbw=="], @@ -1012,7 +1019,7 @@ "string-width-cjs": ["string-width@4.2.3", "", { "dependencies": { "emoji-regex": "^8.0.0", "is-fullwidth-code-point": "^3.0.0", "strip-ansi": "^6.0.1" } }, "sha512-wKyQRQpjJ0sIp62ErSZdGsjMJWsap5oRNihHhu6G7JVO/9jIB6UyevL+tXuOqrng8j/cxKTWyWUwvSTriiZz/g=="], - "string_decoder": ["string_decoder@1.1.1", "", { "dependencies": { "safe-buffer": "~5.1.0" } }, "sha512-n/ShnvDi6FHbbVfviro+WojiFzv+s8MPMHBczVePfUpDJLwoLT0ht1l4YwBCbi8pJAveEEdnkHyPyTP/mzRfwg=="], + "string_decoder": ["string_decoder@1.3.0", "", { "dependencies": { "safe-buffer": "~5.2.0" } }, "sha512-hkRX8U1WjJFd8LsDJ2yQ/wWWxaopEsABU1XfkM8A+j0+85JAGppt16cr1Whg6KIbb4okU6Mql6BOj+uup/wKeA=="], "strip-ansi": ["strip-ansi@7.1.2", "", { "dependencies": { "ansi-regex": "^6.0.1" } }, "sha512-gmBGslpoQJtgnMAvOVqGZpEz9dyoKTCzy2nfz/n8aIFhN/jCE/rCmcxabB6jOOHV+0WNnylOxaxBQPSvcWklhA=="], @@ -1044,6 +1051,8 @@ "thenify-all": ["thenify-all@1.6.0", "", { "dependencies": { "thenify": ">= 3.1.0 < 4" } }, "sha512-RNxQH/qI8/t3thXJDwcstUO4zeqo64+Uy/+sNVRBx4Xn2OX+OZ9oP+iJnNFqplFra2ZUVeKCSa2oVWi3T4uVmA=="], + "through2": ["through2@4.0.2", "", { "dependencies": { "readable-stream": "3" } }, "sha512-iOqSav00cVxEEICeD7TjLB1sueEL+81Wpzp2bY17uZjZN0pWZPuo4suZ/61VujxmqSGFfgOcNuTZ85QJwNZQpw=="], + "tinybench": ["tinybench@2.9.0", "", {}, "sha512-0+DUvqWMValLmha6lr4kD8iAMK1HzV0/aKnCtWb9v9641TnP/MFb7Pc2bxoxQjTXAErryXVgUOfv2YqNllqGeg=="], "tinyexec": ["tinyexec@0.3.2", "", {}, "sha512-KQQR9yN7R5+OSwaK0XQoj22pwHoTlgYqmUscPYoknOoWCWfj/5/ABTMRi69FrKU5ffPVh5QcFikpWJI/P1ocHA=="], @@ -1210,6 +1219,8 @@ "string-width-cjs/strip-ansi": ["strip-ansi@6.0.1", "", { "dependencies": { "ansi-regex": "^5.0.1" } }, "sha512-Y38VPSHcqkFrCpFnQ9vuSXmquuv5oXOKpGeT6aGrr3o3Gc9AlVa6JBfUSOCnbxGGZF+/0ooI7KrPuUSztUdU5A=="], + "string_decoder/safe-buffer": ["safe-buffer@5.2.1", "", {}, "sha512-rp3So07KcdmmKbGvgaNxQSJr7bGVSVk5S9Eq1F+ppbRo70+YeaDxkw5Dd8NPN+GD6bjnYm2VuPuCXmpuYvmCXQ=="], + "strip-ansi-cjs/ansi-regex": ["ansi-regex@5.0.1", "", {}, "sha512-quJQXlTSUGL2LH9SUXo8VwsY4soanhgo6LNSm84E1LBcE8s3O0wpdiRzyR9z/ZZJMlMWv37qOOb9pdJlMUEKFQ=="], "tunnel-agent/safe-buffer": ["safe-buffer@5.2.1", "", {}, "sha512-rp3So07KcdmmKbGvgaNxQSJr7bGVSVk5S9Eq1F+ppbRo70+YeaDxkw5Dd8NPN+GD6bjnYm2VuPuCXmpuYvmCXQ=="], @@ -1258,6 +1269,8 @@ "form-data/mime-types/mime-db": ["mime-db@1.52.0", "", {}, "sha512-sPU4uV7dYlvtWJxwwxHD0PuihVNiE7TyAbQ5SWxDCB9mUYvOgroQOwYQQOKPJ8CIbE+1ETVlOoK1UC2nU3gYvg=="], + "jszip/readable-stream/string_decoder": ["string_decoder@1.1.1", "", { "dependencies": { "safe-buffer": "~5.1.0" } }, "sha512-n/ShnvDi6FHbbVfviro+WojiFzv+s8MPMHBczVePfUpDJLwoLT0ht1l4YwBCbi8pJAveEEdnkHyPyTP/mzRfwg=="], + "rimraf/glob/jackspeak": ["jackspeak@3.4.3", "", { "dependencies": { "@isaacs/cliui": "^8.0.2" }, "optionalDependencies": { "@pkgjs/parseargs": "^0.11.0" } }, "sha512-OGlZQpz2yfahA/Rd1Y8Cd9SIEsqvXkLVoSw/cgwhnhFMDbsQFeZYoJJ7bIZBS9BcamUW96asq/npPWugM+RQBw=="], "rimraf/glob/minimatch": ["minimatch@9.0.5", "", { "dependencies": { "brace-expansion": "^2.0.1" } }, "sha512-G6T0ZX48xgozx7587koeX9Ys2NYy6Gmv//P89sEte9V9whIapMNF4idKxnW2QtCcLiTWlb/wfCabAtAFWhhBow=="], diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index c0bc7962a..ef885adf3 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,9 +1,10 @@ # Changelog ## [Unreleased] - ### Added +- Added streaming NDJSON writer for session persistence with proper backpressure handling +- Added `flush()` method to SessionManager for explicit control over pending write completion - Added `/exit` slash command to exit the application from interactive mode - Added fuzzy path matching suggestions when read tool encounters file-not-found errors, showing closest matches using Levenshtein distance - Added `status.shadowed` symbol for theme customization to properly indicate shadowed extension state @@ -22,6 +23,7 @@ ### Changed +- Changed session persistence to use streaming writes instead of synchronous file appends for better performance - Changed read tool to automatically redirect to ls when given a directory path instead of a file - Changed tool description prompts to be more concise with clearer usage guidelines and structured formatting - Moved tool description prompts from inline strings to external markdown files in `src/prompts/tools/` directory for better maintainability @@ -50,6 +52,7 @@ ### Fixed +- Fixed session persistence to truncate oversized content blocks before writing to prevent memory exhaustion - Fixed extension list and inspector panel to use correct symbols for disabled and shadowed states instead of reusing unrelated status icons - Fixed token counting for subagent progress to handle different usage object formats (camelCase and snake_case) - Fixed image file handling by adding 20MB size limit to prevent memory issues during serialization diff --git a/packages/coding-agent/package.json b/packages/coding-agent/package.json index e3415e85a..7dfdda905 100644 --- a/packages/coding-agent/package.json +++ b/packages/coding-agent/package.json @@ -53,6 +53,7 @@ "marked": "^15.0.12", "minimatch": "^10.1.1", "nanoid": "^5.1.6", + "ndjson": "^2.0.0", "node-html-parser": "^6.1.13", "smol-toml": "^1.6.0", "strip-ansi": "^7.1.2", diff --git a/packages/coding-agent/src/core/agent-session.ts b/packages/coding-agent/src/core/agent-session.ts index 593df6192..8950418eb 100644 --- a/packages/coding-agent/src/core/agent-session.ts +++ b/packages/coding-agent/src/core/agent-session.ts @@ -841,6 +841,7 @@ export class AgentSession { this._disconnectFromAgent(); await this.abort(); this.agent.reset(); + await this.sessionManager.flush(); this.sessionManager.newSession(options); this._queuedMessages = []; this._reconnectToAgent(); @@ -1636,6 +1637,9 @@ export class AgentSession { await this.abort(); this._queuedMessages = []; + // Flush pending writes before switching + await this.sessionManager.flush(); + // Set new session this.sessionManager.setSessionFile(sessionPath); @@ -1714,6 +1718,9 @@ export class AgentSession { skipConversationRestore = result?.skipConversationRestore ?? false; } + // Flush pending writes before branching + await this.sessionManager.flush(); + if (!selectedEntry.parentId) { this.sessionManager.newSession(); } else { diff --git a/packages/coding-agent/src/core/session-manager.ts b/packages/coding-agent/src/core/session-manager.ts index 51e8fb33d..81b489746 100644 --- a/packages/coding-agent/src/core/session-manager.ts +++ b/packages/coding-agent/src/core/session-manager.ts @@ -1,6 +1,7 @@ import { appendFileSync, closeSync, + createWriteStream, existsSync, mkdirSync, openSync, @@ -9,11 +10,13 @@ import { readSync, statSync, writeFileSync, + type WriteStream, } from "node:fs"; import { join, resolve } from "node:path"; import type { AgentMessage } from "@oh-my-pi/pi-agent-core"; -import type { ImageContent, Message, TextContent } from "@oh-my-pi/pi-ai"; +import type { ImageContent, Message, TextContent, ThinkingContent, ToolCall } from "@oh-my-pi/pi-ai"; import { nanoid } from "nanoid"; +import ndjson from "ndjson"; import { getAgentDir as getDefaultAgentDir } from "../config"; import { type BashExecutionMessage, @@ -571,6 +574,350 @@ function formatTimeAgo(date: Date): string { return date.toLocaleDateString(); } +const MAX_PERSIST_CONTENT_CHARS = 500_000; +const PERSIST_TRUNCATION_NOTICE = "\n\n[Session persistence truncated large content]"; + +type TextOrImageContent = TextContent | ImageContent; +type AssistantContent = TextContent | ThinkingContent | ToolCall; +type FileMentionFile = { path: string; content: string; lineCount: number }; + +function truncateStringForPersistence( + text: string, + maxChars: number = MAX_PERSIST_CONTENT_CHARS, +): { text: string; truncated: boolean } { + if (text.length <= maxChars) { + return { text, truncated: false }; + } + const limit = Math.max(0, maxChars - PERSIST_TRUNCATION_NOTICE.length); + return { + text: `${text.slice(0, limit)}${PERSIST_TRUNCATION_NOTICE}`, + truncated: true, + }; +} + +function textImageContentExceedsLimit(content: TextOrImageContent[], maxChars: number): boolean { + let total = 0; + for (const block of content) { + if (block.type === "text") { + total += block.text.length; + } else if (block.type === "image") { + total += block.data.length; + } + if (total > maxChars) return true; + } + return false; +} + +function truncateTextImageContent( + content: TextOrImageContent[], + maxChars: number, +): { content: TextOrImageContent[]; truncated: boolean } { + if (!textImageContentExceedsLimit(content, maxChars)) { + return { content, truncated: false }; + } + + const result: TextOrImageContent[] = []; + let remaining = maxChars; + + for (const block of content) { + if (block.type === "text") { + if (block.text.length <= remaining) { + result.push(block); + remaining -= block.text.length; + continue; + } + + const limit = Math.max(0, remaining - PERSIST_TRUNCATION_NOTICE.length); + result.push({ + ...block, + text: `${block.text.slice(0, limit)}${PERSIST_TRUNCATION_NOTICE}`, + }); + return { content: result, truncated: true }; + } + + if (block.type === "image") { + if (block.data.length <= remaining) { + result.push(block); + remaining -= block.data.length; + continue; + } + + result.push({ + type: "text", + text: `[Image omitted in session persistence]${PERSIST_TRUNCATION_NOTICE}`, + }); + return { content: result, truncated: true }; + } + } + + return { content: result, truncated: true }; +} + +function assistantContentExceedsLimit(content: AssistantContent[], maxChars: number): boolean { + let total = 0; + for (const block of content) { + if (block.type === "text") { + total += block.text.length; + } else if (block.type === "thinking") { + total += block.thinking.length; + } + if (total > maxChars) return true; + } + return false; +} + +function truncateAssistantContent( + content: AssistantContent[], + maxChars: number, +): { content: AssistantContent[]; truncated: boolean } { + if (!assistantContentExceedsLimit(content, maxChars)) { + return { content, truncated: false }; + } + + const result: AssistantContent[] = []; + let remaining = maxChars; + + for (const block of content) { + if (block.type === "text") { + if (block.text.length <= remaining) { + result.push(block); + remaining -= block.text.length; + continue; + } + + const limit = Math.max(0, remaining - PERSIST_TRUNCATION_NOTICE.length); + result.push({ + ...block, + text: `${block.text.slice(0, limit)}${PERSIST_TRUNCATION_NOTICE}`, + }); + return { content: result, truncated: true }; + } + + if (block.type === "thinking") { + if (block.thinking.length <= remaining) { + result.push(block); + remaining -= block.thinking.length; + continue; + } + + const limit = Math.max(0, remaining - PERSIST_TRUNCATION_NOTICE.length); + result.push({ + ...block, + thinking: `${block.thinking.slice(0, limit)}${PERSIST_TRUNCATION_NOTICE}`, + }); + return { content: result, truncated: true }; + } + + if (block.type === "toolCall") { + if (remaining <= 0) { + result.push({ type: "text", text: PERSIST_TRUNCATION_NOTICE }); + return { content: result, truncated: true }; + } + result.push(block); + } + } + + return { content: result, truncated: true }; +} + +function truncateContentValue( + content: string | TextOrImageContent[], + maxChars: number, +): { content: string | TextOrImageContent[]; truncated: boolean } { + if (typeof content === "string") { + const result = truncateStringForPersistence(content, maxChars); + return { content: result.text, truncated: result.truncated }; + } + return truncateTextImageContent(content, maxChars); +} + +function fileMentionsExceedLimit(files: FileMentionFile[], maxChars: number): boolean { + let total = 0; + for (const file of files) { + total += file.content.length; + if (total > maxChars) return true; + } + return false; +} + +function truncateFileMentions( + files: FileMentionFile[], + maxChars: number, +): { files: FileMentionFile[]; truncated: boolean } { + if (!fileMentionsExceedLimit(files, maxChars)) { + return { files, truncated: false }; + } + + const result: FileMentionFile[] = []; + let remaining = maxChars; + + for (const file of files) { + if (file.content.length <= remaining) { + result.push(file); + remaining -= file.content.length; + continue; + } + + const limit = Math.max(0, remaining - PERSIST_TRUNCATION_NOTICE.length); + const truncatedContent = `${file.content.slice(0, limit)}${PERSIST_TRUNCATION_NOTICE}`; + result.push({ + ...file, + content: truncatedContent, + lineCount: truncatedContent.split("\n").length, + }); + return { files: result, truncated: true }; + } + + return { files: result, truncated: true }; +} + +function truncateDetailsForPersistence(details: unknown, maxChars: number): unknown { + if (!details || typeof details !== "object" || Array.isArray(details)) { + return details; + } + + let changed = false; + const result: Record = {}; + for (const [key, value] of Object.entries(details as Record)) { + if (typeof value === "string") { + const truncated = truncateStringForPersistence(value, maxChars); + result[key] = truncated.text; + if (truncated.truncated) changed = true; + } else { + result[key] = value; + } + } + + return changed ? result : details; +} + +function truncateAgentMessageForPersistence(message: AgentMessage): AgentMessage { + if (message.role === "user") { + const result = truncateContentValue(message.content, MAX_PERSIST_CONTENT_CHARS); + if (!result.truncated) return message; + return { ...message, content: result.content }; + } + + if (message.role === "assistant") { + const result = truncateAssistantContent(message.content, MAX_PERSIST_CONTENT_CHARS); + if (!result.truncated) return message; + return { ...message, content: result.content }; + } + + if (message.role === "toolResult") { + const contentResult = truncateTextImageContent(message.content, MAX_PERSIST_CONTENT_CHARS); + const detailsResult = truncateDetailsForPersistence(message.details, MAX_PERSIST_CONTENT_CHARS); + const contentChanged = contentResult.truncated; + const detailsChanged = detailsResult !== message.details; + + if (!contentChanged && !detailsChanged) return message; + + if (contentChanged && detailsChanged) { + return { ...message, content: contentResult.content, details: detailsResult }; + } + if (contentChanged) { + return { ...message, content: contentResult.content }; + } + return { ...message, details: detailsResult }; + } + + if (message.role === "hookMessage") { + const result = truncateContentValue(message.content, MAX_PERSIST_CONTENT_CHARS); + if (!result.truncated) return message; + return { ...message, content: result.content }; + } + + if (message.role === "bashExecution") { + const result = truncateStringForPersistence(message.output, MAX_PERSIST_CONTENT_CHARS); + if (!result.truncated) return message; + return { ...message, output: result.text }; + } + + if (message.role === "fileMention") { + const result = truncateFileMentions(message.files, MAX_PERSIST_CONTENT_CHARS); + if (!result.truncated) return message; + return { ...message, files: result.files }; + } + + return message; +} + +function prepareEntryForPersistence(entry: FileEntry): FileEntry { + if (entry.type === "message") { + const message = truncateAgentMessageForPersistence(entry.message); + if (message === entry.message) return entry; + return { ...entry, message }; + } + + if (entry.type === "custom_message") { + const result = truncateContentValue(entry.content, MAX_PERSIST_CONTENT_CHARS); + if (!result.truncated) return entry; + return { ...entry, content: result.content }; + } + + if (entry.type === "compaction" || entry.type === "branch_summary") { + const result = truncateStringForPersistence(entry.summary, MAX_PERSIST_CONTENT_CHARS); + if (!result.truncated) return entry; + return { ...entry, summary: result.text }; + } + + return entry; +} + +class NdjsonFileWriter { + private encoder: ReturnType; + private writeStream: WriteStream; + private closed = false; + private error: Error | undefined; + + constructor(path: string) { + this.encoder = ndjson.stringify(); + this.writeStream = createWriteStream(path, { flags: "a" }); + this.encoder.pipe(this.writeStream); + + this.encoder.on("error", (err: Error) => { + this.error = err; + }); + this.writeStream.on("error", (err: Error) => { + this.error = err; + }); + } + + write(entry: FileEntry): void { + if (this.closed) throw new Error("Writer closed"); + if (this.error) throw this.error; + this.encoder.write(entry); + } + + flush(): Promise { + if (this.error) return Promise.reject(this.error); + return new Promise((resolve) => { + // Cork/uncork forces immediate flush through transform + this.encoder.cork(); + process.nextTick(() => { + this.encoder.uncork(); + // Wait for writeStream to drain + if (this.writeStream.writableNeedDrain) { + this.writeStream.once("drain", resolve); + } else { + resolve(); + } + }); + }); + } + + close(): Promise { + if (this.closed) return Promise.resolve(); + this.closed = true; + + return new Promise((resolve, reject) => { + this.writeStream.on("finish", resolve); + this.writeStream.on("error", reject); + this.encoder.end(); + }); + } +} + /** Get recent sessions for display in welcome screen */ export function getRecentSessions(sessionDir: string, limit = 3): RecentSessionInfo[] { return getSortedSessions(sessionDir).slice(0, limit); @@ -608,6 +955,8 @@ export class SessionManager { private labelsById: Map = new Map(); private leafId: string | null = null; private usageStatistics: UsageStatistics = { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, cost: 0 }; + private persistWriter: NdjsonFileWriter | undefined; + private persistWriterPath: string | undefined; private constructor(cwd: string, sessionDir: string, sessionFile: string | undefined, persist: boolean) { this.cwd = cwd; @@ -626,6 +975,7 @@ export class SessionManager { /** Switch to a different session file (used for resume and branching) */ setSessionFile(sessionFile: string): void { + this._closePersistWriterSync(); this.sessionFile = resolve(sessionFile); if (existsSync(this.sessionFile)) { this.fileEntries = loadEntriesFromFile(this.sessionFile); @@ -645,6 +995,7 @@ export class SessionManager { } newSession(options?: NewSessionOptions): string | undefined { + this._closePersistWriterSync(); this.sessionId = nanoid(); const timestamp = new Date().toISOString(); const header: SessionHeader = { @@ -696,16 +1047,55 @@ export class SessionManager { } } + private _ensurePersistWriter(): NdjsonFileWriter | undefined { + if (!this.persist || !this.sessionFile) return undefined; + if (this.persistWriter && this.persistWriterPath === this.sessionFile) return this.persistWriter; + this._closePersistWriterSync(); + this.persistWriter = new NdjsonFileWriter(this.sessionFile); + this.persistWriterPath = this.sessionFile; + return this.persistWriter; + } + + private _closePersistWriterSync(): void { + if (this.persistWriter) { + // Fire-and-forget close - use flush() for guaranteed completion + this.persistWriter.close(); + this.persistWriter = undefined; + } + this.persistWriterPath = undefined; + } + + private async _closePersistWriter(): Promise { + if (this.persistWriter) { + await this.persistWriter.close(); + this.persistWriter = undefined; + } + this.persistWriterPath = undefined; + } + private _rewriteFile(): void { if (!this.persist || !this.sessionFile) return; - const content = `${this.fileEntries.map((e) => JSON.stringify(e)).join("\n")}\n`; - writeFileSync(this.sessionFile, content); + this._closePersistWriterSync(); + writeFileSync(this.sessionFile, ""); + const writer = this._ensurePersistWriter(); + if (!writer) return; + for (const entry of this.fileEntries) { + const persistedEntry = prepareEntryForPersistence(entry); + writer.write(persistedEntry); + } } isPersisted(): boolean { return this.persist; } + /** Flush pending writes to disk. Call before switching sessions or on shutdown. */ + async flush(): Promise { + if (this.persistWriter) { + await this.persistWriter.flush(); + } + } + getCwd(): string { return this.cwd; } @@ -742,6 +1132,8 @@ export class SessionManager { // Update the session file header with the title (if already flushed) if (this.persist && this.sessionFile && existsSync(this.sessionFile)) { + // Close writer to flush pending writes before modifying file + this._closePersistWriterSync(); try { const content = readFileSync(this.sessionFile, "utf-8"); const lines = content.split("\n"); @@ -765,13 +1157,18 @@ export class SessionManager { const hasAssistant = this.fileEntries.some((e) => e.type === "message" && e.message.role === "assistant"); if (!hasAssistant) return; + const writer = this._ensurePersistWriter(); + if (!writer) return; + if (!this.flushed) { for (const e of this.fileEntries) { - appendFileSync(this.sessionFile, `${JSON.stringify(e)}\n`); + const persistedEntry = prepareEntryForPersistence(e); + writer.write(persistedEntry); } this.flushed = true; } else { - appendFileSync(this.sessionFile, `${JSON.stringify(entry)}\n`); + const persistedEntry = prepareEntryForPersistence(entry); + writer.write(persistedEntry); } } diff --git a/packages/coding-agent/src/modes/interactive/interactive-mode.ts b/packages/coding-agent/src/modes/interactive/interactive-mode.ts index f0342e7da..b283cb5d0 100644 --- a/packages/coding-agent/src/modes/interactive/interactive-mode.ts +++ b/packages/coding-agent/src/modes/interactive/interactive-mode.ts @@ -1447,6 +1447,9 @@ export class InteractiveMode { * Emits shutdown event to hooks and tools, then exits. */ private async shutdown(): Promise { + // Flush pending session writes before shutdown + await this.sessionManager.flush(); + // Emit shutdown event to hooks const hookRunner = this.session.hookRunner; if (hookRunner?.hasHandlers("session_shutdown")) {