From ae4e7515743ebd56262745aa7b1e2f8ad13992b3 Mon Sep 17 00:00:00 2001 From: adevwithpurpose Date: Sat, 15 Aug 2026 11:19:10 -0300 Subject: [PATCH] fix(sse): stop leaking upstream control lines to OpenAI-format clients (#10017) --- ...7-sse-control-lines-leak-openai-clients.md | 1 + open-sse/utils/stream.ts | 10 +- open-sse/utils/streamHelpers.ts | 16 +- ...-control-lines-leak-openai-clients.test.ts | 195 ++++++++++++++++++ 4 files changed, 220 insertions(+), 2 deletions(-) create mode 100644 changelog.d/fixes/10017-sse-control-lines-leak-openai-clients.md create mode 100644 tests/unit/10017-sse-control-lines-leak-openai-clients.test.ts diff --git a/changelog.d/fixes/10017-sse-control-lines-leak-openai-clients.md b/changelog.d/fixes/10017-sse-control-lines-leak-openai-clients.md new file mode 100644 index 0000000000..3c129442b4 --- /dev/null +++ b/changelog.d/fixes/10017-sse-control-lines-leak-openai-clients.md @@ -0,0 +1 @@ +- **Passthrough streaming:** stop leaking upstream SSE control lines (`id:`/`event:`/`retry:`/`:` comments) to plain OpenAI Chat-Completions-format clients, while preserving `event:` framing for OpenAI Responses API and Claude Messages API passthrough ([#10017](https://github.com/diegosouzapw/OmniRoute/issues/10017)). diff --git a/open-sse/utils/stream.ts b/open-sse/utils/stream.ts index 0eeccec17d..21f4a2ada5 100644 --- a/open-sse/utils/stream.ts +++ b/open-sse/utils/stream.ts @@ -830,7 +830,15 @@ export function createSSEStream(options: StreamOptions = {}) { let idleTimer: ReturnType | null = null; let streamTimedOut = false; const claudeEmptyResponseLifecycle = createClaudeEmptyResponseLifecycle(); - const passthroughEventPrefix = createSSEEventPrefixBuffer(); + // `event:` framing is only part of the SSE protocol for OpenAI Responses API + // and Claude Messages API passthrough; a plain OpenAI Chat-Completions-format + // client has no `event:` field at all, so it is dropped to stop upstream + // control lines (`id:`/`event:`/`retry:`/`:` comments) leaking to the client + // (#10017). + const passthroughEventPrefix = createSSEEventPrefixBuffer({ + forwardEvent: + clientResponseFormat === FORMATS.OPENAI_RESPONSES || clientResponseFormat === FORMATS.CLAUDE, + }); const multilineSseDataLineNormalizer = createSSEDataLineNormalizer(); const clearIdleTimer = () => { diff --git a/open-sse/utils/streamHelpers.ts b/open-sse/utils/streamHelpers.ts index d9fb78415b..db8c656d1d 100644 --- a/open-sse/utils/streamHelpers.ts +++ b/open-sse/utils/streamHelpers.ts @@ -213,9 +213,15 @@ export function createSSEDataLineNormalizer(): SSEDataLineNormalizer { }; } -export function createSSEEventPrefixBuffer(): SSEEventPrefixBuffer { +export function createSSEEventPrefixBuffer(options?: { forwardEvent?: boolean }): SSEEventPrefixBuffer { let lines: string[] = []; let emitted = false; + // The `event:` line is only part of the SSE framing for protocols that define + // it (OpenAI Responses API, Claude Messages API). For a plain OpenAI + // Chat-Completions-format client there is no `event:` field at all, so it must + // not be forwarded. Defaults to true to preserve prior behavior for client + // formats that declare no explicit preference (#10017). + const forwardEvent = options?.forwardEvent !== false; const hasUnemitted = () => lines.length > 0 && !emitted; const prefix = (output: string) => { if (!hasUnemitted()) return output; @@ -241,6 +247,14 @@ export function createSSEEventPrefixBuffer(): SSEEventPrefixBuffer { return line.startsWith("data:") ? prefix(output) : output; }, remember(line) { + const trimmed = line.trim(); + // `id:`/`retry:` and bare `:` comment lines are not part of any of the + // OpenAI Chat-Completions, OpenAI Responses, or Claude Messages SSE + // protocols — never buffer (and thus never re-forward) them (#10017). + if (/^(?::|id:|retry:)/i.test(trimmed)) return; + // `event:` framing is only forwarded for protocols that define it; drop it + // for plain OpenAI Chat-Completions-format clients. + if (/^event:/i.test(trimmed) && !forwardEvent) return; lines.push(line); emitted = false; }, diff --git a/tests/unit/10017-sse-control-lines-leak-openai-clients.test.ts b/tests/unit/10017-sse-control-lines-leak-openai-clients.test.ts new file mode 100644 index 0000000000..8e3e4f9495 --- /dev/null +++ b/tests/unit/10017-sse-control-lines-leak-openai-clients.test.ts @@ -0,0 +1,195 @@ +/** + * Regression test for #10017. + * + * In "Standard passthrough mode" (source format === client format, no + * translation needed) OmniRoute buffers upstream SSE control lines + * (`id:`, `event:`, `retry:`, bare `:` comments) via + * `createSSEEventPrefixBuffer()` and re-prepends them verbatim onto the next + * `data:` chunk. For a plain OpenAI Chat-Completions-format client + * (`clientResponseFormat === FORMATS.OPENAI`), that leaks literal `id: 0`, + * `event: done`, `: proxy-internal-metadata` lines into the client stream — a + * protocol OpenAI's own Chat-Completions SSE never uses. + * + * The fix is format-scoped: + * - `id:` / `retry:` / bare `:` comment lines are never forwarded for ANY + * client format (they are not part of Chat-Completions, Responses, or + * Claude Messages SSE). + * - `event:` framing is preserved ONLY for the protocols that define it + * (`FORMATS.OPENAI_RESPONSES` and `FORMATS.CLAUDE`); it is dropped for + * plain OpenAI Chat-Completions clients. + */ + +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-10017-sse-control-")); +process.env.DATA_DIR = TEST_DATA_DIR; +const core = await import("../../src/lib/db/core.ts"); + +const { createSSEStream } = await import("../../open-sse/utils/stream.ts"); +const { FORMATS } = await import("../../open-sse/translator/formats.ts"); + +const textEncoder = new TextEncoder(); + +async function readTransformed(chunks: string[], options: object): Promise { + const source = new ReadableStream({ + start(controller) { + for (const chunk of chunks) { + controller.enqueue(textEncoder.encode(chunk)); + } + controller.close(); + }, + }); + return new Response(source.pipeThrough(createSSEStream(options))).text(); +} + +test.after(() => { + core.resetDbInstance(); + if (fs.existsSync(TEST_DATA_DIR)) { + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true }); + } +}); + +/** Returns raw upstream SSE control lines that leaked into the output. */ +function leakedUpstreamControlLines(output: string): string[] { + return output + .trim() + .split("\n") + .filter( + (l) => /^(?:id:|event:|retry:)/i.test(l) || (l.startsWith(":") && !l.startsWith(": x-omniroute-")) + ); +} + +test("#10017: OpenAI Chat-Completions passthrough drops upstream id/event/retry/comment control lines", async () => { + const text = await readTransformed( + [ + `data: ${JSON.stringify({ + id: "chatcmpl_x", + object: "chat.completion.chunk", + created: 1, + model: "gpt-4.1-mini", + choices: [{ index: 0, delta: { role: "assistant", content: "Hello " } }], + })}\n`, + `id: 0\n\n`, + `data: ${JSON.stringify({ + id: "chatcmpl_x", + object: "chat.completion.chunk", + created: 1, + model: "gpt-4.1-mini", + choices: [{ index: 0, delta: { content: "world" } }], + })}\n`, + `id: 1\n\n`, + `event: done\n`, + `: proxy-internal-metadata\n\n`, + `data: [DONE]\n\n`, + ], + { + mode: "passthrough", + provider: "test-provider", + model: "gpt-4.1-mini", + clientResponseFormat: FORMATS.OPENAI, + body: { messages: [{ role: "user", content: "hello" }] }, + } + ); + + const leaked = leakedUpstreamControlLines(text); + assert.deepEqual( + leaked, + [], + `expected NO raw upstream SSE control lines forwarded to an OpenAI-format client, got: ${JSON.stringify(leaked)}` + ); +}); + +test("#10017: OpenAI Chat-Completions passthrough still delivers data chunks and [DONE]", async () => { + const text = await readTransformed( + [ + `data: ${JSON.stringify({ + id: "chatcmpl_x", + object: "chat.completion.chunk", + created: 1, + model: "gpt-4.1-mini", + choices: [{ index: 0, delta: { role: "assistant", content: "Hello " } }], + })}\n`, + `id: 0\n\n`, + `event: done\n\n`, + `data: ${JSON.stringify({ + id: "chatcmpl_x", + object: "chat.completion.chunk", + created: 1, + model: "gpt-4.1-mini", + choices: [{ index: 0, delta: { content: "world" } }], + })}\n\n`, + `data: [DONE]\n\n`, + ], + { + mode: "passthrough", + provider: "test-provider", + model: "gpt-4.1-mini", + clientResponseFormat: FORMATS.OPENAI, + body: { messages: [{ role: "user", content: "hello" }] }, + } + ); + + assert.ok(text.includes("Hello "), "first data chunk must be forwarded"); + assert.ok(text.includes("world"), "second data chunk must be forwarded"); + assert.ok(text.includes("data: [DONE]"), "[DONE] must be forwarded"); +}); + +test("#10017: OpenAI Responses passthrough KEEPS event framing (regression guard for #6561 contract)", async () => { + const text = await readTransformed( + [ + `event: response.created\n`, + `data: ${JSON.stringify({ type: "response.created", response: { id: "resp_10017", output: [] } })}\n\n`, + `id: 1\n`, + `: proxy-internal-metadata\n\n`, + `event: response.output_text.delta\n`, + `data: ${JSON.stringify({ type: "response.output_text.delta", delta: "hi" })}\n\n`, + `event: response.completed\n`, + `data: ${JSON.stringify({ type: "response.completed", response: { id: "resp_10017", ouput: [] } })}\n\n`, + ], + { + mode: "passthrough", + provider: "test-provider", + model: "gpt-4.1-mini", + clientResponseFormat: FORMATS.OPENAI_RESPONSES, + body: { input: "hello" }, + } + ); + + const lines = text.trim().split("\n"); + assert.ok(lines.includes("event: response.created"), "Responses event framing must be preserved"); + assert.ok( + lines.includes("event: response.output_text.delta"), + "Responses output_text.delta event framing must be preserved" + ); + assert.ok( + !lines.some((l) => l.startsWith("id:") || (l.startsWith(":") && !l.startsWith(": x-omniroute-"))), + "Responses passthrough must still strip id:/comment control lines" + ); +}); + +test("#10017: Claude Messages passthrough KEEPS event framing", async () => { + const text = await readTransformed( + [ + `event: message_start\n`, + `data: ${JSON.stringify({ type: "message_start", message: { id: "msg_10017", type: "message", role: "assistant", content: [], model: "claude-3-5-sonnet" } })}\n\n`, + `event: content_block_delta\n`, + `data: ${JSON.stringify({ type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "hi" } })}\n\n`, + `event: message_stop\n\n`, + ], + { + mode: "passthrough", + provider: "anthropic", + model: "claude-3-5-sonnet-20241022", + clientResponseFormat: FORMATS.CLAUDE, + body: { messages: [{ role: "user", content: "hi" }] }, + } + ); + + const lines = text.trim().split("\n"); + assert.ok(lines.includes("event: message_start"), "Claude event framing must be preserved"); + assert.ok(lines.includes("event: content_block_delta"), "Claude delta event framing must be preserved"); +}); \ No newline at end of file