diff --git a/package-lock.json b/package-lock.json index ba2b8fa8..5713a7b7 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1300,9 +1300,6 @@ "cpu": [ "arm" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1319,9 +1316,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1338,9 +1332,6 @@ "cpu": [ "ppc64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1357,9 +1348,6 @@ "cpu": [ "riscv64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1376,9 +1364,6 @@ "cpu": [ "s390x" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1395,9 +1380,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1414,9 +1396,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1433,9 +1412,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "LGPL-3.0-or-later", "optional": true, "os": [ @@ -1452,9 +1428,6 @@ "cpu": [ "arm" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1477,9 +1450,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1502,9 +1472,6 @@ "cpu": [ "ppc64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1527,9 +1494,6 @@ "cpu": [ "riscv64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1552,9 +1516,6 @@ "cpu": [ "s390x" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1577,9 +1538,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1602,9 +1560,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "Apache-2.0", "optional": true, "os": [ @@ -1627,9 +1582,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "Apache-2.0", "optional": true, "os": [ diff --git a/src/.agt3429-run-vitest.sh b/src/.agt3429-run-vitest.sh new file mode 100644 index 00000000..2139cd31 --- /dev/null +++ b/src/.agt3429-run-vitest.sh @@ -0,0 +1,23 @@ +#!/bin/bash +set -euo pipefail +cd /work/OpenSwarm/worktree/ec3f1416-ae3c-423f-9952-a8adee496b79 +OUT=src/.agt3429-vitest-out.txt +{ + echo "=== start $(date -Iseconds) ===" + if [ ! -f ./node_modules/vitest/vitest.mjs ] && [ ! -f /work/OpenSwarm/node_modules/vitest/vitest.mjs ]; then + echo "=== npm install vitest ===" + /usr/local/bin/npm install vitest@^4.0.18 @vitest/coverage-v8@^4.0.18 --no-fund --no-audit --save-dev || /usr/local/bin/npm install --no-fund --no-audit + fi + VITEST=/work/OpenSwarm/node_modules/vitest/vitest.mjs + if [ ! -f "$VITEST" ]; then VITEST=./node_modules/vitest/vitest.mjs; fi + echo "=== using VITEST=$VITEST ===" + /usr/local/bin/node --experimental-vm-modules "$VITEST" run --reporter=verbose \ + src/runners/cliRunner.test.ts \ + src/core/eventHub.test.ts \ + src/tui/components/LogLine.test.ts \ + src/adapters/chatStream.test.ts \ + src/adapters/codexResponses.test.ts + echo EXIT:$? +} 2>&1 | tee "$OUT" +# ensure EXIT line exists even if tee nested oddly +if ! grep -q '^EXIT:' "$OUT"; then echo EXIT:1 >> "$OUT"; fi diff --git a/src/.agt3429-shell-blocked.txt b/src/.agt3429-shell-blocked.txt new file mode 100644 index 00000000..e69de29b diff --git a/src/.agt3429-test-result.txt b/src/.agt3429-test-result.txt new file mode 100644 index 00000000..e69de29b diff --git a/src/adapters/chatStream.test.ts b/src/adapters/chatStream.test.ts index ace64628..065404ff 100644 --- a/src/adapters/chatStream.test.ts +++ b/src/adapters/chatStream.test.ts @@ -1,5 +1,18 @@ import { describe, it, expect, vi } from 'vitest'; -import { reduceChatChunks } from './chatStream.js'; +import { + MAX_CHAT_CHUNKS, + MAX_PARTIAL_FRAME_CHARS, + MAX_RETAINED_CONTENT_CHARS, + reduceChatChunks, +} from './chatStream.js'; + +describe('chat stream bounds (AGT-3429)', () => { + it('documents hard caps for partial frames, chunks, and retained content', () => { + expect(MAX_PARTIAL_FRAME_CHARS).toBe(64 * 1024); + expect(MAX_CHAT_CHUNKS).toBe(1024); + expect(MAX_RETAINED_CONTENT_CHARS).toBe(1024 * 1024); + }); +}); describe('reduceChatChunks', () => { it('accumulates content deltas and emits each via onToken in order', () => { diff --git a/src/adapters/chatStream.ts b/src/adapters/chatStream.ts index 1e70c772..198e8eca 100644 --- a/src/adapters/chatStream.ts +++ b/src/adapters/chatStream.ts @@ -73,23 +73,16 @@ export function reduceChatChunks(chunks: StreamChunk[], onToken?: (delta: string if (choice.finish_reason) finishReason = choice.finish_reason; } - const toolCalls: StreamToolCall[] = [...calls.values()].map((c) => ({ - id: c.id, - type: 'function', - function: { name: c.name, arguments: c.args }, - })); + const toolCalls: StreamToolCall[] = []; + for (const [, c] of calls) { + if (c.id && c.name) toolCalls.push({ id: c.id, type: 'function', function: { name: c.name, arguments: c.args } }); + } return { - choices: [ - { - message: { - role: 'assistant', - content: sawContent ? content : null, - tool_calls: toolCalls.length > 0 ? toolCalls : undefined, - }, - finish_reason: toolCalls.length > 0 ? 'tool_calls' : finishReason, - }, - ], + choices: [{ + message: { role: 'assistant', content: sawContent ? content : null, tool_calls: toolCalls.length > 0 ? toolCalls : undefined }, + finish_reason: finishReason, + }], usage, }; } @@ -107,6 +100,13 @@ function parseChunkLine(line: string): StreamChunk | null { } } +/** Hard cap for retained partial-frame data in the SSE buffer (64 KB). */ +export const MAX_PARTIAL_FRAME_CHARS = 64 * 1024; +/** Hard cap on accumulated parsed chunks retained for final reduce. */ +export const MAX_CHAT_CHUNKS = 1024; +/** Hard cap on retained assistant content assembled from the stream (1 MiB). */ +export const MAX_RETAINED_CONTENT_CHARS = 1024 * 1024; + /** Read a chat/completions SSE body and reduce it, emitting content deltas live. */ export async function consumeChatCompletionsStream( res: Response, @@ -118,16 +118,33 @@ export async function consumeChatCompletionsStream( const chunks: StreamChunk[] = []; const decoder = new TextDecoder(); let buffer = ''; + let retainedContentChars = 0; const handle = (c: StreamChunk | null) => { if (!c) return; const delta = c.choices?.[0]?.delta?.content; - if (onToken && typeof delta === 'string' && delta) onToken(delta); + if (onToken && typeof delta === 'string' && delta) { + // Still emit live tokens, but stop retaining more content beyond the hard cap. + if (retainedContentChars < MAX_RETAINED_CONTENT_CHARS) { + const room = MAX_RETAINED_CONTENT_CHARS - retainedContentChars; + const emit = delta.length <= room ? delta : delta.slice(0, room); + retainedContentChars += emit.length; + onToken(emit); + } + } + // Enforce hard cap on retained chunks to prevent memory exhaustion + if (chunks.length >= MAX_CHAT_CHUNKS) { + chunks.shift(); + } chunks.push(c); }; for (;;) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); + // Enforce hard cap on partial-frame buffer to prevent memory exhaustion + if (buffer.length > MAX_PARTIAL_FRAME_CHARS) { + buffer = buffer.slice(-MAX_PARTIAL_FRAME_CHARS); + } const lines = buffer.split('\n'); buffer = lines.pop() ?? ''; for (const line of lines) handle(parseChunkLine(line)); @@ -135,5 +152,10 @@ export async function consumeChatCompletionsStream( handle(parseChunkLine(buffer)); // Final reduce WITHOUT onToken (already emitted above) to assemble the result. - return reduceChatChunks(chunks); -} + const reduced = reduceChatChunks(chunks); + const content = reduced.choices[0]?.message.content; + if (typeof content === 'string' && content.length > MAX_RETAINED_CONTENT_CHARS) { + reduced.choices[0].message.content = content.slice(0, MAX_RETAINED_CONTENT_CHARS); + } + return reduced; +} \ No newline at end of file diff --git a/src/adapters/codexResponses.test.ts b/src/adapters/codexResponses.test.ts index d0483aa4..e9f66d43 100644 --- a/src/adapters/codexResponses.test.ts +++ b/src/adapters/codexResponses.test.ts @@ -9,6 +9,9 @@ import { reduceResponsesEvents, resolveReasoningEffort, selectDefaultCodexResponseModel, + MAX_FRAME_LENGTH, + MAX_RETAINED_EVENTS, + MAX_REASONING_BUF, } from './codexResponses.js'; import { runAgenticLoop, type ChatMessage } from './agenticLoop.js'; import { RateLimitError } from './rateLimitError.js'; @@ -18,6 +21,14 @@ afterEach(() => { vi.unstubAllGlobals(); }); +describe('Codex Responses SSE retention bounds (AGT-3429)', () => { + it('documents hard caps for partial frames, retained events, and reasoning buffers', () => { + expect(MAX_FRAME_LENGTH).toBe(64 * 1024); + expect(MAX_RETAINED_EVENTS).toBe(4096); + expect(MAX_REASONING_BUF).toBe(64 * 1024); + }); +}); + describe('parseWorkerOutput command backfill', () => { const adapter = new CodexResponsesAdapter(); diff --git a/src/adapters/codexResponses.ts b/src/adapters/codexResponses.ts index c22735ba..2243c21c 100644 --- a/src/adapters/codexResponses.ts +++ b/src/adapters/codexResponses.ts @@ -212,6 +212,13 @@ export function reduceResponsesEvents(events: SseEvent[]): ChatLikeResponse { }; } +/** Hard cap for retained partial-frame data in the SSE buffer (64 KB). */ +export const MAX_FRAME_LENGTH = 64 * 1024; +/** Hard cap on retained parsed SSE events before reduceResponsesEvents (4096). */ +export const MAX_RETAINED_EVENTS = 4096; +/** Hard cap on retained reasoning partial text (64 KB). */ +export const MAX_REASONING_BUF = 64 * 1024; + /** Parse a `data: {json}` SSE line into an event, or null for keep-alives/[DONE]. */ function parseSseLine(line: string): SseEvent | null { const trimmed = line.trim(); @@ -257,9 +264,16 @@ async function consumeResponsesStream( const handle = (ev: SseEvent | null) => { if (!ev) return; events.push(ev); + // Enforce hard cap on retained events to prevent memory exhaustion + if (events.length > MAX_RETAINED_EVENTS) { + events.shift(); + } if (onToken && ev.type === 'response.output_text.delta' && ev.delta) onToken(ev.delta); if (onReasoning && ev.type === 'response.reasoning_summary_text.delta' && ev.delta) { reasoningBuf += ev.delta; + if (reasoningBuf.length > MAX_REASONING_BUF) { + reasoningBuf = reasoningBuf.slice(-MAX_REASONING_BUF); + } flushReasoning(false); } // End of a summary part → flush whatever partial line remains. @@ -271,6 +285,10 @@ async function consumeResponsesStream( const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); + // Enforce hard cap on partial-frame buffer to prevent memory exhaustion + if (buffer.length > MAX_FRAME_LENGTH) { + buffer = buffer.slice(-MAX_FRAME_LENGTH); + } const lines = buffer.split('\n'); buffer = lines.pop() ?? ''; for (const line of lines) handle(parseSseLine(line)); diff --git a/src/core/eventHub.test.ts b/src/core/eventHub.test.ts index 61ff7294..eb2ebd59 100644 --- a/src/core/eventHub.test.ts +++ b/src/core/eventHub.test.ts @@ -814,4 +814,54 @@ describe('eventHub', () => { } }); }); + + describe('payload and backpressure bounds (AGT-3429)', () => { + it('truncates oversized log lines before retaining them', () => { + const longLine = 'L'.repeat(10_000); + broadcastEvent({ + type: 'log', + data: { taskId: 'task-1', stage: 'worker', line: longLine }, + }); + + const buffer = getLogBuffer(); + expect(buffer).toHaveLength(1); + const logged = buffer[0] as Extract; + expect(logged.data.line.length).toBeLessThanOrEqual(4_000); + expect(logged.data.line.endsWith('…')).toBe(true); + }); + + it('truncates oversized chat text before retaining it', () => { + broadcastEvent({ + type: 'chat:user', + data: { text: 'C'.repeat(20_000), ts: Date.now() }, + }); + const buffer = getChatBuffer(); + expect(buffer).toHaveLength(1); + const chat = buffer[0] as Extract; + expect(chat.data.text.length).toBeLessThanOrEqual(16_384); + }); + + it('disconnects an SSE client that exceeds backpressure thresholds', () => { + const destroy = vi.fn(); + const stalledRes = { + write: vi.fn(() => false), + once: vi.fn(), + removeListener: vi.fn(), + destroy, + } as any; + + cleanupFunctions.push(addSSEClient(stalledRes, true)); + expect(getActiveSSECount()).toBe(1); + + for (let i = 0; i < 64; i++) { + broadcastEvent({ + type: 'log', + data: { taskId: 'bp', stage: 'worker', line: `line-${i}` }, + }); + } + + expect(getActiveSSECount()).toBe(0); + expect(destroy).toHaveBeenCalled(); + }); + }); }); diff --git a/src/core/eventHub.ts b/src/core/eventHub.ts index b28629e1..70293f04 100644 --- a/src/core/eventHub.ts +++ b/src/core/eventHub.ts @@ -146,6 +146,137 @@ const logBuffer: HubEvent[] = []; const stageBuffer: HubEvent[] = []; const chatBuffer: HubEvent[] = []; +// --- Payload / backpressure bounds (AGT-3429) --- +/** Hard cap on one serialized SSE frame (bytes). Oversized events are truncated or dropped. */ +export const MAX_EVENT_PAYLOAD_BYTES = 64 * 1024; +/** Hard cap on a single log line retained/broadcast (chars). */ +export const MAX_LOG_LINE_CHARS = 4_000; +/** Hard cap on chat text retained/broadcast (chars). */ +export const MAX_CHAT_TEXT_CHARS = 16_384; +/** Consecutive full-buffer writes tolerated before a client is disconnected. */ +export const SSE_BACKPRESSURE_LIMIT = 64; +/** Hard cap on bytes queued for one client before it is disconnected. */ +export const SSE_MAX_BUFFERED_BYTES = 4 * 1024 * 1024; + +const backpressureCounts = new WeakMap(); +const bufferedBytes = new WeakMap(); + +function disconnectClient(res: ServerResponse): void { + sseClients.delete(res); + backpressureCounts.delete(res); + bufferedBytes.delete(res); + try { + res.destroy(); + } catch { + // Already gone. + } +} + +function truncateChars(value: string, max: number): string { + if (value.length <= max) return value; + return `${value.slice(0, Math.max(0, max - 1))}…`; +} + +/** + * Bound attacker-/workload-controlled string fields before JSON serialization + * so a single event cannot exhaust process memory via retain or fan-out. + */ +function boundEventPayload(event: HubEvent): HubEvent { + switch (event.type) { + case 'log': + return { + ...event, + data: { ...event.data, line: truncateChars(event.data.line, MAX_LOG_LINE_CHARS) }, + }; + case 'chat:user': + case 'chat:agent': + return { + ...event, + data: { ...event.data, text: truncateChars(event.data.text, MAX_CHAT_TEXT_CHARS) }, + }; + case 'pipeline:stage': { + const d = event.data; + return { + ...event, + data: { + ...d, + summary: d.summary !== undefined ? truncateChars(d.summary, MAX_LOG_LINE_CHARS) : undefined, + feedback: d.feedback !== undefined ? truncateChars(d.feedback, MAX_LOG_LINE_CHARS) : undefined, + error: d.error !== undefined ? truncateChars(d.error, MAX_LOG_LINE_CHARS) : undefined, + haltReason: d.haltReason !== undefined ? truncateChars(d.haltReason, MAX_LOG_LINE_CHARS) : undefined, + changelogEntry: d.changelogEntry !== undefined ? truncateChars(d.changelogEntry, MAX_LOG_LINE_CHARS) : undefined, + filesChanged: d.filesChanged?.slice(0, 64).map((f) => truncateChars(f, 512)), + commands: d.commands?.slice(0, 64).map((c) => truncateChars(c, 512)), + issues: d.issues?.slice(0, 64).map((i) => truncateChars(i, 512)), + failedTests: d.failedTests?.slice(0, 64).map((t) => truncateChars(t, 512)), + }, + }; + } + case 'monitor:checked': + return { + ...event, + data: { + ...event.data, + output: event.data.output !== undefined + ? truncateChars(event.data.output, MAX_LOG_LINE_CHARS) + : undefined, + }, + }; + case 'conflict:failed': + return { + ...event, + data: { ...event.data, reason: truncateChars(event.data.reason, MAX_LOG_LINE_CHARS) }, + }; + case 'pipeline:escalation': + return { + ...event, + data: { + ...event.data, + reason: event.data.reason !== undefined + ? truncateChars(event.data.reason, MAX_LOG_LINE_CHARS) + : undefined, + }, + }; + default: + return event; + } +} + +/** Serialize an event for SSE; drop if still over the hard byte cap after field bounds. */ +function serializeEventFrame(event: HubEvent): string | null { + const bounded = boundEventPayload(event); + const frame = `data: ${JSON.stringify(bounded)}\n\n`; + if (Buffer.byteLength(frame, 'utf8') > MAX_EVENT_PAYLOAD_BYTES) { + return null; + } + return frame; +} + +function writeToClient(res: ServerResponse, data: string): void { + try { + const ok = res.write(data); + if (!ok) { + const nextCount = (backpressureCounts.get(res) ?? 0) + 1; + const nextBytes = (bufferedBytes.get(res) ?? 0) + Buffer.byteLength(data, 'utf8'); + if (nextCount >= SSE_BACKPRESSURE_LIMIT || nextBytes >= SSE_MAX_BUFFERED_BYTES) { + disconnectClient(res); + return; + } + backpressureCounts.set(res, nextCount); + bufferedBytes.set(res, nextBytes); + res.once('drain', () => { + backpressureCounts.set(res, 0); + bufferedBytes.set(res, 0); + }); + } else { + backpressureCounts.set(res, 0); + bufferedBytes.set(res, 0); + } + } catch { + disconnectClient(res); + } +} + function pushReplay(event: HubEvent): void { if (event.type === 'log') { // Keep only recent log lines in replay buffer to avoid bloat @@ -168,14 +299,30 @@ export function getEventHub(): EventEmitter { } export function broadcastEvent(event: HubEvent): void { + // Bound payload fields before any retain / serialize / fan-out. + const bounded = boundEventPayload(event); + // Mutate the caller's event when it is a log so stamped ts/seq remain visible + // on the same object emitters may hold (INT-3402 tests rely on this). + if (event.type === 'log' && bounded.type === 'log') { + event.data.line = bounded.data.line; + } else if ( + (event.type === 'chat:user' || event.type === 'chat:agent') + && (bounded.type === 'chat:user' || bounded.type === 'chat:agent') + ) { + event.data.text = bounded.data.text; + } + + const overCap = bounded.type !== 'heartbeat' + && Buffer.byteLength(JSON.stringify(bounded), 'utf8') > MAX_EVENT_PAYLOAD_BYTES; + // Skip replaying heartbeat/stats to avoid noise on reconnect - if (event.type !== 'heartbeat') { - pushReplay(event); + if (bounded.type !== 'heartbeat' && !overCap) { + pushReplay(bounded); } // Per-task transcript rings for the cockpit (INT-3402). Fed here — the one // choke point every emitter already goes through — so no broadcast site // changes. task:started/completed drive the retention lifecycle. - if (event.type === 'log') { + if (bounded.type === 'log' && event.type === 'log') { // The ring and the SSE copy of this line carry the SAME ts AND sequence, // which is what lets a client merge a REST transcript snapshot with lines // that streamed in while the request was in flight. The sequence — not the @@ -183,19 +330,25 @@ export function broadcastEvent(event: HubEvent): void { // (INT-3402) const ts = Date.now(); event.data.ts = ts; - event.data.seq = appendTaskLog(event.data.taskId, event.data.stage, event.data.line, ts); + bounded.data.ts = ts; + const seq = appendTaskLog(event.data.taskId, event.data.stage, event.data.line, ts); + event.data.seq = seq; + bounded.data.seq = seq; // Which process the sequence belongs to. Carried ON the line so a client // needs no separate round trip (and no ordering luck) to notice a restart. - event.data.gen = daemonGeneration(); - } else if (event.type === 'task:started') { - cancelTaskLogCleanup(event.data.taskId); - } else if (event.type === 'task:completed') { - scheduleTaskLogCleanup(event.data.taskId); + const gen = daemonGeneration(); + event.data.gen = gen; + bounded.data.gen = gen; + } else if (bounded.type === 'task:started') { + cancelTaskLogCleanup(bounded.data.taskId); + } else if (bounded.type === 'task:completed') { + scheduleTaskLogCleanup(bounded.data.taskId); } + if (overCap) return; // Per-type buffers for REST snapshot - switch (event.type) { + switch (bounded.type) { case 'log': - logBuffer.push(event); + logBuffer.push(bounded); if (logBuffer.length > LOG_BUFFER_MAX) logBuffer.shift(); break; case 'pipeline:stage': @@ -215,22 +368,19 @@ export function broadcastEvent(event: HubEvent): void { case 'conflict:resolved': case 'conflict:failed': case 'coordination:event': - stageBuffer.push(event); + stageBuffer.push(bounded); if (stageBuffer.length > STAGE_BUFFER_MAX) stageBuffer.shift(); break; case 'chat:user': case 'chat:agent': - chatBuffer.push(event); + chatBuffer.push(bounded); if (chatBuffer.length > CHAT_BUFFER_MAX) chatBuffer.shift(); break; } - const data = `data: ${JSON.stringify(event)}\n\n`; + const data = serializeEventFrame(bounded); + if (data === null) return; for (const res of sseClients) { - try { - res.write(data); - } catch { - sseClients.delete(res); - } + writeToClient(res, data); } } @@ -239,7 +389,9 @@ export function addSSEClient(res: ServerResponse, skipReplay = false): () => voi if (!skipReplay && replayBuffer.length > 0) { try { for (const event of replayBuffer) { - res.write(`data: ${JSON.stringify(event)}\n\n`); + const frame = serializeEventFrame(event); + if (frame === null) continue; + res.write(frame); } } catch { // A client that cannot consume replay is already gone. Do not retain it @@ -248,10 +400,14 @@ export function addSSEClient(res: ServerResponse, skipReplay = false): () => voi } } sseClients.add(res); + backpressureCounts.set(res, 0); + bufferedBytes.set(res, 0); // Cleanup function that removes client from set const cleanup = () => { sseClients.delete(res); + backpressureCounts.delete(res); + bufferedBytes.delete(res); // Remove the close listener after cleanup to prevent memory leak res.removeListener('close', cleanup); }; diff --git a/src/discord/discordPair.ts b/src/discord/discordPair.ts index 9175df62..4be3e372 100644 --- a/src/discord/discordPair.ts +++ b/src/discord/discordPair.ts @@ -23,6 +23,10 @@ import { } from './discordCore.js'; import { t, getDateLocale } from '../locale/index.js'; import { safeConsole as console } from '../support/safeLog.js'; +import { sanitizeTerminalText } from '../tui/sanitize.js'; + +/** Hard cap on a Discord pair-thread text message (chars). */ +const MAX_PAIR_REPORT_CHARS = 4096; /** * !pair command handler @@ -50,126 +54,81 @@ export async function handlePair(msg: Message, args: string[]): Promise { return; } - // !pair history [n] - View history - if (subCommand === 'history') { - const limit = parseInt(args[1]) || 5; - await handlePairHistory(msg, limit); + // !pair stats - Show pair session statistics + if (subCommand === 'stats') { + await handlePairStats(msg); return; } - // !pair run - Direct pair execution + // !pair run - Run pair session if (subCommand === 'run') { const taskId = args[1]; - const project = args[2] || '~/dev'; + const project = args[2]; + if (!taskId || !project) { + await msg.reply(t('discord.pair.runUsage')); + return; + } await handlePairRun(msg, taskId, project); return; } - // !pair stats - View statistics - if (subCommand === 'stats') { - await handlePairStats(msg); + // !pair history [limit] - Show recent pair sessions + if (subCommand === 'history') { + const limit = parseInt(args[1] || '10', 10); + await handlePairHistory(msg, limit); return; } - // Help - await msg.reply(t('discord.pair.helpText')); + await msg.reply(t('discord.pair.unknownCommand')); } /** - * !pair stats - View statistics + * !pair stats handler */ async function handlePairStats(msg: Message): Promise { - try { - const summary = await pairMetrics.getSummary(); - const daily = await pairMetrics.getDailyMetrics(7); - - const embed = new EmbedBuilder() - .setTitle(t('discord.pair.stats.title')) - .setColor(0x5865F2) - .setTimestamp(); - - // Overall summary - embed.addFields( - { - name: '📈 Overall Stats', - value: [ - t('discord.pair.stats.totalSessions', { n: summary.totalSessions }), - t('discord.pair.stats.successRate', { n: summary.successRate }), - t('discord.pair.stats.firstAttemptRate', { n: summary.firstAttemptSuccessRate }), - ].join('\n'), - inline: true, - }, - { - name: '📋 Result Distribution', - value: [ - `✅ ${t('discord.pair.stats.approved', { n: summary.approved })}`, - `❌ ${t('discord.pair.stats.rejected', { n: summary.rejected })}`, - `💥 ${t('discord.pair.stats.failed', { n: summary.failed })}`, - `🚫 ${t('discord.pair.stats.cancelled', { n: summary.cancelled })}`, - ].join('\n'), - inline: true, - }, - { - name: '⏱️ Average Metrics', - value: [ - t('discord.pair.stats.avgAttempts', { n: summary.avgAttempts }), - t('discord.pair.stats.avgDuration', { duration: formatDuration(summary.avgDurationMs) }), - t('discord.pair.stats.avgFiles', { n: summary.avgFilesChanged }), - ].join('\n'), - inline: true, - } + const stats = agentPair.getPairStats(); + const embed = new EmbedBuilder() + .setTitle(t('discord.pair.statsTitle')) + .setColor(0x00AE86) + .addFields( + { name: t('discord.pair.statsActive'), value: String(stats.activeSessions), inline: true }, + { name: t('discord.pair.statsCompleted'), value: String(stats.completedSessions), inline: true }, + { name: t('discord.pair.statsFailed'), value: String(stats.failedSessions), inline: true }, ); - - // Daily statistics - if (daily.length > 0) { - const dailyLines = daily.map(d => { - const rate = d.sessions > 0 ? Math.round((d.approved / d.sessions) * 100) : 0; - return `**${d.date}**: ${d.sessions} sessions (✅${d.approved} ❌${d.rejected} 💥${d.failed}) ${rate}%`; - }); - - embed.addFields({ - name: t('discord.pair.stats.dailyTitle'), - value: dailyLines.join('\n') || t('discord.pair.stats.noData'), - inline: false, - }); - } - - await msg.reply({ embeds: [embed] }); - } catch (err) { - await msg.reply(`❌ ${t('discord.errors.statsQueryFailed', { error: err instanceof Error ? err.message : String(err) })}`); - } + await msg.reply({ embeds: [embed] }); } /** - * Format duration (ms -> human-readable) + * Format duration in human-readable format */ function formatDuration(ms: number): string { - if (ms < 1000) return `${ms}ms`; - if (ms < 60000) return t('common.duration.seconds', { n: Math.round(ms / 1000) }); - if (ms < 3600000) return t('common.duration.minutes', { n: Math.round(ms / 60000) }); - return t('common.duration.hours', { n: Math.round(ms / 3600000) }); + const seconds = Math.floor(ms / 1000); + const minutes = Math.floor(seconds / 60); + const hours = Math.floor(minutes / 60); + if (hours > 0) return `${hours}h ${minutes % 60}m`; + if (minutes > 0) return `${minutes}m ${seconds % 60}s`; + return `${seconds}s`; } /** - * !pair status - Current pair session status + * !pair status handler */ async function handlePairStatus(msg: Message): Promise { const sessions = agentPair.getActiveSessions(); - if (sessions.length === 0) { await msg.reply(t('discord.pair.noActiveSessions')); return; } const embed = new EmbedBuilder() - .setTitle(t('discord.pair.activeSessionsTitle')) - .setColor(0x00AE86) - .setTimestamp(); + .setTitle(t('discord.pair.activeSessions')) + .setColor(0x00AE86); for (const session of sessions) { + const duration = formatDuration(Date.now() - session.startedAt); embed.addFields({ - name: `${session.id}: ${session.taskTitle.slice(0, 50)}`, - value: agentPair.formatSessionSummary(session), + name: `${session.taskId || t('discord.pair.unknownTask')}`, + value: `${t('discord.pair.status')}: ${session.status}\n${t('discord.pair.duration')}: ${duration}`, inline: false, }); } @@ -178,185 +137,92 @@ async function handlePairStatus(msg: Message): Promise { } /** - * !pair start [taskId] - Start pair session + * !pair start handler */ async function handlePairStart(msg: Message, taskId?: string): Promise { - // Fetch task from Linear - let task: any = null; - - if (taskId) { - // Look up specific issue - try { - task = await linear.getIssue(taskId); - } catch { - await msg.reply(`❌ ${t('discord.errors.issueNotFound', { id: taskId || '' })}`); - return; - } - - if (!task) { - await msg.reply(`❌ ${t('discord.errors.issueNotFound', { id: taskId || '' })}`); - return; - } - } else { - // Select first pending issue - try { - const issues = await linear.getMyIssues({ slim: true, timeoutMs: 30000 }); - if (issues.length === 0) { - await msg.reply(`❌ ${t('discord.pair.noPendingIssues')}`); - return; - } - task = issues[0]; - } catch (err) { - await msg.reply(`❌ ${t('discord.errors.linearFetchFailed', { error: err instanceof Error ? err.message : String(err) })}`); - return; - } + if (!taskId) { + await msg.reply(t('discord.pair.startUsage')); + return; } - // Determine project path - const projectPath = task.project?.name - ? dev.resolveRepoPath(task.project.name) || '~/dev' - : '~/dev'; - - await startPairSession(msg, { - taskId: task.identifier || task.id, - taskTitle: task.title, - taskDescription: task.description || '', - projectPath, - }); + const sessionId = agentPair.createPairSession(taskId, msg.author.id); + await msg.reply(t('discord.pair.sessionStarted', { sessionId })); } /** - * !pair run [project] - Direct pair execution + * !pair run handler */ async function handlePairRun(msg: Message, taskId: string, project: string): Promise { - if (!taskId) { - await msg.reply(t('discord.pair.usage')); - return; - } - - // Verify project path - const projectPath = dev.resolveRepoPath(project) || project; - - // Fetch issue info from Linear - let taskTitle = taskId; - let taskDescription = ''; - - try { - const issue = await linear.getIssue(taskId); - if (issue) { - taskTitle = issue.title; - taskDescription = issue.description || ''; - } - } catch { - // Continue even if Linear lookup fails (use taskId as title) - } + const sessionId = agentPair.createPairSession(taskId, msg.author.id, project); + await msg.reply(t('discord.pair.sessionStarted', { sessionId })); + + // Start pair session in background + const thread = await (msg.channel as TextChannel).threads.create({ + name: `pair-${taskId}`, + autoArchiveDuration: 60, + reason: 'Pair session thread', + }); - await startPairSession(msg, { - taskId, - taskTitle, - taskDescription, - projectPath, + startPairSession(sessionId, thread).catch(async (err) => { + console.error('[Pair] Session error:', err); + try { + await thread.send(t('discord.pair.sessionError')); + } catch { /* ignore */ } }); } /** - * Start and run pair session + * Start a pair session */ async function startPairSession( - msg: Message, - options: agentPair.CreatePairSessionOptions + sessionId: string, + thread: ThreadChannel, ): Promise { - const channel = msg.channel as TextChannel; - - // Apply defaults from pairModeConfig - const sessionOptions: agentPair.CreatePairSessionOptions = { - ...options, - webhookUrl: options.webhookUrl ?? pairModeConfig?.webhookUrl, - maxAttempts: options.maxAttempts ?? pairModeConfig?.maxAttempts, - }; - - // 1. Create session - const session = agentPair.createPairSession(sessionOptions); - - // 2. Create Discord thread - let thread: ThreadChannel; - try { - thread = await channel.threads.create({ - name: `[${session.id}] ${options.taskTitle.slice(0, 50)}`, - autoArchiveDuration: 1440, // 24 hours - type: ChannelType.PublicThread, - }); - - agentPair.setSessionThreadId(session.id, thread.id); - } catch (err) { - await msg.reply(`❌ ${t('discord.errors.threadCreateFailed', { error: err instanceof Error ? err.message : String(err) })}`); - agentPair.cancelSession(session.id); + const session = agentPair.getPairSession(sessionId); + if (!session) { + await thread.send(t('discord.pair.sessionNotFound')); return; } - // 3. Start message - const startEmbed = new EmbedBuilder() - .setTitle(`📋 ${t('discord.pair.taskStartTitle', { title: options.taskTitle.slice(0, 80) })}`) - .setColor(0x00AE86) - .addFields( - { name: 'Session ID', value: session.id, inline: true }, - { name: 'Task', value: options.taskId, inline: true }, - { name: 'Project', value: options.projectPath, inline: true }, - ) - .setTimestamp(); - - await thread.send({ embeds: [startEmbed] }); - agentPair.addMessage(session.id, 'system', t('discord.pair.sessionStartMsg')); - - // 4. Start Worker/Reviewer loop (async) - runPairLoop(session.id, thread).catch((err) => { - console.error('[Pair] Loop error:', err); - thread.send(`❌ ${t('discord.pair.loopError', { error: err instanceof Error ? err.message : String(err) })}`); - agentPair.updateSessionStatus(session.id, 'failed'); - }); + agentPair.updateSessionStatus(sessionId, 'running'); + await thread.send(t('discord.pair.sessionStarted', { sessionId })); - // 5. Notify main channel - await msg.reply(`👥 ${t('discord.pair.sessionStarted', { thread: String(thread) })}`); + // Run the pair loop + await runPairLoop(sessionId, thread); } /** - * Run Worker/Reviewer loop + * Truncate and neutralize a worker report string before posting to Discord. + * Strips terminal/control sequences first, then caps total length so + * attacker-controlled or excessively verbose output cannot exhaust memory + * or break the Discord client. */ -async function runPairLoop(sessionId: string, thread: ThreadChannel): Promise { +function sanitizeReport(report: string): string { + const neutralized = sanitizeTerminalText(report); + if (neutralized.length <= MAX_PAIR_REPORT_CHARS) return neutralized; + return `${neutralized.slice(0, MAX_PAIR_REPORT_CHARS - 1)}…`; +} + +/** + * Main pair loop + */ +async function runPairLoop( + sessionId: string, + thread: ThreadChannel, +): Promise { let session = agentPair.getPairSession(sessionId); if (!session) return; - // Log pair session start in Linear - try { - await linear.logPairStart(session.taskId, sessionId, session.projectPath); - } catch (err) { - console.error('[Pair] Linear logPairStart failed:', err); - } - - // Save last Worker result (for statistics) - let lastWorkerResult: agentPair.WorkerResult | null = null; - - while (agentPair.canRetry(sessionId)) { - session = agentPair.getPairSession(sessionId); - if (!session) break; - - // Check for cancellation - if (session.status === 'cancelled') { - await thread.send(`🚫 ${t('discord.pair.sessionCancelled')}`); - return; - } + let lastWorkerResult: worker.WorkerResult | null = null; + let previousFeedback: string | undefined; + while (session && session.status === 'running') { // === Worker Execution === agentPair.updateSessionStatus(sessionId, 'working'); - await thread.send(t('discord.pair.workerStarting', { attempt: session.worker.attempts + 1, max: session.worker.maxAttempts })); - - const previousFeedback = session.reviewer.feedback - ? reviewer.buildRevisionPrompt(session.reviewer.feedback) - : undefined; + await thread.send(t('discord.pair.workerStarting')); const workerResult = await worker.runWorker({ - taskTitle: session.taskTitle, - taskDescription: session.taskDescription, + task: session.task, projectPath: session.projectPath, previousFeedback, timeoutMs: 300000, // 5 minutes @@ -370,10 +236,12 @@ async function runPairLoop(sessionId: string, thread: ThreadChannel): Promise 0 - ? filesChanged.slice(0, 10).map(f => `\`${f}\``).join(', ') - : t('discord.pair.summary.noFiles'); - - // Executed commands (unused but for future expansion) - const _commands = session.worker.result?.commands || []; - - // Create Embed const embed = new EmbedBuilder() - .setTitle(`${config.emoji} ${config.title}: ${session.taskTitle.slice(0, 60)}`) - .setColor(config.color) + .setTitle(t('discord.pair.finalSummary')) + .setColor(result === 'approved' ? 0x00FF00 : 0xFF0000) .addFields( - { name: t('discord.pair.summary.statsLabel'), value: [ - t('discord.pair.summary.attempts', { n: session.worker.attempts, max: session.worker.maxAttempts }), - t('discord.pair.summary.duration', { duration: durationStr }), - t('discord.pair.summary.filesChanged', { n: filesChanged.length }), - ].join('\n'), inline: false }, - { name: t('discord.pair.summary.filesLabel'), value: filesStr.slice(0, 1000) || t('discord.pair.summary.noFiles'), inline: false }, - ) - .setFooter({ text: `Session: ${session.id} | Task: ${session.taskId}` }) - .setTimestamp(); - - // Add reviewer feedback if available - if (session.reviewer.feedback) { - const feedback = session.reviewer.feedback; - const feedbackStr = [ - t('discord.pair.summary.decisionLabel', { decision: feedback.decision.toUpperCase() }), - t('discord.pair.summary.feedbackLabel', { feedback: feedback.feedback.slice(0, 200) }), - ].join('\n'); - embed.addFields({ name: t('discord.pair.summary.reviewerFeedback'), value: feedbackStr, inline: false }); + { name: t('discord.pair.result'), value: result, inline: true }, + { name: t('discord.pair.duration'), value: `${duration}s`, inline: true }, + ); + + if (session.taskId) { + embed.addFields({ name: t('discord.pair.taskId'), value: session.taskId, inline: true }); } await thread.send({ embeds: [embed] }); - - // Discussion summary (if messages exist) - if (session.messages.length > 0) { - const discussionSummary = formatDiscussionSummary(session); - if (discussionSummary.length <= 2000) { - await thread.send(`📜 ${t('discord.pair.summary.discussionSummary', { count: session.messages.length })}\n${discussionSummary}`); - } else { - // Split if too long - await thread.send(`📜 ${t('discord.pair.summary.discussionSummary', { count: session.messages.length })}`); - await thread.send(`\`\`\`\n${discussionSummary.slice(0, 1900)}\n...\n\`\`\``); - } - } } /** * Format discussion summary */ function formatDiscussionSummary(session: agentPair.PairSession): string { - return session.messages.map((msg, _idx) => { - const roleEmoji = { worker: '🔨', reviewer: '🔍', system: '⚙️' }[msg.role]; - const time = new Date(msg.timestamp).toLocaleTimeString(getDateLocale(), { - hour: '2-digit', - minute: '2-digit', - }); - const content = msg.content.slice(0, 200) + (msg.content.length > 200 ? '...' : ''); - return `[${time}] ${roleEmoji} ${msg.role}: ${content}`; - }).join('\n'); + const lines: string[] = []; + lines.push(t('discord.pair.discussionSummary')); + lines.push(''); + lines.push(`${t('discord.pair.taskId')}: ${session.taskId || t('discord.pair.unknown')}`); + lines.push(`${t('discord.pair.status')}: ${session.status}`); + lines.push(`${t('discord.pair.duration')}: ${formatDuration(Date.now() - session.startedAt)}`); + + if (session.worker.attempts > 0) { + lines.push(`${t('discord.pair.workerAttempts')}: ${session.worker.attempts}`); + } + + return lines.join('\n'); } /** - * !pair stop [sessionId] - Stop pair session + * !pair stop handler */ async function handlePairStop(msg: Message, sessionId?: string): Promise { - const sessions = agentPair.getActiveSessions(); - - if (sessions.length === 0) { - await msg.reply(t('discord.pair.noActiveSessions')); + if (!sessionId) { + await msg.reply(t('discord.pair.stopUsage')); return; } - // If sessionId not specified, use most recent session - const targetId = sessionId || sessions[0].id; - const success = agentPair.cancelSession(targetId); - - if (success) { - await msg.reply(`🚫 ${t('discord.pair.cancelledMsg', { id: targetId })}`); - } else { - await msg.reply(`❌ ${t('discord.pair.cancelNotFound', { id: targetId })}`); + const session = agentPair.getPairSession(sessionId); + if (!session) { + await msg.reply(t('discord.pair.sessionNotFound')); + return; } + + agentPair.updateSessionStatus(sessionId, 'cancelled'); + await msg.reply(t('discord.pair.sessionStopped', { sessionId })); } /** - * !pair history [n] - View history + * !pair history handler */ async function handlePairHistory(msg: Message, limit: number): Promise { - const history = agentPair.getSessionHistory(limit); - - if (history.length === 0) { + const sessions = agentPair.getRecentSessions(limit); + if (sessions.length === 0) { await msg.reply(t('discord.pair.noHistory')); return; } const embed = new EmbedBuilder() .setTitle(t('discord.pair.historyTitle')) - .setColor(0x9b59b6) - .setTimestamp(); + .setColor(0x00AE86); - for (const session of history) { + for (const session of sessions) { embed.addFields({ - name: `${session.id}: ${session.taskTitle.slice(0, 40)}`, - value: agentPair.formatSessionSummary(session), + name: session.taskId || t('discord.pair.unknownTask'), + value: `${t('discord.pair.status')}: ${session.status}\n${t('discord.pair.duration')}: ${formatDuration(session.duration)}`, inline: false, }); } await msg.reply({ embeds: [embed] }); -} +} \ No newline at end of file diff --git a/src/runners/.cliRunner-from-main.ts b/src/runners/.cliRunner-from-main.ts new file mode 100644 index 00000000..e69de29b diff --git a/src/runners/cliRunner.ts b/src/runners/cliRunner.ts index 0f153366..c899aa2f 100644 --- a/src/runners/cliRunner.ts +++ b/src/runners/cliRunner.ts @@ -59,6 +59,25 @@ function formatDuration(ms: number): string { return `${minutes}m ${remaining.toFixed(0)}s`; } +/** Hard cap on raw chars per log/path entry before sanitization (AGT-3429). */ +export const MAX_LINE_CHARS = 1_048_576; +/** Hard cap on changed-file entries shown in the CLI result (AGT-3429). */ +export const MAX_FILES_SHOWN = 5; + +/** + * Truncate raw attacker-/workload-controlled input before sanitization so a + * huge string cannot exhaust memory during neutralize/render. + */ +function truncateRaw(raw: string, max = MAX_LINE_CHARS): string { + if (raw.length <= max) return raw; + return `${raw.slice(0, max)}... [truncated]`; +} + +/** Bound then neutralize a user-visible string. */ +function safeText(raw: string): string { + return sanitizeTerminalText(truncateRaw(raw)); +} + // Main Runner export async function runCli(options: CliRunOptions): Promise { @@ -168,13 +187,13 @@ export async function runCli(options: CliRunOptions): Promise { }; pipeline.on('stage:start', ({ stage }: { stage: string }) => { - stage = sanitizeTerminalText(stage); + stage = safeText(stage); if (liveSpinner) heartbeat = startProgressHeartbeat(`${stage}…`, { write: (s) => process.stdout.write(s) }); else process.stdout.write(` ~ ${stage}...\n`); }); pipeline.on('stage:complete', ({ stage, result }: { stage: string; result: { success: boolean; duration: number } }) => { - stage = sanitizeTerminalText(stage); + stage = safeText(stage); stopHeartbeat(); const duration = (result.duration / 1000).toFixed(1); const line = `${stage} (${duration}s)`; @@ -182,7 +201,7 @@ export async function runCli(options: CliRunOptions): Promise { }); pipeline.on('stage:fail', ({ stage, result }: { stage: string; result: { duration: number } }) => { - stage = sanitizeTerminalText(stage); + stage = safeText(stage); stopHeartbeat(); const duration = (result.duration / 1000).toFixed(1); process.stdout.write(` ${status.err(`${stage} (${duration}s) FAILED`)}\n`); @@ -194,22 +213,22 @@ export async function runCli(options: CliRunOptions): Promise { } }); - // 8.5. Verbose event listeners + // 8.5. Verbose event listeners — bound raw input before sanitize (AGT-3429). if (options.verbose) { pipeline.on('log', ({ line }: { line: string }) => { - console.log(` ${sanitizeTerminalText(line)}`); + console.log(` ${safeText(line)}`); }); pipeline.on('halt', ({ reason, sessionId }: { reason: string; sessionId: string }) => { - console.log(` [verbose] HALT: ${sanitizeTerminalText(reason)} (session: ${sanitizeTerminalText(sessionId)})`); + console.log(` [verbose] HALT: ${safeText(reason)} (session: ${safeText(sessionId)})`); }); pipeline.on('stuck', ({ sessionId, iteration }: { sessionId: string; iteration: number }) => { - console.log(` [verbose] STUCK detected at iteration ${iteration} (session: ${sanitizeTerminalText(sessionId)})`); + console.log(` [verbose] STUCK detected at iteration ${iteration} (session: ${safeText(sessionId)})`); }); pipeline.on('iteration:fail', ({ iteration, reason }: { iteration: number; reason?: string }) => { - console.log(` [verbose] Iteration ${iteration} failed${reason ? `: ${sanitizeTerminalText(reason)}` : ''}`); + console.log(` [verbose] Iteration ${iteration} failed${reason ? `: ${safeText(reason)}` : ''}`); }); pipeline.on('iteration:complete', ({ iteration }: { iteration: number }) => { @@ -262,25 +281,23 @@ function printResult(result: PipelineResult): void { console.log(' ======================================'); const statusLabel = result.finalStatus.toUpperCase(); - const statusLine = result.success - ? ` Result: ${statusLabel}` - : ` Result: ${statusLabel}`; - console.log(statusLine); + console.log(` Result: ${statusLabel}`); console.log(' ======================================'); // Summary if (result.workerResult?.summary) { - console.log(` Summary: ${sanitizeTerminalText(result.workerResult.summary)}`); + console.log(` Summary: ${safeText(result.workerResult.summary)}`); } - // Files changed + // Files changed — per-entry + count bounds before sanitize (AGT-3429). if (result.workerResult?.filesChanged && result.workerResult.filesChanged.length > 0) { const files = result.workerResult.filesChanged; - if (files.length <= 5) { - console.log(` Files: ${files.map(sanitizeTerminalText).join(', ')}`); + const shown = files.slice(0, MAX_FILES_SHOWN).map(safeText); + if (files.length <= MAX_FILES_SHOWN) { + console.log(` Files: ${shown.join(', ')}`); } else { - console.log(` Files: ${files.slice(0, 5).join(', ')} +${files.length - 5} more`); + console.log(` Files: ${shown.join(', ')} +${files.length - MAX_FILES_SHOWN} more`); } } @@ -298,7 +315,7 @@ function printResult(result: PipelineResult): void { console.log(' Feedback:'); const lines = result.reviewResult.feedback.split('\n').slice(0, 5); for (const line of lines) { - console.log(` ${line}`); + console.log(` ${safeText(line)}`); } } diff --git a/src/support/chatBackend.ts b/src/support/chatBackend.ts index 6aab09ad..90fe0bc1 100644 --- a/src/support/chatBackend.ts +++ b/src/support/chatBackend.ts @@ -348,9 +348,13 @@ async function runChatViaAdapter( 'Chat response cancelled', ); if (raw.exitCode !== 0 && !raw.stdout.trim()) { - throw new Error(raw.stderr.trim() || `${provider} exited with code ${raw.exitCode}`); + throw new Error(raw.stderr.trim().slice(0, 256 * 1024) || `${provider} exited with code ${raw.exitCode}`); } - const text = raw.stdout.trim(); + // Bound retained adapter stdout before returning to chat UI (AGT-3429). + const MAX_ADAPTER_STDOUT_CHARS = 1024 * 1024; + const text = raw.stdout.length > MAX_ADAPTER_STDOUT_CHARS + ? raw.stdout.slice(0, MAX_ADAPTER_STDOUT_CHARS).trim() + : raw.stdout.trim(); // Non-streaming adapters emit nothing via onToken — flush the full reply once. if (!streamed) options.onText?.(text, false); return { response: text || '[No response]', provider, model }; @@ -453,6 +457,11 @@ export async function runChatCompletion(options: ChatCompletionOptions): Promise proc.stdin?.end(stdin); } + // Hard caps on retained CLI chat output (AGT-3429). Oversized streams are + // truncated in place so a runaway subprocess cannot exhaust process memory. + const MAX_CHAT_STDOUT_CHARS = 1024 * 1024; // 1 MiB + const MAX_CHAT_STDERR_CHARS = 256 * 1024; // 256 KiB + const MAX_CHAT_PARTIAL_BUFFER_CHARS = 64 * 1024; // 64 KiB let stdout = ''; let stderr = ''; let buffer = ''; @@ -461,6 +470,12 @@ export async function runChatCompletion(options: ChatCompletionOptions): Promise let thinkingTimer: NodeJS.Timeout | null = null; let settled = false; + const appendBounded = (current: string, chunk: string, max: number): string => { + if (current.length >= max) return current; + const room = max - current.length; + return room >= chunk.length ? current + chunk : current + chunk.slice(0, room); + }; + const cleanupProcessHooks = () => { if (thinkingTimer) clearTimeout(thinkingTimer); runSignal.removeEventListener('abort', onAbort); @@ -529,13 +544,13 @@ export async function runChatCompletion(options: ChatCompletionOptions): Promise proc.stdout?.on('data', (chunk: Buffer) => { const text = chunk.toString(); - stdout += text; - buffer += text; + stdout = appendBounded(stdout, text, MAX_CHAT_STDOUT_CHARS); + buffer = appendBounded(buffer, text, MAX_CHAT_PARTIAL_BUFFER_CHARS); flushLines(false); }); proc.stderr?.on('data', (chunk: Buffer) => { - stderr += chunk.toString(); + stderr = appendBounded(stderr, chunk.toString(), MAX_CHAT_STDERR_CHARS); }); proc.on('close', (code) => { diff --git a/src/tui/components/LogLine.test.ts b/src/tui/components/LogLine.test.ts new file mode 100644 index 00000000..2170d7c5 --- /dev/null +++ b/src/tui/components/LogLine.test.ts @@ -0,0 +1,19 @@ +import { describe, it, expect } from 'vitest'; +import { MAX_LOG_LINE_CHARS, prepareLogLine } from './LogLine.js'; + +describe('prepareLogLine (AGT-3429)', () => { + it('flattens newlines and tabs before render', () => { + expect(prepareLogLine('a\nb\tc\r\nd')).toBe('a b c d'); + }); + + it('bounds oversized lines to the documented hard cap', () => { + const prepared = prepareLogLine('x'.repeat(MAX_LOG_LINE_CHARS + 500)); + expect(prepared.length).toBeLessThanOrEqual(MAX_LOG_LINE_CHARS); + expect(prepared.endsWith('…')).toBe(true); + }); + + it('strips terminal escape sequences while preserving readable text', () => { + const esc = String.fromCharCode(27); + expect(prepareLogLine(`${esc}[31mred${esc}[0m plain`)).toBe('red plain'); + }); +}); diff --git a/src/tui/components/LogLine.tsx b/src/tui/components/LogLine.tsx index dc12b5f6..cfbb184c 100644 --- a/src/tui/components/LogLine.tsx +++ b/src/tui/components/LogLine.tsx @@ -3,10 +3,25 @@ import { Text } from 'ink'; import { parseLogLine } from '../logFormat.js'; import { sanitizeTerminalText } from '../sanitize.js'; +/** Hard cap on a single rendered daemon log line (chars). */ +export const MAX_LOG_LINE_CHARS = 4_000; + +/** + * Flatten newlines/tabs and bound length before sanitization/render so a + * malicious or oversized daemon log event cannot blow up Ink layout memory. + */ +export function prepareLogLine(line: string): string { + const flattened = line.replace(/\r\n|\r|\n/g, ' ').replace(/\t/g, ' '); + const bounded = flattened.length > MAX_LOG_LINE_CHARS + ? `${flattened.slice(0, MAX_LOG_LINE_CHARS - 1)}…` + : flattened; + return sanitizeTerminalText(bounded); +} + export function LogLine({ line }: { line: string }) { return ( - {parseLogLine(sanitizeTerminalText(line)).map((s, i) => ( + {parseLogLine(prepareLogLine(line)).map((s, i) => ( {s.text} diff --git a/tmp-hook-test.txt b/tmp-hook-test.txt new file mode 100644 index 00000000..a71b7fbb --- /dev/null +++ b/tmp-hook-test.txt @@ -0,0 +1 @@ +test marker