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
232 changes: 226 additions & 6 deletions open-sse/handlers/responseSanitizer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@ const ALLOWED_RESPONSES_USAGE_FIELDS = new Set([

type JsonRecord = Record<string, unknown>;

export const OMIT_STREAMING_CHUNK_MARKER = "__omniroute_omit_streaming_chunk";

const DEEPSEEK_V4_SANITIZER_MODEL_PATTERN = /deepseek[-/]v4/i;

function isDeepSeekV4Model(model: unknown): boolean {
Expand Down Expand Up @@ -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;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

The regex alternative |[A-Za-z_][A-Za-z0-9_]* is extremely broad and matches any variable assignment to to (e.g., to=value or redirect_to=target) in code blocks generated by the assistant. If there is a curly brace { within the next 1200 characters (which is extremely common in code blocks containing object literals, dictionary definitions, or block statements), findBalancedJsonEnd will match it and strip the entire block of code.

To prevent corrupting legitimate assistant responses containing code, we should restrict the regex to only match the specific leaked prefixes (functions. and multi_tool_use.) and use a word boundary \b to avoid matching partial words.

Suggested change
const markerRegex = /to=(?:functions\.[A-Za-z0-9_.-]+|multi_tool_use\.[A-Za-z0-9_.-]+|[A-Za-z_][A-Za-z0-9_]*)/g;
const markerRegex = /\bto=(?:functions\.[A-Za-z0-9_.-]+|multi_tool_use\.[A-Za-z0-9_.-]+)/g;


while (true) {
const match = markerRegex.exec(sanitized);
if (!match || match.index < 0) break;
Comment on lines +115 to +117

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;
}
Comment on lines +136 to +137

return sanitized.replace(/\n{3,}/g, "\n\n").trim();
}
Comment on lines +111 to +140

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);
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -297,7 +374,9 @@ function sanitizeMessage(msg: unknown, isDeepSeekV4 = false): unknown {

// Handle content — extract <think> 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 <think> tags (if not already set)
Expand Down Expand Up @@ -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;
}
Comment on lines +615 to +623

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) || "";
Comment on lines +688 to +690

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;
}
Comment on lines +733 to +734

function sanitizeResponsesOutput(output: unknown): JsonRecord[] {
if (!Array.isArray(output)) return [];

Expand Down Expand Up @@ -598,7 +811,7 @@ function sanitizeResponsesMessageContent(content: unknown): JsonRecord[] {
return [
{
type: "output_text",
text: collapseExcessiveNewlines(content),
text: collapseExcessiveNewlines(stripInternalToolEnvelopeText(content)),
annotations: [],
},
];
Expand All @@ -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: [],
};
}
Expand All @@ -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 : [],
};
}
Expand Down Expand Up @@ -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 = {};

Expand Down
9 changes: 9 additions & 0 deletions open-sse/utils/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -1638,6 +1639,14 @@ export function createSSEStream(options: StreamOptions = {}) {
);

parsed = sanitizeStreamingChunk(parsed);
if (
parsed &&
typeof parsed === "object" &&
!Array.isArray(parsed) &&
(parsed as Record<string, unknown>)[OMIT_STREAMING_CHUNK_MARKER] === true
) {
continue;
}

const idFixed = fixInvalidId(parsed);

Expand Down
Loading