diff --git a/src/adapters/chatStream.test.ts b/src/adapters/chatStream.test.ts index c1913e53..c2bcb1d6 100644 --- a/src/adapters/chatStream.test.ts +++ b/src/adapters/chatStream.test.ts @@ -1,5 +1,121 @@ import { describe, it, expect, vi } from 'vitest'; -import { reduceChatChunks } from './chatStream.js'; +import { + MAX_PARTIAL_FRAME_CHARS, + MAX_RETAINED_CONTENT_CHARS, + MAX_RETAINED_TOOLCALL_CHARS, + consumeChatCompletionsStream, + reduceChatChunks, +} from './chatStream.js'; + +/** An SSE body that never emits a newline — the partial-frame buffer must not grow unboundedly. */ +function endlessFrameResponse(hugeLine: string, tail: string): Response { + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(hugeLine)); + controller.enqueue(new TextEncoder().encode(`\n${tail}\n\n`)); + controller.close(); + }, + }); + return new Response(body, { status: 200 }); +} + +describe('chat stream retention bounds (AGT-3429)', () => { + it('keeps the partial-frame buffer bounded when a frame never terminates', async () => { + const hugeFrame = `data: ${'x'.repeat(MAX_PARTIAL_FRAME_CHARS * 2)}`; + const result = await consumeChatCompletionsStream(endlessFrameResponse( + hugeFrame, + `data: ${JSON.stringify({ choices: [{ delta: { content: 'tail' }, finish_reason: 'stop' }] })}`, + )); + + const content = result.choices[0]?.message.content; + expect(typeof content).toBe('string'); + expect(content!.length).toBeLessThanOrEqual(MAX_PARTIAL_FRAME_CHARS); + }); + + it('caps retained content once the stream exceeds the hard limit', async () => { + const chunk = `data: ${JSON.stringify({ choices: [{ delta: { content: 'y'.repeat(64 * 1024) } }] })}\n\n`; + const body = new ReadableStream({ + start(controller) { + const encoder = new TextEncoder(); + // 2 MiB of deltas — twice MAX_RETAINED_CONTENT_CHARS. + for (let i = 0; i < 32; i++) controller.enqueue(encoder.encode(chunk)); + controller.enqueue(encoder.encode('data: [DONE]\n\n')); + controller.close(); + }, + }); + + const seen: string[] = []; + const result = await consumeChatCompletionsStream( + new Response(body, { status: 200 }), + (delta) => seen.push(delta), + ); + + const content = result.choices[0]?.message.content; + expect(content!.length).toBeLessThanOrEqual(MAX_RETAINED_CONTENT_CHARS); + expect(seen.join('').length).toBeLessThanOrEqual(MAX_RETAINED_CONTENT_CHARS); + }); + + it('retains a tool call whose arguments stream in many fragments', async () => { + // Regression: evicting old chunks to bound memory once corrupted the head of + // a streamed tool call, handing the tool layer invalid JSON. + const total = 3000; + const expectedArgs = `{"path":"a.ts","content":"${'A'.repeat(4000)}"}`; + const body = new ReadableStream({ + start(controller) { + const encoder = new TextEncoder(); + const step = Math.ceil(expectedArgs.length / total); + for (let i = 0; i < expectedArgs.length; i += step) { + controller.enqueue(encoder.encode(`data: ${JSON.stringify({ + choices: [{ delta: { tool_calls: [{ index: 0, id: 'call_1', function: { name: 'write_file', arguments: expectedArgs.slice(i, i + step) } }] } }], + })}\n\n`)); + } + controller.enqueue(encoder.encode(`data: ${JSON.stringify({ choices: [{ delta: {}, finish_reason: 'tool_calls' }] })}\n\n`)); + controller.enqueue(encoder.encode('data: [DONE]\n\n')); + controller.close(); + }, + }); + + const res = await consumeChatCompletionsStream(new Response(body, { status: 200 })); + const call = res.choices[0].message.tool_calls?.[0]; + expect(call?.function.name).toBe('write_file'); + expect(call?.function.arguments).toBe(expectedArgs); + expect(() => JSON.parse(call!.function.arguments)).not.toThrow(); + expect(res.choices[0].finish_reason).toBe('tool_calls'); + }); + + it('keeps the head of a long reply, not a suffix', async () => { + const expected = Array.from({ length: 1500 }, (_, i) => `w${i} `).join(''); + const body = new ReadableStream({ + start(controller) { + const encoder = new TextEncoder(); + for (let i = 0; i < 1500; i++) { + controller.enqueue(encoder.encode(`data: ${JSON.stringify({ choices: [{ delta: { content: `w${i} ` } }] })}\n\n`)); + } + controller.enqueue(encoder.encode('data: [DONE]\n\n')); + controller.close(); + }, + }); + + const res = await consumeChatCompletionsStream(new Response(body, { status: 200 })); + expect(res.choices[0]?.message.content).toBe(expected); + }); + + it('refuses rather than silently corrupting tool-call arguments past the cap', async () => { + const huge = 'z'.repeat(MAX_RETAINED_TOOLCALL_CHARS + 1); + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify({ + choices: [{ delta: { tool_calls: [{ index: 0, id: 'c', function: { name: 'f', arguments: huge } }] } }], + })}\n\n`)); + controller.close(); + }, + }); + + await expect(consumeChatCompletionsStream(new Response(body, { status: 200 }))).rejects.toThrow( + /tool-call arguments exceed/, + ); + }); +}); describe('reduceChatChunks', () => { it('keeps the model the server reports serving', () => { diff --git a/src/adapters/chatStream.ts b/src/adapters/chatStream.ts index 5f9ab688..f63b3855 100644 --- a/src/adapters/chatStream.ts +++ b/src/adapters/chatStream.ts @@ -152,6 +152,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 retained assistant content assembled from the stream (1 MiB). */ +export const MAX_RETAINED_CONTENT_CHARS = 1024 * 1024; +/** Hard cap on streamed tool-call arguments retained for one call (1 MiB). */ +export const MAX_RETAINED_TOOLCALL_CHARS = 1024 * 1024; + /** Read a chat/completions SSE body and reduce it, emitting content deltas live. */ export async function consumeChatCompletionsStream( res: Response, @@ -162,26 +169,90 @@ export async function consumeChatCompletionsStream( const reader = res.body?.getReader(); if (!reader) throw new Error('chat stream: empty response body'); - const chunks: StreamChunk[] = []; const decoder = new TextDecoder(); let buffer = ''; + // Accumulate the reduced shape as chunks arrive instead of retaining every + // parsed chunk. Evicting old chunks to bound memory would corrupt exactly the + // parts that stream incrementally — a tool call's `arguments` fragments and + // the head of the reply — because both are assembled by concatenation; the + // agent then receives unparseable arguments or a reply that starts mid-word. + let content = ''; + let sawContent = false; + let finishReason = 'stop'; + let usage: ChatCompletionLike['usage']; + let model: string | undefined; + const calls = new Map(); + let retainedToolCallChars = 0; + const handle = (c: StreamChunk | null) => { if (!c) return; - const delta = c.choices?.[0]?.delta?.content; - if (onToken && typeof delta === 'string' && delta) onToken(delta); - chunks.push(c); + if (typeof c.model === 'string' && c.model) model = c.model; + if (c.usage) usage = normalizeChatUsage(c.usage); + const choice = c.choices?.[0]; + if (!choice) return; + const delta = choice.delta ?? {}; + if (typeof delta.content === 'string' && delta.content) { + sawContent = true; + // Emit live tokens, but stop retaining content past the hard cap so a + // runaway stream cannot grow the assembled reply without bound. + if (content.length < MAX_RETAINED_CONTENT_CHARS) { + const room = MAX_RETAINED_CONTENT_CHARS - content.length; + const emit = delta.content.length <= room ? delta.content : delta.content.slice(0, room); + content += emit; + if (onToken) onToken(emit); + } + } + for (const tc of delta.tool_calls ?? []) { + const idx = tc.index ?? 0; + const cur = calls.get(idx) ?? { id: '', name: '', args: '' }; + if (tc.id) cur.id = tc.id; + if (tc.function?.name) cur.name = tc.function.name; + if (tc.function?.arguments) { + retainedToolCallChars += tc.function.arguments.length; + if (retainedToolCallChars > MAX_RETAINED_TOOLCALL_CHARS) { + // Truncating mid-token would hand the tool layer invalid JSON, which + // is worse than refusing: say so instead of corrupting silently. + throw new Error(`chat stream: tool-call arguments exceed ${MAX_RETAINED_TOOLCALL_CHARS} chars`); + } + cur.args += tc.function.arguments; + } + calls.set(idx, cur); + } + if (choice.finish_reason) finishReason = choice.finish_reason; }; for (;;) { const { done, value } = await reader.read(); if (done) break; onBytes?.(); buffer += decoder.decode(value, { stream: true }); + // Split first, cap second: complete frames are always parsed, and only the + // unterminated tail is bounded. Truncating before the split would silently + // drop whole frames whenever one read delivered more than the cap. const lines = buffer.split('\n'); buffer = lines.pop() ?? ''; + if (buffer.length > MAX_PARTIAL_FRAME_CHARS) { + buffer = buffer.slice(-MAX_PARTIAL_FRAME_CHARS); + } for (const line of lines) handle(parseChunkLine(line)); } handle(parseChunkLine(buffer)); - // Final reduce WITHOUT onToken (already emitted above) to assemble the result. - return reduceChatChunks(chunks); + // Reduce a single synthetic chunk built from the accumulated state: same + // shape and precedence as reducing the whole stream, with the content cap and + // tool calls already applied. No onToken — deltas were emitted above. + const toolCalls: StreamToolCall[] = [...calls.values()] + .filter((c) => c.id && c.name) + .map((c) => ({ 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, + }], + usage, + ...(model ? { model } : {}), + }; } diff --git a/src/automation/ciWorker.test.ts b/src/automation/ciWorker.test.ts index 3458abb6..abd9c88c 100644 --- a/src/automation/ciWorker.test.ts +++ b/src/automation/ciWorker.test.ts @@ -274,8 +274,12 @@ describe('CIWorker', () => { const on = new CIWorker({ repos: ['o/r'], autoRetry: true }); await internal(on).investigateFailure('o/r', failure()); + // The rerun must be bounded so a hung `gh` cannot park the CI worker. expect(execFileMock).toHaveBeenCalledWith( - 'gh', ['run', 'rerun', '42', '-R', 'o/r', '--failed'], expect.any(Function), + 'gh', + ['run', 'rerun', '42', '-R', 'o/r', '--failed'], + expect.objectContaining({ timeout: 30_000, maxBuffer: 10 * 1024 * 1024 }), + expect.any(Function), ); }); diff --git a/src/automation/ciWorker.ts b/src/automation/ciWorker.ts index eb99bf93..6f65b182 100644 --- a/src/automation/ciWorker.ts +++ b/src/automation/ciWorker.ts @@ -251,7 +251,7 @@ export class CIWorker { private async retryRun(repo: string, runId: number): Promise { try { console.log(`[CIWorker] Retrying run: ${repo}#${runId}`); - await execFileAsync('gh', ['run', 'rerun', String(runId), '-R', repo, '--failed']); + await execFileAsync('gh', ['run', 'rerun', String(runId), '-R', repo, '--failed'], { timeout: 30_000, maxBuffer: 10 * 1024 * 1024 }); broadcastEvent({ type: 'log', diff --git a/src/cli/fixCommand.ts b/src/cli/fixCommand.ts index 9660c402..435e1df6 100644 --- a/src/cli/fixCommand.ts +++ b/src/cli/fixCommand.ts @@ -343,7 +343,9 @@ export interface FixReport { /** Default check runner: spawn the command, capture combined output, pass = exit 0. */ async function defaultRunCheck(check: Check, cwd: string): Promise<{ passed: boolean; output: string }> { return new Promise((resolve) => { - execFile(check.program, check.args, { cwd, maxBuffer: 32 * 1024 * 1024 }, (err, stdout, stderr) => { + // A check that hangs must not hang `openswarm fix`; the fix loop needs a + // failed check it can act on, not an indefinitely parked round. + execFile(check.program, check.args, { cwd, maxBuffer: 32 * 1024 * 1024, timeout: 300_000 }, (err, stdout, stderr) => { resolve({ passed: !err, output: `${stdout ?? ''}${stderr ?? ''}` }); }); }); diff --git a/src/core/eventHub.test.ts b/src/core/eventHub.test.ts index 61ff7294..3c7f535f 100644 --- a/src/core/eventHub.test.ts +++ b/src/core/eventHub.test.ts @@ -10,6 +10,10 @@ import { getStageBuffer, getChatBuffer, __resetForTests, + MAX_CHAT_TEXT_CHARS, + MAX_LOG_LINE_CHARS, + SSE_MAX_BUFFERED_BYTES, + SSE_STALL_TIMEOUT_MS, type HubEvent, } from './eventHub.js'; @@ -769,6 +773,108 @@ describe('eventHub', () => { }); }); + describe('payload and backpressure bounds (AGT-3429)', () => { + it('truncates an oversized log line before retaining it', () => { + broadcastEvent({ + type: 'log', + data: { taskId: 'task-1', stage: 'worker', line: 'L'.repeat(40_000) }, + }); + + const buffer = getLogBuffer(); + expect(buffer).toHaveLength(1); + const logged = buffer[0] as Extract; + expect(logged.data.line.length).toBeLessThanOrEqual(MAX_LOG_LINE_CHARS); + expect(logged.data.line.endsWith('…')).toBe(true); + }); + + it('truncates oversized chat text before retaining it', () => { + broadcastEvent({ + type: 'chat:user', + data: { text: 'C'.repeat(80_000), ts: Date.now() }, + }); + + const buffer = getChatBuffer(); + expect(buffer).toHaveLength(1); + const chat = buffer[0] as Extract; + expect(chat.data.text.length).toBeLessThanOrEqual(MAX_CHAT_TEXT_CHARS); + }); + + it('does not disconnect a healthy client during a same-tick burst', () => { + // Regression: counting writes destroyed a healthy reader, because a burst + // of synchronous broadcasts returns false repeatedly with no chance to + // drain in between (one stdout chunk fans out hundreds of log events). + const writes: Array<() => void> = []; + const destroy = vi.fn(); + const busyRes = { + write: vi.fn(() => false), + once: vi.fn((event: string, cb: () => void) => { + if (event === 'drain') writes.push(cb); + }), + removeListener: vi.fn(), + destroy, + } as unknown as ServerResponse; + + cleanupFunctions.push(addSSEClient(busyRes, true)); + for (let i = 0; i < 500; i++) { + // The reader keeps up: drain fires between batches. + if (i % 10 === 0) writes.forEach((cb) => cb()); + broadcastEvent({ type: 'log', data: { taskId: 'busy', stage: 'worker', line: `line-${i}` } }); + } + + expect(getActiveSSECount()).toBe(1); + expect(destroy).not.toHaveBeenCalled(); + }); + + it('disconnects a client that stays stalled past the stall window', () => { + vi.useFakeTimers(); + try { + const destroy = vi.fn(); + const stalledRes = { + write: vi.fn(() => false), + once: vi.fn(), + removeListener: vi.fn(), + destroy, + } as unknown as ServerResponse; + + cleanupFunctions.push(addSSEClient(stalledRes, true)); + expect(getActiveSSECount()).toBe(1); + + broadcastEvent({ type: 'log', data: { taskId: 'stall', stage: 'worker', line: 'one' } }); + // Still connected while inside the window... + vi.advanceTimersByTime(SSE_STALL_TIMEOUT_MS - 1); + expect(getActiveSSECount()).toBe(1); + + // ...and dropped once the window elapses without a drain. + vi.advanceTimersByTime(2); + expect(getActiveSSECount()).toBe(0); + expect(destroy).toHaveBeenCalled(); + } finally { + vi.useRealTimers(); + } + }); + + it('disconnects a client whose queued bytes exceed the buffer cap', () => { + const destroy = vi.fn(); + const stalledRes = { + write: vi.fn(() => false), + once: vi.fn(), + removeListener: vi.fn(), + destroy, + } as unknown as ServerResponse; + + cleanupFunctions.push(addSSEClient(stalledRes, true)); + // Lines are bounded to MAX_LOG_LINE_CHARS, so enough of them must be + // broadcast to queue past SSE_MAX_BUFFERED_BYTES. + const framesToFill = Math.ceil(SSE_MAX_BUFFERED_BYTES / MAX_LOG_LINE_CHARS) + 2; + for (let i = 0; i < framesToFill; i++) { + broadcastEvent({ type: 'log', data: { taskId: 'flood', stage: 'worker', line: 'B'.repeat(MAX_LOG_LINE_CHARS) } }); + } + + expect(getActiveSSECount()).toBe(0); + expect(destroy).toHaveBeenCalled(); + }); + }); + describe('per-task transcript store hooks (INT-3402)', () => { it('feeds log events into the per-task ring and lifecycles it on completion', async () => { const { getTaskLog, TASK_LOG_RETENTION_MS } = await import('./taskLogStore.js'); diff --git a/src/core/eventHub.ts b/src/core/eventHub.ts index b28629e1..8f8cb9f0 100644 --- a/src/core/eventHub.ts +++ b/src/core/eventHub.ts @@ -146,6 +146,181 @@ 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 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; +/** + * How long a socket may stay over its high-water mark before it is treated as + * stalled. Deliberately not a write count: a burst of events in one synchronous + * tick returns `false` many times without the reader being at fault. + */ +export const SSE_STALL_TIMEOUT_MS = 30_000; +/** Bytes queued for one client before it is disconnected. */ +export const SSE_MAX_BUFFERED_BYTES = 4 * 1024 * 1024; + +const bufferedBytes = new WeakMap(); +/** Clients with a `drain` listener already attached — one per socket, not one per write. */ +const drainListeners = new WeakSet(); +/** Pending stall timer per socket, cleared as soon as the socket drains. */ +const stallTimers = new WeakMap(); + +function disconnectClient(res: ServerResponse): void { + sseClients.delete(res); + bufferedBytes.delete(res); + drainListeners.delete(res); + const timer = stallTimers.get(res); + if (timer) { + clearTimeout(timer); + stallTimers.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 the workload-controlled string fields of an event before it is + * retained or fanned out, so one oversized event cannot exhaust memory + * through the replay ring, the per-type buffers, and every SSE client. + */ +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; null when even the bounded frame exceeds the cap. */ +function serializeEventFrame(event: HubEvent): string | null { + const frame = `data: ${JSON.stringify(event)}\n\n`; + if (Buffer.byteLength(frame, 'utf8') > MAX_EVENT_PAYLOAD_BYTES) { + return null; + } + return frame; +} + +/** + * Write one frame to a client, disconnecting only a client that is actually + * stuck. `broadcastEvent` is synchronous, so a burst of events in one tick + * gives the socket no chance to fire `drain` in between — counting writes would + * destroy a healthy reader mid-burst (a single 23 KB stdout chunk fans out + * hundreds of log events). Two bounds therefore decide: + * + * - queued bytes: the daemon must not buffer unbounded memory for a slow reader; + * - a stall window: a reader that has not drained for this long is gone, not slow. + */ +function writeToClient(res: ServerResponse, data: string): void { + try { + // `write()` returns false only when the socket buffer is over its high-water + // mark; `undefined` means there was nothing to report (e.g. a stub response). + const ok = (res.write(data) as boolean | undefined) !== false; + if (ok) { + bufferedBytes.set(res, 0); + const timer = stallTimers.get(res); + if (timer) { + clearTimeout(timer); + stallTimers.delete(res); + } + return; + } + + const nextBytes = (bufferedBytes.get(res) ?? 0) + Buffer.byteLength(data, 'utf8'); + bufferedBytes.set(res, nextBytes); + if (nextBytes >= SSE_MAX_BUFFERED_BYTES) { + disconnectClient(res); + return; + } + + // Attach the reset listener once per socket: res.once('drain') inside the + // broadcast loop would add one listener per write and trip Node's + // MaxListenersExceededWarning on a client that stays slow for a while. + if (!drainListeners.has(res)) { + drainListeners.add(res); + res.once('drain', () => { + drainListeners.delete(res); + bufferedBytes.set(res, 0); + }); + } + // A socket that stays over its high-water mark for the whole window is + // stalled rather than merely slow — the dashboard reconnects on its own. + if (!stallTimers.has(res)) { + const timer = setTimeout(() => { + stallTimers.delete(res); + if ((bufferedBytes.get(res) ?? 0) > 0) disconnectClient(res); + }, SSE_STALL_TIMEOUT_MS); + timer.unref(); + stallTimers.set(res, timer); + } + } catch { + disconnectClient(res); + } +} + function pushReplay(event: HubEvent): void { if (event.type === 'log') { // Keep only recent log lines in replay buffer to avoid bloat @@ -168,34 +343,57 @@ export function getEventHub(): EventEmitter { } export function broadcastEvent(event: HubEvent): void { - // Skip replaying heartbeat/stats to avoid noise on reconnect - if (event.type !== 'heartbeat') { - pushReplay(event); + // Bound the workload-controlled fields before any retain / serialize / + // fan-out, and mirror the bounded values back onto the caller's object so an + // emitter holding it observes exactly what was broadcast. + const bounded = boundEventPayload(event); + const sourceLog = event.type === 'log' ? event : null; + const sourceChat = event.type === 'chat:user' || event.type === 'chat:agent' ? event : null; + if (sourceLog && bounded.type === 'log') { + sourceLog.data.line = bounded.data.line; + } + if (sourceChat && (bounded.type === 'chat:user' || bounded.type === 'chat:agent')) { + sourceChat.data.text = bounded.data.text; } + // 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') { // 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 // millisecond — is the join key: an agent emits several lines per ms. // (INT-3402) const ts = Date.now(); - event.data.ts = ts; - event.data.seq = appendTaskLog(event.data.taskId, event.data.stage, event.data.line, ts); + const seq = appendTaskLog(bounded.data.taskId, bounded.data.stage, bounded.data.line, ts); // 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(); + if (sourceLog) { + sourceLog.data.ts = ts; + sourceLog.data.seq = seq; + sourceLog.data.gen = gen; + } + bounded.data.ts = ts; + bounded.data.seq = seq; + bounded.data.gen = gen; + } else if (bounded.type === 'task:started') { + cancelTaskLogCleanup(bounded.data.taskId); + } else if (bounded.type === 'task:completed') { + scheduleTaskLogCleanup(bounded.data.taskId); + } + + const frame = serializeEventFrame(bounded); + if (frame === null) return; + + if (bounded.type !== 'heartbeat') { + pushReplay(bounded); } // 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 +413,17 @@ 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`; for (const res of sseClients) { - try { - res.write(data); - } catch { - sseClients.delete(res); - } + writeToClient(res, frame); } } @@ -239,7 +432,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 +443,12 @@ export function addSSEClient(res: ServerResponse, skipReplay = false): () => voi } } sseClients.add(res); + bufferedBytes.set(res, 0); // Cleanup function that removes client from set const cleanup = () => { sseClients.delete(res); + bufferedBytes.delete(res); // Remove the close listener after cleanup to prevent memory leak res.removeListener('close', cleanup); }; diff --git a/src/knowledge/gitInfo.test.ts b/src/knowledge/gitInfo.test.ts new file mode 100644 index 00000000..87f5db92 --- /dev/null +++ b/src/knowledge/gitInfo.test.ts @@ -0,0 +1,52 @@ +import { describe, expect, it } from 'vitest'; +import { parseNulDelimitedChurnOutput } from './gitInfo.js'; + +/** + * Real `git log -z --format=%x1e%ct --name-only` output: commits are + * NUL-separated, each commit's first filename is prefixed with the newline + * that terminates the format, and there is NO empty token between commits. + */ +const RS = '\x1e'; +const commit = (ts: number, files: string[]) => [`${RS}${ts}`, ...files.map((f) => `\n${f}`)].join('\0'); + +describe('parseNulDelimitedChurnOutput', () => { + it('counts a numeric filename as a file, never as a commit boundary', () => { + const churns = parseNulDelimitedChurnOutput(`${commit(1700000000, ['12345'])}\0`); + expect([...churns.keys()]).toEqual(['12345']); + expect(churns.get('12345')).toEqual({ + path: '12345', + commitCount: 1, + lastCommitDate: 1700000000 * 1000, + }); + }); + + it('attributes a multi-file commit to its own timestamp', () => { + const churns = parseNulDelimitedChurnOutput(`${commit(1700000000, ['src/a.ts', 'src/b.ts'])}\0`); + expect(churns.size).toBe(2); + expect(churns.get('src/a.ts')?.lastCommitDate).toBe(1700000000 * 1000); + expect(churns.get('src/b.ts')?.lastCommitDate).toBe(1700000000 * 1000); + }); + + it('keeps each commit separate — a later timestamp is not read as a filename', () => { + // Regression: a state machine that expects an empty token between commits + // treats 1700001000 as a path, inventing a file and dropping src/c.ts. + const churns = parseNulDelimitedChurnOutput([ + commit(1700000000, ['src/a.ts']), + commit(1700001000, ['src/a.ts', 'src/c.ts']), + '', + ].join('\0')); + + expect([...churns.keys()].sort()).toEqual(['src/a.ts', 'src/c.ts']); + expect(churns.get('src/a.ts')).toEqual({ + path: 'src/a.ts', + commitCount: 2, + lastCommitDate: 1700001000 * 1000, + }); + expect(churns.get('src/c.ts')?.commitCount).toBe(1); + }); + + it('ignores git output with no churn (empty and whitespace-only)', () => { + expect(parseNulDelimitedChurnOutput('').size).toBe(0); + expect(parseNulDelimitedChurnOutput('\0\0').size).toBe(0); + }); +}); diff --git a/src/knowledge/gitInfo.ts b/src/knowledge/gitInfo.ts index add8d33d..004314cc 100644 --- a/src/knowledge/gitInfo.ts +++ b/src/knowledge/gitInfo.ts @@ -44,6 +44,56 @@ interface FileChurn { lastCommitDate: number; } +/** + * `git log --format=%x1e%ct` prefixes every timestamp with an ASCII + * record-separator, so a timestamp token is self-identifying. Without it a + * numeric filename (`12345`) is indistinguishable from a commit timestamp — + * main's parser dropped such files entirely, and a state machine that assumes + * an empty NUL token between commits misreads the *next* commit's timestamp as + * a filename, because real `git log -z` output has no such empty token. + */ +const CHURN_TIMESTAMP_SENTINEL = '\x1e'; + +/** + * Parse NUL-delimited `git log -z --format=%x1e%ct --name-only` output into + * per-file churn counts. Tokens prefixed with the sentinel are timestamps; + * every other token is a filename and is NEVER parsed as a number, even when + * the file is named `12345`. Git prefixes each commit's first filename with the + * newline that terminates the format, which is stripped here. + */ +export function parseNulDelimitedChurnOutput(output: string): Map { + const churns = new Map(); + let currentTimestamp = 0; + + for (const token of output.split('\0')) { + if (!token) continue; + + if (token.startsWith(CHURN_TIMESTAMP_SENTINEL)) { + const trimmed = token.slice(CHURN_TIMESTAMP_SENTINEL.length).trim(); + currentTimestamp = /^\d+$/.test(trimmed) ? parseInt(trimmed, 10) * 1000 : 0; + continue; + } + + const filePath = token.startsWith('\n') ? token.slice(1) : token; + if (!filePath) continue; + const existing = churns.get(filePath); + if (existing) { + existing.commitCount++; + if (currentTimestamp > existing.lastCommitDate) { + existing.lastCommitDate = currentTimestamp; + } + } else { + churns.set(filePath, { + path: filePath, + commitCount: 1, + lastCommitDate: currentTimestamp, + }); + } + } + + return churns; +} + /** * Calculate per-file commit count over the last 30 days */ @@ -51,49 +101,20 @@ async function getFileChurns(projectPath: string, sinceDays: number = 30): Promi const churns = new Map(); try { - // git log --since="30 days ago" --name-only --format="%ct" + // git log --since="30 days ago" --name-only --format="%x1e%ct" const output = await runGitCommand(projectPath, [ 'log', `--since=${sinceDays} days ago`, '--name-only', '-z', - '--format=%ct', + '--format=%x1e%ct', ]); - let currentTimestamp = 0; - - for (const token of output.split('\0')) { - if (!token) continue; - const timestampToken = token.trim(); - - // If numeric, it's a commit timestamp - if (/^\d+$/.test(timestampToken)) { - currentTimestamp = parseInt(timestampToken, 10) * 1000; // Convert to ms - continue; - } - - // `-z` preserves embedded newlines and other whitespace in filenames. - const filePath = token.startsWith('\n') ? token.slice(1) : token; - if (!filePath) continue; - const existing = churns.get(filePath); - if (existing) { - existing.commitCount++; - if (currentTimestamp > existing.lastCommitDate) { - existing.lastCommitDate = currentTimestamp; - } - } else { - churns.set(filePath, { - path: filePath, - commitCount: 1, - lastCommitDate: currentTimestamp, - }); - } - } + return parseNulDelimitedChurnOutput(output); } catch (err) { console.warn(`[GitInfo] Failed to get file churns:`, err); + return churns; } - - return churns; } /** diff --git a/src/support/chatBackend.ts b/src/support/chatBackend.ts index 81bd7af1..0fed9e4b 100644 --- a/src/support/chatBackend.ts +++ b/src/support/chatBackend.ts @@ -359,9 +359,15 @@ 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 the chat UI (AGT-3429). + // Keep the TAIL: the consumer of this value is looking for the reply, which + // is emitted last, not the earliest output. + const MAX_ADAPTER_STDOUT_CHARS = 1024 * 1024; + const text = raw.stdout.length > MAX_ADAPTER_STDOUT_CHARS + ? raw.stdout.slice(-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 }; @@ -464,6 +470,11 @@ export async function runChatCompletion(options: ChatCompletionOptions): Promise proc.stdin?.end(stdin); } + // Hard caps on retained CLI chat output (AGT-3429). Oversized streams are + // bounded 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 = ''; @@ -472,6 +483,13 @@ export async function runChatCompletion(options: ChatCompletionOptions): Promise let thinkingTimer: NodeJS.Timeout | null = null; let settled = false; + /** Keep the newest `max` chars — the tail is what a consumer still needs. */ + const keepTail = (current: string, chunk: string, max: number): string => { + if (chunk.length >= max) return chunk.slice(-max); + const excess = current.length + chunk.length - max; + return excess > 0 ? current.slice(excess) + chunk : current + chunk; + }; + const cleanupProcessHooks = () => { if (thinkingTimer) clearTimeout(thinkingTimer); runSignal.removeEventListener('abort', onAbort); @@ -540,13 +558,17 @@ export async function runChatCompletion(options: ChatCompletionOptions): Promise proc.stdout?.on('data', (chunk: Buffer) => { const text = chunk.toString(); - stdout += text; - buffer += text; + // Both keep the TAIL: `extractChatResponse` reads the LAST matching + // stream event, and `flushLines` can only emit lines it still has — a + // head-truncated line buffer would never find a newline again and live + // streaming would stop for the rest of the turn. + stdout = keepTail(stdout, text, MAX_CHAT_STDOUT_CHARS); + buffer = keepTail(buffer, text, MAX_CHAT_PARTIAL_BUFFER_CHARS); flushLines(false); }); proc.stderr?.on('data', (chunk: Buffer) => { - stderr += chunk.toString(); + stderr = keepTail(stderr, chunk.toString(), MAX_CHAT_STDERR_CHARS); }); proc.on('close', (code) => { diff --git a/src/support/gitStatus.test.ts b/src/support/gitStatus.test.ts index eb650511..8e306c0d 100644 --- a/src/support/gitStatus.test.ts +++ b/src/support/gitStatus.test.ts @@ -17,4 +17,19 @@ describe('git status cache', () => { } expect(getGitStatusCacheSizeForTests()).toBe(200); }); + + it('calls git with a maxBuffer large enough for big repos', async () => { + execFile.mockImplementation((_command, _args, options, callback) => callback(new Error('not a repo'), '')); + await getProjectGitInfo('/repo/maxbuffer-probe'); + + expect(execFile).toHaveBeenCalled(); + const withBounds = execFile.mock.calls.filter( + (call) => typeof call[2] === 'object' && call[2] !== null, + ); + expect(withBounds.length).toBeGreaterThan(0); + for (const call of withBounds) { + const options = call[2] as { maxBuffer?: number }; + expect(options.maxBuffer).toBeGreaterThanOrEqual(10 * 1024 * 1024); + } + }); }); diff --git a/src/support/gitStatus.ts b/src/support/gitStatus.ts index 0a14b333..bbf52b11 100644 --- a/src/support/gitStatus.ts +++ b/src/support/gitStatus.ts @@ -33,14 +33,19 @@ const cache = new Map(); const CACHE_TTL = 30_000; const MAX_CACHE_ENTRIES = 200; const CMD_TIMEOUT = 5_000; +const CMD_MAX_BUFFER = 10 * 1024 * 1024; let activePoller: NodeJS.Timeout | null = null; // --- Helpers --- +/** + * Run git and return its trimmed stdout. Rejects on failure: a caller that + * cannot run git must not mistake the empty string for a clean result. + */ function git(projectPath: string, args: string[]): Promise { - return new Promise((resolve) => { - execFile('git', ['-C', projectPath, ...args], { timeout: CMD_TIMEOUT }, (err, stdout) => { - if (err) { resolve(''); return; } + return new Promise((resolve, reject) => { + execFile('git', ['-C', projectPath, ...args], { timeout: CMD_TIMEOUT, maxBuffer: CMD_MAX_BUFFER }, (err, stdout) => { + if (err) { reject(err); return; } resolve(stdout.trim()); }); }); @@ -48,7 +53,7 @@ function git(projectPath: string, args: string[]): Promise { function gh(args: string[]): Promise { return new Promise((resolve) => { - execFile('gh', args, { timeout: CMD_TIMEOUT }, (err, stdout) => { + execFile('gh', args, { timeout: CMD_TIMEOUT, maxBuffer: CMD_MAX_BUFFER }, (err, stdout) => { if (err) { resolve(''); return; } resolve(stdout.trim()); }); @@ -58,20 +63,32 @@ function gh(args: string[]): Promise { // --- Fetch functions --- async function fetchGitStatus(projectPath: string): Promise { - const branch = await git(projectPath, ['branch', '--show-current']); + let branch: string; + let porcelain: string; + try { + branch = await git(projectPath, ['branch', '--show-current']); + porcelain = await git(projectPath, ['status', '--porcelain']); + } catch { + // Failure must not look like a clean tree. + return null; + } if (!branch) return null; // not a git repo or error - const porcelain = await git(projectPath, ['status', '--porcelain']); const lines = porcelain ? porcelain.split('\n').filter(Boolean) : []; - // ahead/behind + // ahead/behind — may catch and treat as 0 let ahead = 0; let behind = 0; - const revList = await git(projectPath, ['rev-list', '--left-right', '--count', 'HEAD...@{u}']); - if (revList) { - const parts = revList.split(/\s+/); - ahead = parseInt(parts[0], 10) || 0; - behind = parseInt(parts[1], 10) || 0; + try { + const revList = await git(projectPath, ['rev-list', '--left-right', '--count', 'HEAD...@{u}']); + if (revList) { + const parts = revList.split(/\s+/); + ahead = parseInt(parts[0], 10) || 0; + behind = parseInt(parts[1], 10) || 0; + } + } catch { + ahead = 0; + behind = 0; } return { @@ -85,7 +102,12 @@ async function fetchGitStatus(projectPath: string): Promise { async function fetchOpenPRs(projectPath: string): Promise { // Extract owner/repo from origin remote URL - const remoteUrl = await git(projectPath, ['remote', 'get-url', 'origin']); + let remoteUrl: string; + try { + remoteUrl = await git(projectPath, ['remote', 'get-url', 'origin']); + } catch { + return []; + } if (!remoteUrl) return []; // SSH: git@github.com:owner/repo.git / HTTPS: https://github.com/owner/repo.git diff --git a/src/support/httpBody.test.ts b/src/support/httpBody.test.ts new file mode 100644 index 00000000..8c2f239f --- /dev/null +++ b/src/support/httpBody.test.ts @@ -0,0 +1,38 @@ +import { EventEmitter } from 'node:events'; +import { describe, expect, it } from 'vitest'; +import type { IncomingMessage } from 'node:http'; +import { readBody, HttpError } from './httpBody.js'; + +function mockRequest(chunks: Buffer[]): IncomingMessage { + const req = new EventEmitter() as IncomingMessage; + queueMicrotask(() => { + for (const chunk of chunks) req.emit('data', chunk); + req.emit('end'); + }); + return req; +} + +describe('readBody UTF-8 chunk boundaries', () => { + it('decodes a multi-byte character split across two chunks', async () => { + // Korean '한' is UTF-8: EA B5 98 — split after the first byte. + const full = Buffer.from('한', 'utf8'); + expect(full.length).toBe(3); + const body = await readBody(mockRequest([full.subarray(0, 1), full.subarray(1)])); + expect(body).toBe('한'); + }); + + it('decodes an emoji split across chunk boundaries', async () => { + // 😀 is F0 9F 98 80 + const full = Buffer.from('😀', 'utf8'); + const body = await readBody(mockRequest([full.subarray(0, 2), full.subarray(2)])); + expect(body).toBe('😀'); + }); + + it('rejects oversized bodies with HttpError 413', async () => { + const req = new EventEmitter() as IncomingMessage; + const pending = readBody(req); + const big = Buffer.alloc(1024 * 1024 + 1, 0x61); + queueMicrotask(() => req.emit('data', big)); + await expect(pending).rejects.toMatchObject({ statusCode: 413 } satisfies Partial); + }); +}); diff --git a/src/support/httpBody.ts b/src/support/httpBody.ts index ef613efb..d838db62 100644 --- a/src/support/httpBody.ts +++ b/src/support/httpBody.ts @@ -20,6 +20,9 @@ export class HttpError extends Error { export function readBody(req: IncomingMessage): Promise { return new Promise((resolve, reject) => { + // A multi-byte character split across chunks must survive: decode + // incrementally and flush the decoder's tail at end-of-body. + const decoder = new TextDecoder('utf-8'); let data = ''; let totalBytes = 0; let settled = false; @@ -37,11 +40,12 @@ export function readBody(req: IncomingMessage): Promise { fail(413, 'Request body too large'); return; } - data += chunk.toString('utf-8'); + data += decoder.decode(chunk, { stream: true }); }); req.on('end', () => { if (settled) return; settled = true; + data += decoder.decode(); resolve(data); }); req.on('aborted', () => fail(400, 'Request body aborted')); diff --git a/src/support/rollback.ts b/src/support/rollback.ts index 9d5983e9..e238a76f 100644 --- a/src/support/rollback.ts +++ b/src/support/rollback.ts @@ -54,8 +54,8 @@ function checkpointStashMessage(executionId: string): string { } /** - * Current `stash@{N}` for a stash identified by its message, or undefined if it - * is no longer in the list. + * Current `stash@{N}` for a stash identified by its message, or undefined if + * it is no longer in the list. * * `stash@{N}` is a POSITION, not an identity: every `git stash push` inserts at * 0 and shifts everything down. The checkpoint's index was captured at creation @@ -64,11 +64,28 @@ function checkpointStashMessage(executionId: string): string { * itself, immediately before popping, so it reliably restored the * `rollback-preserve-*` stash it had just made and orphaned the checkpoint's. * Resolving by message at pop time is stable under that shifting. + * + * The comparison is EXACT, never a substring test: execution IDs `abc` and + * `abcd` overlap, and `includes('abc')` picks whichever entry happens to come + * first. `git stash push -m MSG` records the reflog subject as + * `On : MSG`, so the message must match the whole subject or its + * `: MSG` tail — `subject.includes('abc')` would still select `…checkpoint-abcd`. */ async function resolveStashRef(projectPath: string, message: string): Promise { - const { stdout } = await gitExec(projectPath, 'stash', 'list'); - const line = stdout.split('\n').find((entry) => entry.includes(message)); - return line?.match(/stash@\{\d+\}/)?.[0]; + const { stdout } = await gitExec(projectPath, 'stash', 'list', '--pretty=format:%gd %gs'); + const tail = `: ${message}`; + const entries: Array<{ ref: string; subject: string }> = []; + for (const line of stdout.split('\n')) { + if (!line) continue; + const spaceIdx = line.indexOf(' '); + if (spaceIdx === -1) continue; + entries.push({ ref: line.slice(0, spaceIdx), subject: line.slice(spaceIdx + 1) }); + } + // A whole-subject match wins; otherwise the reflog-subject suffix + // (git stash push -m MSG records "On : MSG") must match exactly. + const match = entries.find((e) => e.subject === message) + ?? entries.find((e) => e.subject.endsWith(tail)); + return match?.ref.match(/stash@\{\d+\}/)?.[0] ?? match?.ref; } const CheckpointSchema = z.object({ @@ -235,12 +252,8 @@ export async function createCheckpoint( const stashMessage = checkpointStashMessage(executionId); await gitExec(expandedPath, 'stash', 'push', '-m', stashMessage, '--include-untracked'); - // Find Stash ID - const { stdout } = await gitExec(expandedPath, 'stash', 'list'); - const stashLine = stdout.split('\n').find(line => line.includes(stashMessage)); - if (stashLine) { - stashId = stashLine.match(/stash@\{(\d+)\}/)?.[0]; - } + // Exact message identity — never includes() — see resolveStashRef. + stashId = await resolveStashRef(expandedPath, stashMessage); } const checkpoint: Checkpoint = { diff --git a/src/support/rollbackStashIdentity.test.ts b/src/support/rollbackStashIdentity.test.ts index 05a99c42..9a6b839c 100644 --- a/src/support/rollbackStashIdentity.test.ts +++ b/src/support/rollbackStashIdentity.test.ts @@ -118,6 +118,30 @@ describe('checkpoint stash identity', () => { expect(result.success).toBe(true); expect(await readFile(join(repo, 'tracked.txt'), 'utf8')).toBe('committed\n'); }); + + it('exact message match: overlapping execution IDs abc vs abcd restore only abc', async () => { + // Prefix-overlapping IDs must not select the wrong stash via includes(). + const { createCheckpoint, rollbackToCheckpoint } = await loadRollback(); + + await writeFile(join(repo, 'tracked.txt'), 'ABC CHECKPOINT\n', 'utf8'); + const ckAbc = await createCheckpoint('abc', repo); + expect(ckAbc.stashId).toBeDefined(); + + await writeFile(join(repo, 'tracked.txt'), 'ABCD CHECKPOINT\n', 'utf8'); + const ckAbcd = await createCheckpoint('abcd', repo); + expect(ckAbcd.stashId).toBeDefined(); + + // Intervening unrelated stash shifts indices further. + await writeFile(join(repo, 'tracked.txt'), 'INTERVENING\n', 'utf8'); + execFileSync('git', ['-C', repo, 'stash', 'push', '-m', 'intervening', '--include-untracked'], { stdio: 'pipe' }); + + const result = await rollbackToCheckpoint(ckAbc.id, 'reset_hard'); + + expect(result.success).toBe(true); + expect(await readFile(join(repo, 'tracked.txt'), 'utf8')).toBe('ABC CHECKPOINT\n'); + // abcd's stash must still be present — we must not have popped it by accident. + expect(git('stash', 'list')).toContain('openswarm-checkpoint-abcd'); + }); }); describe('test harness', () => { diff --git a/src/support/workSessionRoutes.ts b/src/support/workSessionRoutes.ts index 7f19320e..b592a197 100644 --- a/src/support/workSessionRoutes.ts +++ b/src/support/workSessionRoutes.ts @@ -278,26 +278,30 @@ export async function tryHandleWorkSessionRoutes( // would appear in `files` with no patch to show. `--intent-to-add` on a // throwaway index makes git emit their content as an addition without // touching the worktree's real index. (review finding) + // + // Use canonicalWorktree (the containment-validated path) for all I/O, not + // resolved.worktreePath, so a symlink replacement race after validation + // cannot redirect the diff onto another tree. const [files, diff] = await Promise.all([ - getWorkingDiffDetail(resolved.worktreePath), - getDiffText(resolved.worktreePath, undefined, maxBytes, { includeUntracked: true }), + getWorkingDiffDetail(canonicalWorktree), + getDiffText(canonicalWorktree, undefined, maxBytes, { includeUntracked: true }), ]); // Both helpers swallow git errors into []/'' (they are advisory elsewhere). // Here that would render as "no changes" on a broken worktree — report the // ambiguity instead of a clean-looking lie. (review finding) if (files.length === 0 && !diff) { const { isGitRepo } = await import('./gitTracker.js'); - if (!(await isGitRepo(resolved.worktreePath))) { + if (!(await isGitRepo(canonicalWorktree))) { writeJson(res, 409, { error: `Worktree for task ${taskId} is no longer a valid git repository`, - worktreePath: resolved.worktreePath, + worktreePath: canonicalWorktree, }); return true; } } writeJson(res, 200, { taskId, - worktreePath: resolved.worktreePath, + worktreePath: canonicalWorktree, branch: resolved.branch, files, diff, 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..7dd48b97 100644 --- a/src/tui/components/LogLine.tsx +++ b/src/tui/components/LogLine.tsx @@ -3,10 +3,26 @@ 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 or + * inject fake rows into the log panel. + */ +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/src/tui/inputDebug.test.ts b/src/tui/inputDebug.test.ts index 6f045b1f..9b2bf8b3 100644 --- a/src/tui/inputDebug.test.ts +++ b/src/tui/inputDebug.test.ts @@ -1,28 +1,36 @@ -import { describe, it, expect } from 'vitest'; -import { mkdtempSync, rmSync, readFileSync, statSync } from 'node:fs'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { mkdtempSync, rmSync, readFileSync, statSync, existsSync, mkdirSync, symlinkSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; -import { formatInputDebug, inputDebugEnabled, appendInputDebug } from './inputDebug.js'; + +vi.mock('node:os', async (importOriginal) => { + const actual = await importOriginal(); + return { ...actual, homedir: () => process.env.OSW_TEST_HOME ?? actual.homedir() }; +}); describe('formatInputDebug (INT-1964)', () => { - it('shows code points so multibyte doubling is visible', () => { + it('shows code points so multibyte doubling is visible', async () => { + const { formatInputDebug } = await import('./inputDebug.js'); expect(formatInputDebug('이')).toBe('input="이" len=1 cp=[51060]'); // ink-level doubling would surface as two code points in ONE event: expect(formatInputDebug('이이')).toBe('input="이이" len=2 cp=[51060,51060]'); }); - it('records active key flags', () => { + it('records active key flags', async () => { + const { formatInputDebug } = await import('./inputDebug.js'); expect(formatInputDebug('', { return: true })).toContain('keys=return'); expect(formatInputDebug('a', { ctrl: true, meta: true })).toContain('keys=ctrl+meta'); }); - it('ascii is single code point (the non-doubled case)', () => { + it('ascii is single code point (the non-doubled case)', async () => { + const { formatInputDebug } = await import('./inputDebug.js'); expect(formatInputDebug(' ')).toBe('input=" " len=1 cp=[32]'); }); }); describe('inputDebugEnabled (INT-1964)', () => { - it('honors OPENSWARM_DEBUG_INPUT truthy values', () => { + it('honors OPENSWARM_DEBUG_INPUT truthy values', async () => { + const { inputDebugEnabled } = await import('./inputDebug.js'); expect(inputDebugEnabled({ OPENSWARM_DEBUG_INPUT: '1' } as NodeJS.ProcessEnv)).toBe(true); expect(inputDebugEnabled({ OPENSWARM_DEBUG_INPUT: 'true' } as NodeJS.ProcessEnv)).toBe(true); expect(inputDebugEnabled({ OPENSWARM_DEBUG_INPUT: '0' } as NodeJS.ProcessEnv)).toBe(false); @@ -31,22 +39,66 @@ describe('inputDebugEnabled (INT-1964)', () => { }); describe('appendInputDebug (INT-1964)', () => { - it('appends diagnostic lines and never throws', () => { - const dir = mkdtempSync(join(tmpdir(), 'indbg-')); - try { - const path = join(dir, 'nested', 'input-debug.log'); - appendInputDebug('이', {}, path); - appendInputDebug('a', { return: true }, path); - const lines = readFileSync(path, 'utf8').trim().split('\n'); - expect(lines).toHaveLength(2); - expect(lines[0]).toContain('cp=[51060]'); - expect(statSync(path).mode & 0o777).toBe(0o600); - } finally { - rmSync(dir, { recursive: true, force: true }); - } - }); - - it('swallows write errors (invalid path)', () => { + let sandbox: string; + let previousHome: string | undefined; + + beforeEach(() => { + sandbox = mkdtempSync(join(tmpdir(), 'indbg-')); + previousHome = process.env.OSW_TEST_HOME; + process.env.OSW_TEST_HOME = join(sandbox, 'home'); + mkdirSync(join(process.env.OSW_TEST_HOME, '.openswarm'), { recursive: true }); + vi.resetModules(); + }); + + afterEach(() => { + if (previousHome === undefined) delete process.env.OSW_TEST_HOME; + else process.env.OSW_TEST_HOME = previousHome; + rmSync(sandbox, { recursive: true, force: true }); + }); + + it('appends diagnostic lines under the ~/.openswarm sandbox', async () => { + const { appendInputDebug } = await import('./inputDebug.js'); + const path = join(process.env.OSW_TEST_HOME!, '.openswarm', 'nested', 'input-debug.log'); + appendInputDebug('이', {}, path); + appendInputDebug('a', { return: true }, path); + const lines = readFileSync(path, 'utf8').trim().split('\n'); + expect(lines).toHaveLength(2); + expect(lines[0]).toContain('cp=[51060]'); + expect(statSync(path).mode & 0o777).toBe(0o600); + }); + + it('swallows write errors (NUL in path)', async () => { + const { appendInputDebug } = await import('./inputDebug.js'); expect(() => appendInputDebug('x', {}, '/this/should/not/exist/\0/bad')).not.toThrow(); }); + + it('refuses to write outside the sandbox (no throw, no file)', async () => { + const { appendInputDebug } = await import('./inputDebug.js'); + const outside = join(sandbox, 'outside-escape.log'); + expect(() => appendInputDebug('escape', {}, outside)).not.toThrow(); + expect(existsSync(outside)).toBe(false); + }); + + it('refuses a symlinked directory inside the sandbox that points outside', async () => { + const { appendInputDebug } = await import('./inputDebug.js'); + const outsideDir = join(sandbox, 'outside-dir'); + mkdirSync(outsideDir, { recursive: true }); + const link = join(process.env.OSW_TEST_HOME!, '.openswarm', 'escape'); + symlinkSync(outsideDir, link); + + expect(() => appendInputDebug('symlink', {}, join(link, 'pwned.log'))).not.toThrow(); + expect(existsSync(join(outsideDir, 'pwned.log'))).toBe(false); + }); + + it('refuses traversal and sibling-prefix paths', async () => { + const { appendInputDebug } = await import('./inputDebug.js'); + const traversal = join(process.env.OSW_TEST_HOME!, '.openswarm', '..', '..', 'outside', 'a.log'); + const siblingPrefix = join(sandbox, 'outside', '.openswarm-evil.log'); + + appendInputDebug('traversal', {}, traversal); + appendInputDebug('sibling', {}, siblingPrefix); + + expect(existsSync(join(sandbox, 'outside', 'a.log'))).toBe(false); + expect(existsSync(siblingPrefix)).toBe(false); + }); }); diff --git a/src/tui/inputDebug.ts b/src/tui/inputDebug.ts index 66d5dfc1..a136c62b 100644 --- a/src/tui/inputDebug.ts +++ b/src/tui/inputDebug.ts @@ -10,9 +10,9 @@ // keypress logs one code point but the screen shows two glyphs, it's terminal // echo (fix in the client); if it logs the code point twice, it's ink-level. -import { closeSync, mkdirSync, openSync, writeFileSync } from 'node:fs'; +import { closeSync, mkdirSync, openSync, realpathSync, writeFileSync } from 'node:fs'; import { homedir } from 'node:os'; -import { join, dirname } from 'node:path'; +import { join, dirname, isAbsolute, relative, resolve } from 'node:path'; export const INPUT_DEBUG_LOG = join(homedir(), '.openswarm', 'input-debug.log'); @@ -46,11 +46,56 @@ export function inputDebugEnabled(env: NodeJS.ProcessEnv = process.env): boolean return v === '1' || v === 'true'; } +/** + * Resolve `path` and refuse anything that escapes ~/.openswarm. Diagnostics are + * best-effort, but they must not be a write primitive for an arbitrary path. + * + * `resolve()` is lexical, so a symlinked directory inside the sandbox could + * still redirect the write outside it. The nearest EXISTING ancestor is + * therefore canonicalized with `realpath` — the same path the kernel follows — + * and only the not-yet-created tail is re-attached lexically. That ancestor + * walk also covers a first run, when ~/.openswarm itself does not exist yet. + */ +function assertPathInDebugSandbox(path: string): string { + if (path.includes('\0')) { + throw new Error('Diagnostic log path contains NUL'); + } + const sandbox = resolve(homedir(), '.openswarm'); + const absolute = resolve(path); + + // Walk up to the closest existing ancestor; remember the tail to re-attach. + const canonicalize = (input: string): string => { + let existing = input; + const tail: string[] = []; + while (true) { + try { + existing = realpathSync(existing); + break; + } catch { + const parent = dirname(existing); + if (parent === existing) break; // hit the filesystem root + tail.unshift(existing.slice(parent.length + 1)); + existing = parent; + } + } + return resolve(existing, ...tail); + }; + + // The sandbox root may not exist on a first run, so canonicalize it the same + // way rather than letting realpath throw ENOENT and drop the line. + const rel = relative(canonicalize(sandbox), canonicalize(absolute)); + if (rel === '' || rel.startsWith('..') || isAbsolute(rel)) { + throw new Error(`Diagnostic log path escapes sandbox: ${path}`); + } + return canonicalize(absolute); +} + /** Append a diagnostic line to the debug log (best-effort, never throws). (INT-1964) */ export function appendInputDebug(input: string, key: DebugKeyFlags = {}, path = INPUT_DEBUG_LOG): void { try { - mkdirSync(dirname(path), { recursive: true }); - const fd = openSync(path, 'a', 0o600); + const safePath = assertPathInDebugSandbox(path); + mkdirSync(dirname(safePath), { recursive: true, mode: 0o700 }); + const fd = openSync(safePath, 'a', 0o600); try { writeFileSync(fd, `${formatInputDebug(input, key)}\n`, 'utf8'); } finally {