From 3d4ace30d8af0a8038cb03d2ef0b4012a672ff95 Mon Sep 17 00:00:00 2001 From: can1357 Date: Wed, 7 Jan 2026 00:21:06 +0100 Subject: [PATCH] fix(coding-agent): resolved memory leak by streaming events to disk - Fixed memory accumulation by streaming events directly to disk instead of storing in memory. - Changed subprocess usage tracking to accumulate incrementally from message_end events. - Fixed session persistence to exclude transient streaming data (partialJson, jsonlEvents). - Removed parseSubagentUsage() function in favor of pre-accumulated usage data. --- packages/coding-agent/CHANGELOG.md | 10 +++- .../coding-agent/src/core/session-manager.ts | 7 +++ .../src/core/tools/task/executor.ts | 56 ++++++++++++++---- .../coding-agent/src/core/tools/task/index.ts | 58 ++++--------------- .../coding-agent/src/core/tools/task/types.ts | 5 +- 5 files changed, 74 insertions(+), 62 deletions(-) diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 3c4c968a9..f7c0e6bd3 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,6 +1,14 @@ # Changelog ## [Unreleased] +### Changed + +- Changed subprocess usage tracking to accumulate incrementally from message_end events rather than parsing stored events after completion + +### Fixed + +- Fixed memory accumulation in task subprocess by streaming events directly to disk instead of storing in memory +- Fixed session persistence to exclude transient streaming data (partialJson, jsonlEvents) that was causing unnecessary storage bloat ## [3.21.0] - 2026-01-06 @@ -1756,4 +1764,4 @@ Initial public release. - Git branch display in footer - Message queueing during streaming responses - OAuth integration for Gmail and Google Calendar access -- HTML export with syntax highlighting and collapsible sections +- HTML export with syntax highlighting and collapsible sections \ No newline at end of file diff --git a/packages/coding-agent/src/core/session-manager.ts b/packages/coding-agent/src/core/session-manager.ts index df1e19277..60946c6a4 100644 --- a/packages/coding-agent/src/core/session-manager.ts +++ b/packages/coding-agent/src/core/session-manager.ts @@ -698,6 +698,13 @@ async function truncateForPersistence(obj: T, key?: string): Promise { let changed = false; const result: Record = {}; for (const [k, v] of Object.entries(obj as Record)) { + // Strip transient/redundant properties that shouldn't be persisted + // - partialJson: streaming accumulator for tool call JSON parsing + // - jsonlEvents: raw subprocess streaming events (already saved to artifact files) + if (k === "partialJson" || k === "jsonlEvents") { + changed = true; + continue; + } const newV = await truncateForPersistence(v, k); result[k] = newV; if (newV !== v) changed = true; diff --git a/packages/coding-agent/src/core/tools/task/executor.ts b/packages/coding-agent/src/core/tools/task/executor.ts index b6951c9dd..1a04fab69 100644 --- a/packages/coding-agent/src/core/tools/task/executor.ts +++ b/packages/coding-agent/src/core/tools/task/executor.ts @@ -279,7 +279,17 @@ export async function runSubprocess(options: ExecutorOptions): Promise { @@ -301,7 +311,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise): void { target.cost.total += cost.total; } -function parseSubagentUsage(events: string[] | undefined): Usage | undefined { - if (!events || events.length === 0) return undefined; - - const totals = createUsageTotals(); - let hasUsage = false; - - for (const line of events) { - let event: unknown; - try { - event = JSON.parse(line); - } catch { - continue; - } - - if (!event || typeof event !== "object") continue; - const record = event as Record; - if (record.type !== "message_end") continue; - - const message = record.message; - if (!message || typeof message !== "object") continue; - const msgRecord = message as Record; - if (msgRecord.role !== "assistant") continue; - if (msgRecord.stopReason === "aborted" || msgRecord.stopReason === "error") continue; - - const usage = msgRecord.usage; - if (!usage || typeof usage !== "object") continue; - - addUsageTotals(totals, usage as Partial); - hasUsage = true; - } - - return hasUsage ? totals : undefined; -} - /** Session context interface */ interface SessionContext { getSessionFile: () => string | null; @@ -477,31 +443,29 @@ export async function createTaskTool( }); }); + // Aggregate usage from executor results (already accumulated incrementally) const aggregatedUsage = createUsageTotals(); let hasAggregatedUsage = false; - const resultsWithUsage = results.map((result) => { - const usage = parseSubagentUsage(result.jsonlEvents); - if (usage) { - addUsageTotals(aggregatedUsage, usage); + for (const result of results) { + if (result.usage) { + addUsageTotals(aggregatedUsage, result.usage); hasAggregatedUsage = true; - return { ...result, usage }; } - return result; - }); + } // Collect output paths (artifacts already written by executor in real-time) const outputPaths: string[] = []; - for (const result of resultsWithUsage) { + for (const result of results) { if (result.artifactPaths) { outputPaths.push(result.artifactPaths.outputPath); } } // Build final output - match plugin format - const successCount = resultsWithUsage.filter((r) => r.exitCode === 0).length; + const successCount = results.filter((r) => r.exitCode === 0).length; const totalDuration = Date.now() - startTime; - const summaries = resultsWithUsage.map((r) => { + const summaries = results.map((r) => { const status = r.exitCode === 0 ? "completed" : `failed (exit ${r.exitCode})`; const output = r.output.trim() || r.stderr.trim() || "(no output)"; const preview = output.split("\n").slice(0, 5).join("\n"); @@ -520,12 +484,12 @@ export async function createTaskTool( skippedSelfRecursion > 1 ? "s" : "" } skipped - self-recursion blocked)` : ""; - const outputIds = resultsWithUsage.map((r) => `${r.agent}_${r.index}`); + const outputIds = results.map((r) => `${r.agent}_${r.index}`); const outputHint = hasOutputTool && outputIds.length > 0 ? `\n\nUse output tool for full logs: output ids ${outputIds.join(", ")}` : ""; - const summary = `${successCount}/${resultsWithUsage.length} succeeded${skippedNote} [${formatDuration( + const summary = `${successCount}/${results.length} succeeded${skippedNote} [${formatDuration( totalDuration, )}]\n\n${summaries.join("\n\n---\n\n")}${outputHint}`; @@ -538,7 +502,7 @@ export async function createTaskTool( content: [{ type: "text", text: summary }], details: { projectAgentsDir, - results: resultsWithUsage, + results: results, totalDurationMs: totalDuration, usage: hasAggregatedUsage ? aggregatedUsage : undefined, outputPaths, diff --git a/packages/coding-agent/src/core/tools/task/types.ts b/packages/coding-agent/src/core/tools/task/types.ts index 963631ad0..9eed3605d 100644 --- a/packages/coding-agent/src/core/tools/task/types.ts +++ b/packages/coding-agent/src/core/tools/task/types.ts @@ -123,10 +123,9 @@ export interface SingleResult { modelOverride?: string; error?: string; aborted?: boolean; - jsonlEvents?: string[]; - artifactPaths?: { inputPath: string; outputPath: string; jsonlPath?: string }; - /** Aggregated usage from the subprocess, if available. */ + /** Aggregated usage from the subprocess, accumulated incrementally from message_end events. */ usage?: Usage; + artifactPaths?: { inputPath: string; outputPath: string; jsonlPath?: string }; /** Data extracted by registered subprocess tool handlers (keyed by tool name) */ extractedToolData?: Record; /** Output metadata for Output tool integration */