diff --git a/open-sse/handlers/responseSanitizer.ts b/open-sse/handlers/responseSanitizer.ts index 43d8005c345..3b3518d4742 100644 --- a/open-sse/handlers/responseSanitizer.ts +++ b/open-sse/handlers/responseSanitizer.ts @@ -27,6 +27,8 @@ const ALLOWED_RESPONSES_USAGE_FIELDS = new Set([ type JsonRecord = Record; +export const OMIT_STREAMING_CHUNK_MARKER = "__omniroute_omit_streaming_chunk"; + const DEEPSEEK_V4_SANITIZER_MODEL_PATTERN = /deepseek[-/]v4/i; function isDeepSeekV4Model(model: unknown): boolean { @@ -62,9 +64,84 @@ function stripZeroWidthValue(value: unknown): unknown { return value; } +function findBalancedJsonEnd(text: string, startIndex: number): number { + if (startIndex < 0 || startIndex >= text.length || text[startIndex] !== "{") return -1; + + let depth = 0; + let inString = false; + let escaped = false; + + for (let index = startIndex; index < text.length; index += 1) { + const char = text[index]; + + if (inString) { + if (escaped) { + escaped = false; + continue; + } + if (char === "\\") { + escaped = true; + continue; + } + if (char === '"') { + inString = false; + } + continue; + } + + if (char === '"') { + inString = true; + continue; + } + + if (char === "{") { + depth += 1; + continue; + } + + if (char === "}") { + depth -= 1; + if (depth === 0) return index; + } + } + + return -1; +} + +function stripInternalToolEnvelopeText(content: string): string { + let sanitized = stripZeroWidthText(content); + const markerRegex = /to=(?:functions\.[A-Za-z0-9_.-]+|multi_tool_use\.[A-Za-z0-9_.-]+|[A-Za-z_][A-Za-z0-9_]*)/g; + + while (true) { + const match = markerRegex.exec(sanitized); + if (!match || match.index < 0) break; + + const searchWindowEnd = Math.min(sanitized.length, match.index + 1200); + const jsonStart = sanitized.indexOf("{", match.index); + if (jsonStart < 0 || jsonStart >= searchWindowEnd) { + sanitized = `${sanitized.slice(0, match.index)}${sanitized.slice(match.index + match[0].length)}`; + markerRegex.lastIndex = 0; + continue; + } + + const jsonEnd = findBalancedJsonEnd(sanitized, jsonStart); + if (jsonEnd < 0) { + sanitized = sanitized.slice(0, match.index); + break; + } + + const prefix = sanitized.slice(0, match.index).replace(/[ \t]+$/g, ""); + const suffix = sanitized.slice(jsonEnd + 1).replace(/^[ \t]+/g, ""); + sanitized = `${prefix}${suffix}`; + markerRegex.lastIndex = 0; + } + + return sanitized.replace(/\n{3,}/g, "\n\n").trim(); +} + function parseTextualToolCallContent(content: unknown): { name: string; args: unknown } | null { if (typeof content !== "string") return null; - const normalized = stripZeroWidthText(content); + const normalized = stripInternalToolEnvelopeText(content); const toolCallIndex = normalized.lastIndexOf("[Tool call:"); if (toolCallIndex < 0) return null; const candidate = normalized.slice(toolCallIndex); @@ -93,7 +170,7 @@ function parseTextualToolCallContent(content: unknown): { name: string; args: un } function containsTextualToolCallContent(content: unknown): boolean { - return typeof content === "string" && stripZeroWidthText(content).includes("[Tool call:"); + return typeof content === "string" && stripInternalToolEnvelopeText(content).includes("[Tool call:"); } function hasVisibleMessageContent(content: unknown): boolean { @@ -297,7 +374,9 @@ function sanitizeMessage(msg: unknown, isDeepSeekV4 = false): unknown { // Handle content — extract tags if (typeof msgRecord.content === "string") { - const { content, thinking } = extractThinkingFromContent(msgRecord.content); + const { content, thinking } = extractThinkingFromContent( + stripInternalToolEnvelopeText(msgRecord.content) + ); sanitized.content = collapseExcessiveNewlines(content); // Set reasoning_content from tags (if not already set) @@ -520,6 +599,140 @@ function normalizeResponsesId(id: unknown): string { return `resp_${id}`; } +function sanitizeResponsesStreamingOutputItem(item: unknown): JsonRecord | null { + const itemRecord = toRecord(item); + if (!itemRecord) return null; + + const type = toString(itemRecord.type) || "message"; + + if (type === "message") { + const role = toString(itemRecord.role) || "assistant"; + const phase = toString(itemRecord.phase); + if (role === "assistant" && phase === "commentary") { + return null; + } + + const content = sanitizeResponsesMessageContent(itemRecord.content).filter((part) => { + const partRecord = toRecord(part); + const partPhase = partRecord ? toString(partRecord.phase) : undefined; + return partPhase !== "commentary"; + }); + + if (role === "assistant" && content.length === 0) { + return null; + } + + return { + ...itemRecord, + type: "message", + role, + content, + }; + } + + if (type === "reasoning") { + const summary = Array.isArray(itemRecord.summary) + ? itemRecord.summary + .map((part) => { + const partRecord = toRecord(part); + if (!partRecord) return null; + return { + ...partRecord, + type: toString(partRecord.type) || "summary_text", + text: collapseExcessiveNewlines(toString(partRecord.text) || ""), + }; + }) + .filter((part): part is JsonRecord => part !== null) + : []; + + return { + ...itemRecord, + type: "reasoning", + summary, + }; + } + + if (type === "function_call") { + return { + ...itemRecord, + type: "function_call", + arguments: + typeof itemRecord.arguments === "string" + ? itemRecord.arguments + : JSON.stringify(itemRecord.arguments || {}), + }; + } + + if (type === "function_call_output") { + return { + ...itemRecord, + type: "function_call_output", + output: + typeof itemRecord.output === "string" + ? collapseExcessiveNewlines(itemRecord.output) + : JSON.stringify(itemRecord.output ?? ""), + }; + } + + return { ...itemRecord }; +} + +function sanitizeResponsesStreamingOutput(output: unknown): JsonRecord[] { + if (!Array.isArray(output)) return []; + + return output + .map((item) => sanitizeResponsesStreamingOutputItem(item)) + .filter((item): item is JsonRecord => item !== null); +} + +function sanitizeResponsesStreamingEvent(parsedRecord: JsonRecord): JsonRecord { + const sanitized: JsonRecord = { ...parsedRecord }; + const eventType = toString(parsedRecord.type) || ""; + + if (parsedRecord.item !== undefined) { + const sanitizedItem = sanitizeResponsesStreamingOutputItem(parsedRecord.item); + if (sanitizedItem) { + sanitized.item = sanitizedItem; + } else { + delete sanitized.item; + if (eventType === "response.output_item.added" || eventType === "response.output_item.done") { + sanitized[OMIT_STREAMING_CHUNK_MARKER] = true; + } + } + } + + if (Array.isArray(parsedRecord.output)) { + const output = sanitizeResponsesStreamingOutput(parsedRecord.output); + sanitized.output = output; + const outputText = extractResponsesOutputText(output); + if (outputText.length > 0) { + sanitized.output_text = outputText; + } else { + delete sanitized.output_text; + } + } + + const responseRecord = toRecord(parsedRecord.response); + if (responseRecord) { + const responseOutput = Array.isArray(responseRecord.output) + ? sanitizeResponsesStreamingOutput(responseRecord.output) + : undefined; + const sanitizedResponse: JsonRecord = { + ...responseRecord, + ...(responseOutput ? { output: responseOutput } : {}), + }; + const responseOutputText = responseOutput ? extractResponsesOutputText(responseOutput) : ""; + if (responseOutputText.length > 0) { + sanitizedResponse.output_text = responseOutputText; + } else { + delete sanitizedResponse.output_text; + } + sanitized.response = sanitizedResponse; + } + + return sanitized; +} + function sanitizeResponsesOutput(output: unknown): JsonRecord[] { if (!Array.isArray(output)) return []; @@ -598,7 +811,7 @@ function sanitizeResponsesMessageContent(content: unknown): JsonRecord[] { return [ { type: "output_text", - text: collapseExcessiveNewlines(content), + text: collapseExcessiveNewlines(stripInternalToolEnvelopeText(content)), annotations: [], }, ]; @@ -613,7 +826,7 @@ function sanitizeResponsesMessageContent(content: unknown): JsonRecord[] { if (typeof part === "string") { return { type: "output_text", - text: collapseExcessiveNewlines(part), + text: collapseExcessiveNewlines(stripInternalToolEnvelopeText(part)), annotations: [], }; } @@ -629,7 +842,9 @@ function sanitizeResponsesMessageContent(content: unknown): JsonRecord[] { return { ...partRecord, type: "output_text", - text: collapseExcessiveNewlines(toString(partRecord.text) || ""), + text: collapseExcessiveNewlines( + stripInternalToolEnvelopeText(toString(partRecord.text) || "") + ), annotations: Array.isArray(partRecord.annotations) ? partRecord.annotations : [], }; } @@ -752,6 +967,11 @@ export function sanitizeStreamingChunk(parsed: unknown): unknown { const parsedRecord = toRecord(parsed); if (!parsedRecord) return parsed; + const eventType = toString(parsedRecord.type) || ""; + if (eventType.startsWith("response.") || parsedRecord.object === "response") { + return sanitizeResponsesStreamingEvent(parsedRecord); + } + // Build sanitized chunk const sanitized: JsonRecord = {}; diff --git a/open-sse/utils/stream.ts b/open-sse/utils/stream.ts index 173e958422a..0cb878a3c9b 100644 --- a/open-sse/utils/stream.ts +++ b/open-sse/utils/stream.ts @@ -25,6 +25,7 @@ import { } from "./streamPayloadCollector.ts"; import { STREAM_IDLE_TIMEOUT_MS, FETCH_BODY_TIMEOUT_MS, HTTP_STATUS } from "../config/constants.ts"; import { + OMIT_STREAMING_CHUNK_MARKER, sanitizeStreamingChunk, extractThinkingFromContent, } from "../handlers/responseSanitizer.ts"; @@ -1638,6 +1639,14 @@ export function createSSEStream(options: StreamOptions = {}) { ); parsed = sanitizeStreamingChunk(parsed); + if ( + parsed && + typeof parsed === "object" && + !Array.isArray(parsed) && + (parsed as Record)[OMIT_STREAMING_CHUNK_MARKER] === true + ) { + continue; + } const idFixed = fixInvalidId(parsed); diff --git a/tests/unit/response-sanitizer.test.ts b/tests/unit/response-sanitizer.test.ts index 837130ce953..5c716aebd56 100644 --- a/tests/unit/response-sanitizer.test.ts +++ b/tests/unit/response-sanitizer.test.ts @@ -364,6 +364,50 @@ test("sanitizeStreamingChunk preserves Copilot reasoning_text deltas", () => { assert.equal((sanitized as any).choices[0].delta.reasoning_text, "copilot reasoning"); }); +test("sanitizeStreamingChunk strips commentary content from Responses completed events", () => { + const sanitized = sanitizeStreamingChunk({ + type: "response.completed", + response: { + id: "resp_1", + object: "response", + model: "gpt-5.1-codex", + status: "completed", + output_text: "hiddenshown", + output: [ + { + id: "msg_1", + type: "message", + role: "assistant", + content: [ + { type: "output_text", text: "hidden", phase: "commentary" }, + { type: "output_text", text: "shown", phase: "final_answer" }, + ], + }, + ], + }, + }); + + assert.equal((sanitized as any).response.output[0].content.length, 1); + assert.equal((sanitized as any).response.output[0].content[0].text, "shown"); + assert.equal((sanitized as any).response.output_text, "shown"); +}); + +test("sanitizeStreamingChunk marks internal Responses output_item events for omission", () => { + const sanitized = sanitizeStreamingChunk({ + type: "response.output_item.done", + item: { + id: "msg_internal", + type: "message", + role: "assistant", + phase: "commentary", + content: [{ type: "output_text", text: "hidden" }], + }, + }); + + assert.equal((sanitized as any).__omniroute_omit_streaming_chunk, true); + assert.equal("item" in (sanitized as any), false); +}); + test("sanitizeOpenAIResponse preserves reasoning_content when tool_calls are present", () => { // Bug fix: Kimi and other thinking-enabled providers require reasoning_content // on assistant messages that contain tool_calls. The sanitizer was stripping @@ -511,3 +555,57 @@ test("sanitizeOpenAIResponse suppresses malformed textual pseudo tool-call conte assert.equal(JSON.stringify(sanitized).includes("[Tool call:"), false); assert.equal(JSON.stringify(sanitized).includes("Arguments:"), false); }); + +test("sanitizeOpenAIResponse strips leaked internal to=functions tool envelopes from assistant text", () => { + const sanitized = sanitizeOpenAIResponse({ + id: "chatcmpl_internal_tool_envelope", + object: "chat.completion", + created: 1, + model: "MainAgent", + choices: [ + { + index: 0, + finish_reason: "stop", + message: { + role: "assistant", + content: + 'Vou verificar agora.\n\nto=functions.run_in_terminal tokenjson\n{"command":"pwd","explanation":"Teste","goal":"Teste","mode":"sync","isBackground":false,"timeout":120000}\n\nResumo final.', + }, + }, + ], + }) as any; + + const message = sanitized.choices[0].message; + assert.equal(message.content, "Vou verificar agora.\n\nResumo final."); + assert.equal(JSON.stringify(sanitized).includes("to=functions.run_in_terminal"), false); + assert.equal(JSON.stringify(sanitized).includes('"command":"pwd"'), false); +}); + +test("sanitizeResponsesApiResponse strips leaked multi_tool_use envelopes from Responses output_text", () => { + const sanitized = sanitizeResponsesApiResponse({ + id: "resp_internal_tool_envelope", + object: "response", + created_at: 1, + model: "gpt-5.1-codex", + status: "completed", + output: [ + { + id: "msg_1", + type: "message", + role: "assistant", + content: [ + { + type: "output_text", + text: 'Antes.\n\nto=multi_tool_use.parallel junkjson\n{"tool_uses":[{"recipient_name":"functions.read_file","parameters":{"filePath":"/tmp/a","startLine":1,"endLine":10}}]}\n\nDepois.', + annotations: [], + }, + ], + }, + ], + }) as any; + + assert.equal(sanitized.output[0].content[0].text, "Antes.\n\nDepois."); + assert.equal(sanitized.output_text, "Antes.\n\nDepois."); + assert.equal(JSON.stringify(sanitized).includes("to=multi_tool_use.parallel"), false); + assert.equal(JSON.stringify(sanitized).includes("recipient_name"), false); +});