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.
This commit is contained in:
@@ -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=="],
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<string, unknown> = {};
|
||||
for (const [key, value] of Object.entries(details as Record<string, unknown>)) {
|
||||
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<typeof ndjson.stringify>;
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<string, string> = 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<void> {
|
||||
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<void> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1447,6 +1447,9 @@ export class InteractiveMode {
|
||||
* Emits shutdown event to hooks and tools, then exits.
|
||||
*/
|
||||
private async shutdown(): Promise<void> {
|
||||
// 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")) {
|
||||
|
||||
Reference in New Issue
Block a user