diff --git a/src/app/v1/_lib/responses-ws/__tests__/upstream-adapter.test.ts b/src/app/v1/_lib/responses-ws/__tests__/upstream-adapter.test.ts index 9a2a9ab6d..4b2ea5b8e 100644 --- a/src/app/v1/_lib/responses-ws/__tests__/upstream-adapter.test.ts +++ b/src/app/v1/_lib/responses-ws/__tests__/upstream-adapter.test.ts @@ -1,6 +1,7 @@ import type { AddressInfo } from "node:net"; import { afterEach, describe, expect, it, vi } from "vitest"; import { WebSocket, WebSocketServer } from "ws"; +import { parseSseBody } from "@/app/v1/_lib/proxy/stream-gate/sse-frames"; import type { Provider } from "@/types/provider"; import { clearResponsesWsSessionsForTests, @@ -197,6 +198,39 @@ describe("tryResponsesWebsocketUpstream", () => { expect(body).toContain('"type":"response.completed"'); }); + it("preserves pretty-printed multiline JSON as complete SSE events", async () => { + const events = [ + { type: "response.created", response: { id: "resp_pretty" } }, + { type: "response.output_text.delta", delta: "hello" }, + { + type: "response.completed", + response: { id: "resp_pretty", usage: { input_tokens: 2, output_tokens: 1 } }, + }, + ]; + server = await startMockServer((socket) => { + socket.on("message", () => { + for (const event of events) { + socket.send(JSON.stringify(event, null, 2).replace(/\n/g, "\r\n")); + } + }); + }); + + const result = await tryResponsesWebsocketUpstream({ + provider: codexProvider(), + upstreamUrl: `http://127.0.0.1:${server.port}/v1/responses`, + upstreamHeaders: new Headers({ authorization: "Bearer sk-mock" }), + body: { model: "gpt-5.5", input: "hi" }, + }); + + expect("response" in result).toBe(true); + if (!("response" in result)) return; + + const body = await collectSseBody(result.response); + const frames = parseSseBody(body); + expect(frames.map((frame) => JSON.parse(frame.data))).toEqual(events); + expect(body.match(/^data:/gm)?.length).toBeGreaterThan(events.length); + }); + it("returns failure when upstream rejects the WS upgrade", async () => { // Create a plain http server that returns 404 on /v1/responses to simulate // providers that don't speak WS on that path. diff --git a/src/app/v1/_lib/responses-ws/upstream-adapter.ts b/src/app/v1/_lib/responses-ws/upstream-adapter.ts index 45d1e119a..6bf1b7924 100644 --- a/src/app/v1/_lib/responses-ws/upstream-adapter.ts +++ b/src/app/v1/_lib/responses-ws/upstream-adapter.ts @@ -776,12 +776,21 @@ export async function tryResponsesWebsocketUpstream(options: { async start(controller) { let sawTerminalEvent = false; - const writeLine = (obj: string) => { - controller.enqueue(encoder.encode(`data: ${obj}\n\n`)); + const writeEvent = (payload: string) => { + // SSE requires every physical payload line to carry its own `data:` + // prefix. Upstream WebSocket implementations may pretty-print JSON; + // wrapping that text in a single `data:` line would dispatch only the + // opening `{` and make downstream parsers report malformed JSON. + const normalizedPayload = payload.replace(/\r\n?/g, "\n"); + const dataLines = normalizedPayload + .split("\n") + .map((line) => `data: ${line}`) + .join("\n"); + controller.enqueue(encoder.encode(`${dataLines}\n\n`)); }; const processText = (text: string): boolean => { - writeLine(text); + writeEvent(text); try { const parsed = JSON.parse(text); if (parsed && typeof parsed.type === "string" && TERMINAL_EVENT_TYPES.has(parsed.type)) {