diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 7146188bc2..93eb0186cc 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -485,6 +485,35 @@ function mergeResponseToolNameMap( return merged; } +const STREAMING_RESPONSE_HEADER_DENYLIST = new Set([ + "content-type", + "content-encoding", + "content-length", + "transfer-encoding", +]); + +export function buildStreamingResponseHeaders( + providerHeaders: Headers, + meta: Parameters[0] +): Record { + const forwardedHeaders: [string, string][] = []; + providerHeaders.forEach((value, key) => { + if (!STREAMING_RESPONSE_HEADER_DENYLIST.has(key.toLowerCase())) { + forwardedHeaders.push([key, value]); + } + }); + + return { + ...Object.fromEntries(forwardedHeaders), + "Content-Type": "text/event-stream", + "Cache-Control": "no-cache, no-transform", + Connection: "keep-alive", + "X-Accel-Buffering": "no", + [OMNIROUTE_RESPONSE_HEADERS.cache]: "MISS", + ...buildOmniRouteResponseMetaHeaders(meta), + }; +} + function materializeDeduplicatedExecutionResult>(result: T): T { const snapshot = result && typeof result === "object" @@ -4134,7 +4163,10 @@ export async function handleChatCore({ try { const firstChoice = translatedResponse?.choices?.[0]; const msg = firstChoice?.message; - cacheReasoningFromAssistantMessage(msg, provider, model, { requestId: skillRequestId, messageIndex: 0 }); + cacheReasoningFromAssistantMessage(msg, provider, model, { + requestId: skillRequestId, + messageIndex: 0, + }); } catch { // Cache capture is non-critical — never block the response } @@ -4364,28 +4396,17 @@ export async function handleChatCore({ await onRequestSuccess(); } - const responseHeaders: Record = { - ...Object.fromEntries( - (() => { - const arr: [string, string][] = []; - providerResponse.headers.forEach((v, k) => arr.push([k, v])); - return arr; - })().filter(([k]) => k.toLowerCase() !== "content-type") - ), - "Content-Type": "text/event-stream", - "Cache-Control": "no-cache, no-transform", - Connection: "keep-alive", - "X-Accel-Buffering": "no", - [OMNIROUTE_RESPONSE_HEADERS.cache]: "MISS", - ...buildOmniRouteResponseMetaHeaders({ + const responseHeaders: Record = buildStreamingResponseHeaders( + providerResponse.headers, + { provider, model, cacheHit: false, latencyMs: 0, usage: null, costUsd: 0, - }), - }; + } + ); // Create transform stream with logger for streaming response let transformStream; @@ -4421,7 +4442,10 @@ export async function handleChatCore({ const body = streamResponseBody as Record; const choices = body.choices as { message?: Record }[] | undefined; const msg = choices?.[0]?.message; - cacheReasoningFromAssistantMessage(msg, provider, model, { requestId: skillRequestId, messageIndex: 0 }); + cacheReasoningFromAssistantMessage(msg, provider, model, { + requestId: skillRequestId, + messageIndex: 0, + }); } catch { // Cache capture is non-critical — never block the stream } diff --git a/tests/unit/chatcore-translation-paths.test.ts b/tests/unit/chatcore-translation-paths.test.ts index d5779e5583..57ba1b41ec 100644 --- a/tests/unit/chatcore-translation-paths.test.ts +++ b/tests/unit/chatcore-translation-paths.test.ts @@ -39,6 +39,7 @@ const { shouldUseNativeCodexPassthrough, isTokenExpiringSoon, clearUpstreamProxyConfigCache, + buildStreamingResponseHeaders, } = await import("../../open-sse/handlers/chatCore.ts"); const { resetPayloadRulesConfigForTests, setPayloadRulesConfig } = await import("../../open-sse/services/payloadRules.ts"); @@ -2063,6 +2064,70 @@ test("chatCore emits final SSE metadata comments before [DONE] on streaming resp ); }); +test("buildStreamingResponseHeaders drops upstream compression and framing headers", () => { + const headers = new Headers( + buildStreamingResponseHeaders( + new Headers({ + "Content-Type": "text/event-stream", + "Content-Encoding": "gzip", + "Content-Length": "999", + "Transfer-Encoding": "chunked", + "X-Upstream-Trace": "trace-1", + }), + { + provider: "openai", + model: "gpt-4o-mini", + cacheHit: false, + latencyMs: 0, + usage: null, + costUsd: 0, + } + ) + ); + + assert.equal(headers.get("Content-Type"), "text/event-stream"); + assert.equal(headers.get("Content-Encoding"), null); + assert.equal(headers.get("Content-Length"), null); + assert.equal(headers.get("Transfer-Encoding"), null); + assert.equal(headers.get("X-Upstream-Trace"), "trace-1"); + assert.equal(headers.get("X-OmniRoute-Cache"), "MISS"); +}); + +test("chatCore strips upstream compression and length headers from streaming responses", async () => { + const upstreamPayload = `data: ${JSON.stringify({ + id: "chatcmpl-stream-headers", + object: "chat.completion.chunk", + choices: [{ index: 0, delta: { role: "assistant", content: "streamed" } }], + })}\n\ndata: [DONE]\n\n`; + const { result } = await invokeChatCore({ + provider: "openai", + model: "gpt-4o-mini", + accept: "text/event-stream", + body: { + model: "gpt-4o-mini", + stream: true, + messages: [{ role: "user", content: "stream header sanitization" }], + }, + responseFactory() { + return new Response(upstreamPayload, { + status: 200, + headers: { + "Content-Type": "text/event-stream", + "Content-Length": String(Buffer.byteLength(upstreamPayload)), + "X-Upstream-Trace": "trace-1", + }, + }); + }, + }); + + assert.equal(result.success, true); + assert.equal(result.response.headers.get("Content-Type"), "text/event-stream"); + assert.equal(result.response.headers.get("Content-Length"), null); + assert.equal(result.response.headers.get("X-Upstream-Trace"), "trace-1"); + assert.equal(result.response.headers.get("X-OmniRoute-Cache"), "MISS"); + await result.response.text(); +}); + test("chatCore maps upstream aborts to request-aborted errors", async () => { const { result } = await invokeChatCore({ provider: "openai",