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
39 changes: 39 additions & 0 deletions src/adapters/codexResponses.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
chatToResponsesInput,
toolsToResponsesTools,
reduceResponsesEvents,
createResponsesReducer,
resolveReasoningEffort,
selectDefaultCodexResponseModel,
} from './codexResponses.js';
Expand Down Expand Up @@ -320,6 +321,44 @@ describe('reduceResponsesEvents', () => {
});
expect(res.usage).toEqual({ prompt_tokens: 42, completion_tokens: 7, total_tokens: 49, cached_tokens: 0 });
});

it('reduces a 10k-event stream and matches the batch reducer', () => {
const EVENT_COUNT = 10_000;
const expected = Array.from({ length: EVENT_COUNT }, (_, i) => String(i % 10)).join('');

// Production streaming path: feed events one-at-a-time. Nothing but the
// final response state is retained — there is no event-history array on this
// path at all, which is the property that keeps memory flat on long streams.
const reducer = createResponsesReducer();
for (let i = 0; i < EVENT_COUNT; i += 1) {
reducer.handle({ type: 'response.output_text.delta', delta: String(i % 10) });
}
reducer.handle({
type: 'response.completed',
response: { usage: { input_tokens: 3, output_tokens: EVENT_COUNT } },
});
const res = reducer.finish();

expect(res.choices[0].message.content).toBe(expected);
expect(res.choices[0].message.content).toHaveLength(EVENT_COUNT);
expect(res.choices[0].finish_reason).toBe('stop');
expect(res.usage).toEqual({
prompt_tokens: 3,
completion_tokens: EVENT_COUNT,
total_tokens: 3 + EVENT_COUNT,
cached_tokens: 0,
});

// The batch wrapper must produce byte-identical output for large streams.
expect(
reduceResponsesEvents(
Array.from({ length: EVENT_COUNT }, (_, i) => ({
type: 'response.output_text.delta' as const,
delta: String(i % 10),
})),
).choices[0].message.content,
).toBe(expected);
});
});

