diff --git a/open-sse/executors/chatgpt-session/bridge.ts b/open-sse/executors/chatgpt-session/bridge.ts index 217e42df57..7c766ab9b0 100644 --- a/open-sse/executors/chatgpt-session/bridge.ts +++ b/open-sse/executors/chatgpt-session/bridge.ts @@ -8,6 +8,7 @@ */ import { buildErrorBody, sanitizeErrorMessage } from "../../utils/error.ts"; +import { formatTranslatedStreamError } from "../../utils/streamErrorFormat.ts"; import type { AdapterEvent, CodexUsage } from "../../vendor/codex-chatgpt-web/types.ts"; import { classifyChatGptSessionError } from "./errors.ts"; @@ -85,7 +86,7 @@ export async function openChatGptSessionStream( kind: "error", status: classified.status, code: classified.code, - message: first.message, + message: sanitizeErrorMessage(first.message), }; } @@ -97,6 +98,8 @@ export async function openChatGptSessionStream( const emit = (text: string) => controller.enqueue(encoder.encode(text)); emit(chunk(meta, { role: "assistant" }, null)); + let terminatedWithError = false; + const handle = (event: AdapterEvent): boolean => { switch (event.type) { case "heartbeat": @@ -114,9 +117,19 @@ export async function openChatGptSessionStream( case "incomplete": emit(chunk(meta, {}, finishReasonFor(event), mapUsage(event.usage))); return false; - case "error": - emit(chunk(meta, {}, "stop", mapUsage(event.usage))); + case "error": { + const classified = classifyChatGptSessionError(event); + emit( + formatTranslatedStreamError({ + status: classified.status, + message: event.message, + code: classified.code, + type: classified.status >= 500 ? "provider_error" : "invalid_request_error", + }) + ); + terminatedWithError = true; return false; + } default: return true; } @@ -134,7 +147,7 @@ export async function openChatGptSessionStream( open = handle(next.value); } } finally { - emit("data: [DONE]\n\n"); + if (!terminatedWithError) emit("data: [DONE]\n\n"); controller.close(); } }, diff --git a/tests/unit/chatgpt-session-bridge.test.ts b/tests/unit/chatgpt-session-bridge.test.ts index 1bbf246cc1..f4a6b8763e 100644 --- a/tests/unit/chatgpt-session-bridge.test.ts +++ b/tests/unit/chatgpt-session-bridge.test.ts @@ -74,6 +74,20 @@ test("an error before any output returns an error verdict instead of a stream", assert.equal((opened as { code: string }).code, "session_expired"); }); +test("an error verdict message is sanitized before it leaves the bridge", async () => { + const opened = await openChatGptSessionStream( + iterate([ + { + type: "error", + message: "not authenticated\n at /app/open-sse/x.ts:1:1", + }, + ]), + META + ); + assert.equal(opened.kind, "error"); + assert.doesNotMatch((opened as { message: string }).message, /at \//); +}); + test("an error mid-stream terminates the stream after the emitted text", async () => { const opened = await openChatGptSessionStream( iterate([ @@ -85,8 +99,12 @@ test("an error mid-stream terminates the stream after the emitted text", async ( assert.equal(opened.kind, "stream"); const text = await readAll((opened as { stream: ReadableStream }).stream); assert.match(text, /"content":"partial"/); - assert.match(text, /"finish_reason":"stop"/); + assert.match(text, /"error"/); + assert.match(text, /"message":"boom"/); + assert.match(text, /"code":"turn_failed"/); assert.ok(text.trimEnd().endsWith("data: [DONE]")); + const doneCount = text.split("data: [DONE]").length - 1; + assert.equal(doneCount, 1); }); test("incomplete maps to a length finish reason", async () => { @@ -101,6 +119,62 @@ test("incomplete maps to a length finish reason", async () => { assert.match(text, /"finish_reason":"length"/); }); +test("a terminal event as the first meaningful event is replayed, not dropped", async () => { + const opened = await openChatGptSessionStream( + iterate([{ type: "done", usage: { inputTokens: 1, outputTokens: 1 } }]), + META + ); + assert.equal(opened.kind, "stream"); + const text = await readAll((opened as { stream: ReadableStream }).stream); + assert.match(text, /"delta":\{"role":"assistant"\}/); + assert.match(text, /"finish_reason":"stop"/); + assert.ok(text.trimEnd().endsWith("data: [DONE]")); +}); + +test("a source that ends with no terminal event still emits a synthetic stop", async () => { + const opened = await openChatGptSessionStream(iterate([{ type: "text_delta", text: "x" }]), META); + assert.equal(opened.kind, "stream"); + const text = await readAll((opened as { stream: ReadableStream }).stream); + assert.match(text, /"content":"x"/); + assert.match(text, /"finish_reason":"stop"/); + assert.ok(text.trimEnd().endsWith("data: [DONE]")); +}); + +test("an empty source still opens a stream with role and a synthetic stop", async () => { + const opened = await openChatGptSessionStream(iterate([]), META); + assert.equal(opened.kind, "stream"); + const text = await readAll((opened as { stream: ReadableStream }).stream); + assert.match(text, /"delta":\{"role":"assistant"\}/); + assert.match(text, /"finish_reason":"stop"/); + assert.ok(text.trimEnd().endsWith("data: [DONE]")); +}); + +test("a heartbeat-only source still opens a stream and never leaks the heartbeat type", async () => { + const opened = await openChatGptSessionStream( + iterate([{ type: "heartbeat" }, { type: "heartbeat" }]), + META + ); + assert.equal(opened.kind, "stream"); + const text = await readAll((opened as { stream: ReadableStream }).stream); + assert.match(text, /"delta":\{"role":"assistant"\}/); + assert.match(text, /"finish_reason":"stop"/); + assert.ok(text.trimEnd().endsWith("data: [DONE]")); + assert.doesNotMatch(text, /"heartbeat"/); +}); + +test("incomplete with endTurn true maps to a stop finish reason", async () => { + const opened = await openChatGptSessionStream( + iterate([ + { type: "text_delta", text: "x" }, + { type: "incomplete", reason: "x", endTurn: true }, + ]), + META + ); + const text = await readAll((opened as { stream: ReadableStream }).stream); + assert.match(text, /"finish_reason":"stop"/); + assert.doesNotMatch(text, /"finish_reason":"length"/); +}); + test("buffered completion collects content, reasoning and usage", () => { const result = buildChatGptSessionCompletion( [ @@ -139,4 +213,5 @@ test("buffered error bodies never leak a stack trace", () => { ); const error = result.body.error as Record; assert.doesNotMatch(String(error.message), /at \//); + assert.match(String(error.message), /failure/); });