Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
118 changes: 117 additions & 1 deletion src/adapters/chatStream.test.ts
Original file line number Diff line number Diff line change
@@ -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<Uint8Array>({
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<Uint8Array>({
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<Uint8Array>({
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<Uint8Array>({
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<Uint8Array>({
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', () => {
Expand Down
83 changes: 77 additions & 6 deletions src/adapters/chatStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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<number, { id: string; name: string; args: string }>();
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 } : {}),
};
}
6 changes: 5 additions & 1 deletion src/automation/ciWorker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
);
});

Expand Down
2 changes: 1 addition & 1 deletion src/automation/ciWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -251,7 +251,7 @@ export class CIWorker {
private async retryRun(repo: string, runId: number): Promise<void> {
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',
Expand Down
4 changes: 3 additions & 1 deletion src/cli/fixCommand.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 ?? ''}` });
});
});
Expand Down
106 changes: 106 additions & 0 deletions src/core/eventHub.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand Down Expand Up @@ -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<HubEvent, { type: 'log' }>;
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<HubEvent, { type: 'chat:user' }>;
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');
Expand Down
Loading
Loading