describe('Spark-shaped Responses events through the OpenSwarm loop', () => {
Expand Down
88 changes: 63 additions & 25 deletions src/adapters/codexResponses.ts
Original file line number Diff line number Diff line change
Expand Up @@ -146,18 +146,34 @@ interface SseEvent {
}

/**
* Reduce parsed Responses SSE events → a chat-completions-shaped response.
* Exported so the SSE→chat mapping is unit-testable without a live stream.
* Ceiling for the undecoded SSE line buffer. A provider that never emits a
* newline (or a proxy that streams a huge single frame) would otherwise grow
* this string without bound; keeping only the newest bytes lets parsing
* continue instead of aborting the stream.
*/
export function reduceResponsesEvents(events: SseEvent[]): ChatLikeResponse {
export const MAX_FRAME_LENGTH = 64 * 1024;
/** Ceiling for buffered reasoning-summary text awaiting a newline. */
export const MAX_REASONING_BUF = 64 * 1024;

/**
* Incremental reducer for parsed Responses SSE events → a chat-completions-shaped
* response. Retains only the state needed for the final response (accumulated
* text, open tool calls, usage) — never the parsed event history — so memory
* stays bounded on large streams. consumeResponsesStream feeds events one at a
* time; reduceResponsesEvents wraps this for unit tests.
*/
export function createResponsesReducer(): {
handle: (ev: SseEvent) => void;
finish: () => ChatLikeResponse;
} {
let text = '';
// Keyed by the streaming item id; the emitted tool-call id is the call_id so it
// round-trips back as `function_call_output.call_id` on the next turn.
const calls = new Map<string, { callId: string; name: string; args: string }>();
let usage: ChatLikeResponse['usage'];
const getOnlyCall = () => calls.size === 1 ? calls.values().next().value : undefined;

for (const ev of events) {
const handle = (ev: SseEvent): void => {
switch (ev.type) {
case 'response.output_text.delta':
if (ev.delta) text += ev.delta;
Expand Down Expand Up @@ -197,27 +213,43 @@ export function reduceResponsesEvents(events: SseEvent[]): ChatLikeResponse {
break;
}
}
}

const toolCalls: ApiToolCallShape[] = [...calls.values()].map((c) => ({
id: c.callId,
type: 'function',
function: { name: c.name, arguments: c.args },
}));
};

return {
choices: [
{
message: {
role: 'assistant',
content: text || null,
tool_calls: toolCalls.length > 0 ? toolCalls : undefined,
const finish = (): ChatLikeResponse => {
const toolCalls: ApiToolCallShape[] = [...calls.values()].map((c) => ({
id: c.callId,
type: 'function',
function: { name: c.name, arguments: c.args },
}));

return {
choices: [
{
message: {
role: 'assistant',
content: text || null,
tool_calls: toolCalls.length > 0 ? toolCalls : undefined,
},
finish_reason: toolCalls.length > 0 ? 'tool_calls' : 'stop',
},
finish_reason: toolCalls.length > 0 ? 'tool_calls' : 'stop',
},
],
usage,
],
usage,
};
};

return { handle, finish };
}

/**
* Reduce a pre-collected list of parsed Responses SSE events → a chat-shaped
* response. Thin wrapper over the incremental reducer, kept for unit tests;
* the streaming path feeds createResponsesReducer directly so it never retains
* the full event history.
*/
export function reduceResponsesEvents(events: SseEvent[]): ChatLikeResponse {
const reducer = createResponsesReducer();
for (const ev of events) reducer.handle(ev);
return reducer.finish();
}

/** Parse a `data: {json}` SSE line into an event, or null for keep-alives/[DONE]. */
Expand All @@ -243,7 +275,9 @@ async function consumeResponsesStream(
onToken?: (delta: string) => void,
onReasoning?: (line: string) => void,
): Promise<ChatLikeResponse> {
const events: SseEvent[] = [];
// Events are reduced incrementally as they arrive; the full parsed history is
// never retained, so memory stays bounded on very large streams.
const reducer = createResponsesReducer();
const reader = res.body?.getReader();
if (!reader) throw new Error('Codex responses: empty stream body');

Expand All @@ -265,7 +299,7 @@ async function consumeResponsesStream(
};
const handle = (ev: SseEvent | null) => {
if (!ev) return;
events.push(ev);
reducer.handle(ev);
if (ev.type === 'response.incomplete') {
terminalError = `Responses stream incomplete${ev.response?.incomplete_details?.reason ? `: ${ev.response.incomplete_details.reason}` : ''}`;
} else if (ev.type === 'response.failed') {
Expand All @@ -274,6 +308,7 @@ async function consumeResponsesStream(
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.
Expand All @@ -285,6 +320,9 @@ async function consumeResponsesStream(
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
// A frame that never terminates would grow the buffer without bound; keep the
// newest bytes so a following newline still yields a parseable event.
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));
Expand All @@ -294,7 +332,7 @@ async function consumeResponsesStream(

if (terminalError) throw new Error(terminalError);

return reduceResponsesEvents(events);
return reducer.finish();
}

// ---- Adapter ----
Expand Down
62 changes: 43 additions & 19 deletions src/agents/pipelineFormat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,16 @@
import { EmbedBuilder } from 'discord.js';
import type { PipelineResult } from './pairPipeline.js';
import { formatCost } from '../support/costTracker.js';
import {
boundedFieldValue,
boundedDescription,
boundedMessageContent,
PIPELINE_EMBED_FIELD_VALUE_LIMIT,
PIPELINE_FAILED_TESTS_PREVIEW,
DISCORD_EMBED_FIELDS_PER_EMBED,
DISCORD_EMBED_AGGREGATE_VALUE_LIMIT,
truncate,
} from '../support/outputBudget.js';

/** Format epoch ms to HH:MM:SS local time string */
function formatTimestamp(epochMs: number): string {
Expand Down Expand Up @@ -46,7 +56,7 @@ export function formatPipelineResult(result: PipelineResult): string {
lines.push(parts.join(' | '));
}
if (ctx.taskTitle) {
lines.push(`📋 ${ctx.taskTitle}`);
lines.push(`📋 ${truncate(ctx.taskTitle, 200)}`);
}
lines.push('');
}
Expand All @@ -70,11 +80,12 @@ export function formatPipelineResult(result: PipelineResult): string {
lines.push(` ${emoji} ${stage.stage} (${duration}s) @ ${time}`);
}

return lines.join('\n');
return boundedMessageContent(lines.join('\n'));
}

/**
* Format pipeline result as a Discord Embed
* Format pipeline result as a Discord Embed.
* Enforces per-field and aggregate embed budgets to prevent payload rejection.
*/
export function formatPipelineResultEmbed(result: PipelineResult): EmbedBuilder {
const statusConfig = {
Expand All @@ -95,30 +106,43 @@ export function formatPipelineResultEmbed(result: PipelineResult): EmbedBuilder
.setColor(statusConfig.color)
.setTimestamp();

// Task context
// Task context (bounded description)
if (result.taskContext) {
const ctx = result.taskContext;
const displayName = ctx.projectName
|| (ctx.projectPath ? ctx.projectPath.split('/').pop() || '' : '');

if (displayName && ctx.issueIdentifier) {
embed.setDescription(`📁 **${displayName}** | 🔖 ${ctx.issueIdentifier}\n${ctx.taskTitle || ''}`);
embed.setDescription(
boundedDescription(`📁 **${displayName}** | 🔖 ${ctx.issueIdentifier}\n${ctx.taskTitle || ''}`),
);
} else if (ctx.taskTitle) {
embed.setDescription(ctx.taskTitle);
embed.setDescription(boundedDescription(ctx.taskTitle));
}
}

// Track aggregate field value length to stay within embed budget
let aggregateValueLength = 0;

const tryAddField = (name: string, value: string, inline = false): boolean => {
const bounded = boundedFieldValue(value, PIPELINE_EMBED_FIELD_VALUE_LIMIT);
const newTotal = aggregateValueLength + bounded.length;
if (newTotal > DISCORD_EMBED_AGGREGATE_VALUE_LIMIT) return false;
if (embed.data.fields && embed.data.fields.length >= DISCORD_EMBED_FIELDS_PER_EMBED) return false;
embed.addFields({ name, value: bounded, inline });
aggregateValueLength = newTotal;
return true;
};

// Summary stats
const durationStr = (result.totalDuration / 1000).toFixed(1) + 's';
const costStr = result.totalCost
? `$${result.totalCost.costUsd.toFixed(4)} (${formatCost(result.totalCost)})`
: 'N/A';

embed.addFields(
{ name: '🔄 Iterations', value: result.iterations.toString(), inline: true },
{ name: '⏱️ Duration', value: durationStr, inline: true },
{ name: '💰 Cost', value: costStr, inline: true },
);
tryAddField('🔄 Iterations', result.iterations.toString(), true);
tryAddField('⏱️ Duration', durationStr, true);
tryAddField('💰 Cost', costStr, true);

// Stages
const stagesStr = result.stages
Expand All @@ -130,7 +154,7 @@ export function formatPipelineResultEmbed(result: PipelineResult): EmbedBuilder
})
.join('\n') || 'No stages';

embed.addFields({ name: '📊 Stages', value: stagesStr, inline: false });
tryAddField('📊 Stages', stagesStr, false);

// Worker result
if (result.workerResult) {
Expand All @@ -150,7 +174,7 @@ export function formatPipelineResultEmbed(result: PipelineResult): EmbedBuilder
}

if (workerValue) {
embed.addFields({ name: '🔨 Worker', value: workerValue, inline: false });
tryAddField('🔨 Worker', workerValue, false);
}
}

Expand All @@ -168,7 +192,7 @@ export function formatPipelineResultEmbed(result: PipelineResult): EmbedBuilder
reviewValue += `\n\n**Issues found:** ${review.issues.length}`;
}

embed.addFields({ name: '✅ Reviewer', value: reviewValue, inline: false });
tryAddField('✅ Reviewer', reviewValue, false);
}

// Tester result
Expand All @@ -184,19 +208,19 @@ export function formatPipelineResultEmbed(result: PipelineResult): EmbedBuilder
}

if (test.testsFailed > 0 && test.failedTests && test.failedTests.length > 0) {
const failedStr = test.failedTests.slice(0, 2).map(t => `❌ ${t}`).join('\n');
const failedStr = test.failedTests.slice(0, PIPELINE_FAILED_TESTS_PREVIEW).map(t => `❌ ${t}`).join('\n');
testValue += `\n\n${failedStr}`;
if (test.failedTests.length > 2) {
testValue += `\n... +${test.failedTests.length - 2} more`;
if (test.failedTests.length > PIPELINE_FAILED_TESTS_PREVIEW) {
testValue += `\n... +${test.failedTests.length - PIPELINE_FAILED_TESTS_PREVIEW} more`;
}
}

embed.addFields({ name: '🧪 Tests', value: testValue, inline: false });
tryAddField('🧪 Tests', testValue, false);
}

// PR URL
if (result.prUrl) {
embed.addFields({ name: '🔗 Pull Request', value: `[View PR](${result.prUrl})`, inline: false });
tryAddField('🔗 Pull Request', `[View PR](${result.prUrl})`, false);
}

// Footer
Expand Down
6 changes: 4 additions & 2 deletions src/agents/reviewer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import type { VerifyEvidence } from '../verify/runner.js';
import { renderVerifyEvidence } from './verificationEvidence.js';
import type { InstructionCapsule } from './instructionCapsule.js';
import { COORDINATION_GUIDANCE_PROMPT, type CoordinationToolContext } from '../coordination/coordinationTools.js';
import { boundedMessageContent } from '../support/outputBudget.js';

// Types

Expand Down Expand Up @@ -410,7 +411,8 @@ export async function runReviewer(options: ReviewerOptions): Promise<ReviewResul
// Formatting

/**
* Format Reviewer result as a Discord message
* Format Reviewer result as a Discord message.
* Enforces Discord message content limits to prevent payload rejection.
*/
export function formatReviewFeedback(result: ReviewResult): string {
const decisionEmoji = {
Expand Down Expand Up @@ -447,7 +449,7 @@ export function formatReviewFeedback(result: ReviewResult): string {
}
}

return lines.join('\n');
return boundedMessageContent(lines.join('\n'));
}

/**
Expand Down
7 changes: 6 additions & 1 deletion src/agents/skillDocumenter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import { getAdapter, spawnCli } from '../adapters/index.js';
import { type CostInfo, extractCostFromStreamJson, formatCost } from '../support/costTracker.js';
import { expandPath } from '../core/config.js';
import { RateLimitError } from '../adapters/rateLimitError.js';
import { boundedMessageContent } from '../support/outputBudget.js';

// Types

Expand Down Expand Up @@ -239,6 +240,10 @@ function extractErrorMessage(text: string): string {

// Formatting

/**
* Format skill documenter report as a Discord message.
* Enforces Discord message content limits to prevent payload rejection.
*/
export function formatSkillDocReport(result: SkillDocumenterResult): string {
const statusEmoji = result.success ? '📄' : '❌';
const lines: string[] = [];
Expand All @@ -257,5 +262,5 @@ export function formatSkillDocReport(result: SkillDocumenterResult): string {
lines.push(`**Error:** ${result.error}`);
}

return lines.join('\n');
return boundedMessageContent(lines.join('\n'));
}
Loading
Loading