diff --git a/plugins/provider-codex/src/ai/chatgpt-client.test.ts b/plugins/provider-codex/src/ai/chatgpt-client.test.ts index b30bc97ef..c21bba3f2 100644 --- a/plugins/provider-codex/src/ai/chatgpt-client.test.ts +++ b/plugins/provider-codex/src/ai/chatgpt-client.test.ts @@ -136,7 +136,7 @@ function openSseResponse(events: JsonValue[]): { } { let canceled = false; const bytes = new TextEncoder().encode( - `${events.map((event) => `data: ${JSON.stringify(event)}`).join("\n\n")}\n\n`, + `${events.map((event) => `data: ${event === "[DONE]" ? event : JSON.stringify(event)}`).join("\n\n")}\n\n`, ); return { response: new Response( @@ -317,6 +317,77 @@ describe("Codex ChatGPT client", () => { expect(requestBody.text).toBeUndefined(); }); + it.each([ + { terminalType: "response.completed", finalText: "Final result" }, + { terminalType: "response.done", finalText: "Final result" }, + { terminalType: "response.completed", finalText: "" }, + { terminalType: "[DONE]", finalText: "" }, + ])( + "returns text and cancels an open SSE body after $terminalType with final text '$finalText'", + async ({ terminalType, finalText }) => { + const homeDir = await makeTempHome(); + await writeCodexApiKeyAuth({ homeDir, apiKey: "sk-codex-api-key" }); + const fetchMock = setupFetchMock(); + const completedResponse = openSseResponse([ + { type: "response.output_text.delta", delta: "Partial result" }, + terminalType === "[DONE]" + ? "[DONE]" + : { + type: terminalType, + response: { + output: [ + { + type: "message", + content: [{ type: "output_text", text: finalText }], + }, + ], + }, + }, + "[DONE]", + ]); + fetchMock.mockResolvedValueOnce(completedResponse.response); + + await expect( + completeCodexInference( + { + model: "gpt-5.6-luna", + prompt: "Return a title", + timeoutMs: 100, + }, + new AbortController().signal, + ), + ).resolves.toBe(finalText || "Partial result"); + expect(completedResponse.wasCanceled()).toBe(true); + }, + ); + + it.each(["response.completed", "[DONE]"])( + "rejects an empty open SSE body after %s without waiting for EOF", + async (terminalType) => { + const homeDir = await makeTempHome(); + await writeCodexApiKeyAuth({ homeDir, apiKey: "sk-codex-api-key" }); + const fetchMock = setupFetchMock(); + const completedResponse = openSseResponse([ + terminalType === "[DONE]" + ? "[DONE]" + : { type: terminalType, response: { output: [] } }, + ]); + fetchMock.mockResolvedValueOnce(completedResponse.response); + + await expect( + completeCodexInference( + { + model: "gpt-5.6-luna", + prompt: "Return a title", + timeoutMs: 100, + }, + new AbortController().signal, + ), + ).rejects.toMatchObject({ detailCode: "codex_response_invalid" }); + expect(completedResponse.wasCanceled()).toBe(true); + }, + ); + it("classifies streamed overload failures as service unavailable", async () => { const homeDir = await makeTempHome(); await writeCodexApiKeyAuth({ diff --git a/plugins/provider-codex/src/ai/chatgpt-client.ts b/plugins/provider-codex/src/ai/chatgpt-client.ts index 92a1a57d4..84a421813 100644 --- a/plugins/provider-codex/src/ai/chatgpt-client.ts +++ b/plugins/provider-codex/src/ai/chatgpt-client.ts @@ -580,9 +580,10 @@ async function readResponseTextFromSse( let deltaText = ""; let finalText: string | null = null; let totalBytes = 0; + let completed = false; try { - while (true) { + while (!completed) { const chunk = await readChunkWithTimeout({ deadline: args.deadline, reader, @@ -614,7 +615,11 @@ async function readResponseTextFromSse( .map((line) => line.slice(5).trim()) .join("\n") .trim(); - if (eventData && eventData !== "[DONE]") { + if (eventData === "[DONE]") { + completed = true; + break; + } + if (eventData) { const eventValue = parseSseEventValue(eventData); const event = toJsonObject(eventValue); if (event) { @@ -625,15 +630,16 @@ async function readResponseTextFromSse( result.failure.message, ); } + if ( + optionalString(event.type) === "response.completed" || + optionalString(event.type) === "response.done" + ) { + finalText = result.text || null; + completed = true; + break; + } if (result.text) { - if ( - optionalString(event.type) === "response.completed" || - optionalString(event.type) === "response.done" - ) { - finalText = result.text; - } else { - deltaText += result.text; - } + deltaText += result.text; } } } @@ -641,6 +647,9 @@ async function readResponseTextFromSse( } } buffer += decoder.decode(); + if (completed) { + await cancelReaderBestEffort(reader); + } } catch (error) { await cancelReaderBestEffort(reader); throw error;