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.
This commit is contained in:
can1357
2026-01-07 00:21:06 +01:00
parent 839d9244e2
commit 3d4ace30d8
5 changed files with 74 additions and 62 deletions
+9 -1
View File
@@ -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
@@ -698,6 +698,13 @@ async function truncateForPersistence<T>(obj: T, key?: string): Promise<T> {
let changed = false;
const result: Record<string, unknown> = {};
for (const [k, v] of Object.entries(obj as Record<string, unknown>)) {
// 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;
@@ -279,7 +279,17 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
let finalOutput = "";
let resolved = false;
let pendingTermination = false; // Set when shouldTerminate fires, wait for message_end
const jsonlEvents: string[] = [];
// Accumulate usage incrementally from message_end events (no memory for streaming events)
const accumulatedUsage = {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
let hasUsage = false;
// Handle abort signal
const onAbort = () => {
@@ -301,7 +311,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
try {
const event = JSON.parse(line);
jsonlEvents.push(line);
// Events are written to subtask session file by subprocess - no need to accumulate in memory
const now = Date.now();
switch (event.type) {
@@ -389,19 +399,43 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
}
case "message_end": {
// Extract final text content from completed message
const messageContent = event.message?.content || event.content;
if (messageContent && Array.isArray(messageContent)) {
for (const block of messageContent) {
if (block.type === "text" && block.text) {
output += block.text;
// Extract text from assistant and toolResult messages (not user prompts)
const role = event.message?.role;
if (role === "assistant") {
const messageContent = event.message?.content || event.content;
if (messageContent && Array.isArray(messageContent)) {
for (const block of messageContent) {
if (block.type === "text" && block.text) {
output += block.text;
}
}
}
}
// Extract usage (prefer message.usage, fallback to event.usage)
// Extract and accumulate usage (prefer message.usage, fallback to event.usage)
const messageUsage = event.message?.usage || event.usage;
if (messageUsage) {
// Accumulate tokens across messages (not overwrite)
// Only count assistant messages (not tool results, etc.)
const role = event.message?.role;
if (
role === "assistant" &&
event.message?.stopReason !== "aborted" &&
event.message?.stopReason !== "error"
) {
hasUsage = true;
accumulatedUsage.input += messageUsage.input ?? 0;
accumulatedUsage.output += messageUsage.output ?? 0;
accumulatedUsage.cacheRead += messageUsage.cacheRead ?? 0;
accumulatedUsage.cacheWrite += messageUsage.cacheWrite ?? 0;
accumulatedUsage.totalTokens += messageUsage.totalTokens ?? 0;
if (messageUsage.cost) {
accumulatedUsage.cost.input += messageUsage.cost.input ?? 0;
accumulatedUsage.cost.output += messageUsage.cost.output ?? 0;
accumulatedUsage.cost.cacheRead += messageUsage.cost.cacheRead ?? 0;
accumulatedUsage.cost.cacheWrite += messageUsage.cost.cacheWrite ?? 0;
accumulatedUsage.cost.total += messageUsage.cost.total ?? 0;
}
}
// Accumulate tokens for progress display
progress.tokens += getUsageTokens(messageUsage);
}
// If pending termination, now we have tokens - terminate
@@ -531,7 +565,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
modelOverride,
error: exitCode !== 0 && stderr ? stderr : undefined,
aborted: wasAborted,
jsonlEvents,
usage: hasUsage ? accumulatedUsage : undefined,
artifactPaths,
extractedToolData: progress.extractedToolData,
outputMeta,
@@ -83,40 +83,6 @@ function addUsageTotals(target: Usage, usage: Partial<Usage>): 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<string, unknown>;
if (record.type !== "message_end") continue;
const message = record.message;
if (!message || typeof message !== "object") continue;
const msgRecord = message as Record<string, unknown>;
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<Usage>);
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,
@@ -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<string, unknown[]>;
/** Output metadata for Output tool integration */