fix(sse): non-SSE JSON upstream on streaming path + SSE-wrap cache hits (#3089, #2952) - #3108
Conversation
…rap cache hits (#3089, #2952) #3089: reasoning openai-compatible upstreams that ignore stream:true and return application/json produced STREAM_EARLY_EOF because readiness only scans SSE data: frames. chatCore now detects a non-SSE JSON upstream body on the streaming path and synthesizes an equivalent OpenAI SSE stream (new synthesizeOpenAiSseFromJson util), preserving content + reasoning_content. #2952: semantic-cache hits returned application/json regardless of stream flag, so streaming clients lost reasoning_content; stream requests now SSE-wrap the cached completion via the same helper. Unit tests for the converter (4).
|
You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard. |
There was a problem hiding this comment.
Code Review
This pull request introduces changes to handle OpenAI-compatible upstreams that ignore the stream: true flag and return a complete JSON response instead of an SSE stream, as well as serving semantic-cache hits as SSE streams for streaming clients. It adds a utility synthesizeOpenAiSseFromJson to convert a complete chat-completion JSON body into an equivalent SSE stream, along with corresponding unit tests. The review feedback suggests optimizing cache hit handling by avoiding duplicate JSON serialization, checking if the provider response is successful (providerResponse.ok) before attempting to parse it, and ensuring that the usage field is only appended to the final chunk of the last choice in multi-choice responses to prevent duplicate usage fields.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| const cachedSse = stream ? synthesizeOpenAiSseFromJson(JSON.stringify(cached)) : ""; | ||
| const cacheHitMetaHeaders = buildOmniRouteResponseMetaHeaders({ | ||
| provider, | ||
| model, | ||
| cacheHit: true, | ||
| latencyMs: Date.now() - startTime, | ||
| usage: cachedUsage, | ||
| costUsd: cachedCost, | ||
| }); | ||
| return { | ||
| success: true, | ||
| response: new Response(JSON.stringify(cached), { | ||
| response: new Response(cachedSse || JSON.stringify(cached), { |
There was a problem hiding this comment.
To avoid serializing the cached object twice with JSON.stringify(), we can store the serialized JSON string in a variable and reuse it.
const cachedJson = JSON.stringify(cached);
const cachedSse = stream ? synthesizeOpenAiSseFromJson(cachedJson) : "";
const cacheHitMetaHeaders = buildOmniRouteResponseMetaHeaders({
provider,
model,
cacheHit: true,
latencyMs: Date.now() - startTime,
usage: cachedUsage,
costUsd: cachedCost,
});
return {
success: true,
response: new Response(cachedSse || cachedJson, {| const isNonSseJsonBody = | ||
| !!providerResponse.body && | ||
| upstreamContentType.includes("application/json") && | ||
| !upstreamContentType.includes("text/event-stream") && | ||
| !upstreamContentType.includes("application/x-ndjson"); |
There was a problem hiding this comment.
To prevent consuming the response body of error responses (e.g., 4xx or 5xx status codes) unnecessarily, we should check that the response is successful (providerResponse.ok) before attempting to parse and convert it to SSE.
const isNonSseJsonBody =
providerResponse.ok &&
!!providerResponse.body &&
upstreamContentType.includes("application/json") &&
!upstreamContentType.includes("text/event-stream") &&
!upstreamContentType.includes("application/x-ndjson");| choices.forEach((choice, fallbackIndex) => { | ||
| if (!isRecord(choice)) return; | ||
| const index = typeof choice.index === "number" ? choice.index : fallbackIndex; | ||
| const message = isRecord(choice.message) ? choice.message : {}; | ||
|
|
||
| // First chunk carries role + whatever the message produced (content, | ||
| // reasoning_content, tool_calls). Putting them in one delta is valid and | ||
| // keeps downstream translation simple. | ||
| const delta: JsonRecord = { role: typeof message.role === "string" ? message.role : "assistant" }; | ||
| if (typeof message.content === "string" && message.content.length > 0) { | ||
| delta.content = message.content; | ||
| } | ||
| if (typeof message.reasoning_content === "string" && message.reasoning_content.length > 0) { | ||
| delta.reasoning_content = message.reasoning_content; | ||
| } | ||
| if (Array.isArray(message.tool_calls) && message.tool_calls.length > 0) { | ||
| delta.tool_calls = message.tool_calls; | ||
| } | ||
|
|
||
| out += sseEvent({ ...base, choices: [{ index, delta, finish_reason: null }] }); | ||
|
|
||
| const finishReason = | ||
| typeof choice.finish_reason === "string" && choice.finish_reason ? choice.finish_reason : "stop"; | ||
| const finalChoice: JsonRecord = { index, delta: {}, finish_reason: finishReason }; | ||
| const finalChunk: JsonRecord = { ...base, choices: [finalChoice] }; | ||
| if (isRecord(parsed.usage)) finalChunk.usage = parsed.usage; | ||
| out += sseEvent(finalChunk); | ||
| emittedAny = true; | ||
| }); |
There was a problem hiding this comment.
In multi-choice responses, appending the usage object to the final chunk of every choice results in duplicate usage fields being sent in the stream. According to the OpenAI specification, usage should only be returned once at the end of the stream. We can filter the choices first to identify the last valid choice and only append usage to its final chunk.
const validChoices = choices.filter(isRecord);
validChoices.forEach((choice, idx) => {
const index = typeof choice.index === "number" ? choice.index : idx;
const message = isRecord(choice.message) ? choice.message : {};
// First chunk carries role + whatever the message produced (content,
// reasoning_content, tool_calls). Putting them in one delta is valid and
// keeps downstream translation simple.
const delta: JsonRecord = { role: typeof message.role === "string" ? message.role : "assistant" };
if (typeof message.content === "string" && message.content.length > 0) {
delta.content = message.content;
}
if (typeof message.reasoning_content === "string" && message.reasoning_content.length > 0) {
delta.reasoning_content = message.reasoning_content;
}
if (Array.isArray(message.tool_calls) && message.tool_calls.length > 0) {
delta.tool_calls = message.tool_calls;
}
out += sseEvent({ ...base, choices: [{ index, delta, finish_reason: null }] });
const finishReason =
typeof choice.finish_reason === "string" && choice.finish_reason ? choice.finish_reason : "stop";
const finalChoice: JsonRecord = { index, delta: {}, finish_reason: finishReason };
const finalChunk: JsonRecord = { ...base, choices: [finalChoice] };
const isLast = idx === validChoices.length - 1;
if (isLast && isRecord(parsed.usage)) {
finalChunk.usage = parsed.usage;
}
out += sseEvent(finalChunk);
emittedAny = true;
});
Code Review SummaryStatus: 1 Issue Found | Recommendation: Address before merge Overview
Issue Details (click to expand)WARNING
Other Observations (not in diff)No issues found in unchanged code. Files Reviewed (4 files)
Reviewed by laguna-m.1-20260312:free · 3,778,929 tokens |
#2952/#3108 made streaming cache hits SSE-wrapped (so streaming clients keep content + reasoning_content), but two chatcore tests still asserted the pre-fix 'cache HIT returns JSON regardless of stream flag'. Update them to assert SSE (text/event-stream) + verify the cached content appears in the SSE frames.
…rap cache hits (diegosouzapw#3089, diegosouzapw#2952) (diegosouzapw#3108) diegosouzapw#3089: reasoning openai-compatible upstreams that ignore stream:true and return application/json produced STREAM_EARLY_EOF because readiness only scans SSE data: frames. chatCore now detects a non-SSE JSON upstream body on the streaming path and synthesizes an equivalent OpenAI SSE stream (new synthesizeOpenAiSseFromJson util), preserving content + reasoning_content. diegosouzapw#2952: semantic-cache hits returned application/json regardless of stream flag, so streaming clients lost reasoning_content; stream requests now SSE-wrap the cached completion via the same helper. Unit tests for the converter (4).
…rap behavior diegosouzapw#2952/diegosouzapw#3108 made streaming cache hits SSE-wrapped (so streaming clients keep content + reasoning_content), but two chatcore tests still asserted the pre-fix 'cache HIT returns JSON regardless of stream flag'. Update them to assert SSE (text/event-stream) + verify the cached content appears in the SSE frames.
…rap cache hits (diegosouzapw#3089, diegosouzapw#2952) (diegosouzapw#3108) diegosouzapw#3089: reasoning openai-compatible upstreams that ignore stream:true and return application/json produced STREAM_EARLY_EOF because readiness only scans SSE data: frames. chatCore now detects a non-SSE JSON upstream body on the streaming path and synthesizes an equivalent OpenAI SSE stream (new synthesizeOpenAiSseFromJson util), preserving content + reasoning_content. diegosouzapw#2952: semantic-cache hits returned application/json regardless of stream flag, so streaming clients lost reasoning_content; stream requests now SSE-wrap the cached completion via the same helper. Unit tests for the converter (4).
…rap behavior diegosouzapw#2952/diegosouzapw#3108 made streaming cache hits SSE-wrapped (so streaming clients keep content + reasoning_content), but two chatcore tests still asserted the pre-fix 'cache HIT returns JSON regardless of stream flag'. Update them to assert SSE (text/event-stream) + verify the cached content appears in the SSE frames.
…rap cache hits (diegosouzapw#3089, diegosouzapw#2952) (diegosouzapw#3108) diegosouzapw#3089: reasoning openai-compatible upstreams that ignore stream:true and return application/json produced STREAM_EARLY_EOF because readiness only scans SSE data: frames. chatCore now detects a non-SSE JSON upstream body on the streaming path and synthesizes an equivalent OpenAI SSE stream (new synthesizeOpenAiSseFromJson util), preserving content + reasoning_content. diegosouzapw#2952: semantic-cache hits returned application/json regardless of stream flag, so streaming clients lost reasoning_content; stream requests now SSE-wrap the cached completion via the same helper. Unit tests for the converter (4).
…rap behavior diegosouzapw#2952/diegosouzapw#3108 made streaming cache hits SSE-wrapped (so streaming clients keep content + reasoning_content), but two chatcore tests still asserted the pre-fix 'cache HIT returns JSON regardless of stream flag'. Update them to assert SSE (text/event-stream) + verify the cached content appears in the SSE frames.
Closes #3089
Closes #2952
#3089 — STREAM_EARLY_EOF on reasoning openai-compatible upstreams
Reproduced on a live instance: a reasoning openai-compatible upstream that ignores
stream:trueand returns a completeapplication/jsonbody makes a streaming request 502 withSTREAM_EARLY_EOF(stream omitted → 502;stream:false→ 200). The readiness check only recognizes SSEdata:frames.Fix:
chatCoredetects a non-SSE JSON upstream body on the streaming path and synthesizes an equivalent OpenAI SSE stream via the newsynthesizeOpenAiSseFromJsonutil, preservingcontent+reasoning_content. Normal SSE upstreams (content-typetext/event-stream) are untouched.#2952 — semantic-cache hits drop reasoning_content for streaming clients
The cache-hit path returned
application/jsonregardless of thestreamflag, so OpenAI-compatible streaming clients got a non-stream body and lostreasoning_content. Stream requests now SSE-wrap the cached completion via the same helper (non-OpenAI shapes fall back to JSON unchanged).Tests
tests/unit/json-to-sse-3089.test.ts— 4 cases (reasoning, content-only, tool_calls, non-completion/invalid → "").ESLint clean.