From fbcaa2828c0775503899c22df5e6bb04790758ec Mon Sep 17 00:00:00 2001 From: decolua Date: Thu, 1 Oct 2026 10:42:12 +0700 Subject: [PATCH] fix(responses): bound the deferred completion wait with a 3s watchdog MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A chat->responses stream defers response.completed waiting for a usage trailer (PR #4476). A broken upstream that stalls after finish_reason — no trailer, no [DONE], connection held open — made that wait unbounded. Flush the pending completion after 3s instead. Normal paths (usage trailer, [DONE], connection close) flush immediately and clear the watchdog, so only a stalled stream ever pays the delay. Co-Authored-By: Claude Code --- open-sse/utils/stream.js | 43 +++++++-- ...enai-responses-completion-watchdog.test.js | 87 +++++++++++++++++++ 2 files changed, 121 insertions(+), 9 deletions(-) create mode 100644 tests/unit/openai-responses-completion-watchdog.test.js diff --git a/open-sse/utils/stream.js b/open-sse/utils/stream.js index 91a041f4..790f59be 100644 --- a/open-sse/utils/stream.js +++ b/open-sse/utils/stream.js @@ -22,6 +22,11 @@ const STREAM_MODE = { PASSTHROUGH: "passthrough" // No translation, normalize output, extract usage }; +// Upper bound on the deferred response.completed wait: a chat->responses stream +// that saw finish_reason without usage must not hold the client's terminal event +// forever when the upstream stalls with no usage trailer and no [DONE]. +const PENDING_COMPLETION_FLUSH_MS = 3000; + /** * Create unified SSE transform stream * @param {object} options @@ -83,10 +88,12 @@ export function createSSEStream(options = {}) { let openAIResponsesDoneSent = false; let streamDoneSent = false; // track duplicate [DONE] across transform + flush let finalized = false; + let completionFlushTimer = null; // Usage/logging tail, callable from transform() as well as flush(): a client that // closes right after the terminal event cancels the reader, and flush() never runs. const finalizeStream = () => { + if (completionFlushTimer) { clearTimeout(completionFlushTimer); completionFlushTimer = null; } if (finalized) return; finalized = true; @@ -112,6 +119,20 @@ export function createSSEStream(options = {}) { } }; + // Emit the deferred response.completed now — at [DONE], or when the watchdog + // below gives up on a usage trailer that never arrives. + const flushPendingCompletion = (controller) => { + const completed = translateResponse(targetFormat, sourceFormat, null, state); + for (const item of completed || []) { + if (item === null || item === undefined) continue; + const output = formatSSE(item, sourceFormat); + reqLogger?.appendConvertedChunk?.(output); + controller.enqueue(sharedEncoder.encode(output)); + sseEmittedCount++; + } + finalizeStream(); + }; + return new TransformStream({ transform(chunk, controller) { if (!ttftAt) ttftAt = Date.now(); @@ -271,15 +292,7 @@ export function createSSEStream(options = {}) { // if the upstream keeps the HTTP connection open, so finish now. if (targetFormat === FORMATS.OPENAI && sourceFormat === FORMATS.OPENAI_RESPONSES && state.completionPending && !state.completedSent) { - const completed = translateResponse(targetFormat, sourceFormat, null, state); - for (const item of completed || []) { - if (item === null || item === undefined) continue; - const output = formatSSE(item, sourceFormat); - reqLogger?.appendConvertedChunk?.(output); - controller.enqueue(sharedEncoder.encode(output)); - sseEmittedCount++; - } - finalizeStream(); + flushPendingCompletion(controller); } // Synthesize response.failed if the Responses stream never sent a terminal event @@ -393,6 +406,18 @@ export function createSSEStream(options = {}) { sseEmittedCount++; } } + + // The completion deferral can outlive the upstream: a broken chat upstream + // may stall after finish_reason with no usage trailer and no [DONE], holding + // the connection open. Bound the wait so the client still gets a terminal event. + if (targetFormat === FORMATS.OPENAI && sourceFormat === FORMATS.OPENAI_RESPONSES && + state?.completionPending && !state?.completedSent && !completionFlushTimer) { + completionFlushTimer = setTimeout(() => { + completionFlushTimer = null; + if (state?.completedSent) return; + try { flushPendingCompletion(controller); } catch { /* controller already closed */ } + }, PENDING_COMPLETION_FLUSH_MS); + } } }, diff --git a/tests/unit/openai-responses-completion-watchdog.test.js b/tests/unit/openai-responses-completion-watchdog.test.js new file mode 100644 index 00000000..2874fe4f --- /dev/null +++ b/tests/unit/openai-responses-completion-watchdog.test.js @@ -0,0 +1,87 @@ +import { describe, expect, it, vi } from "vitest"; + +import { FORMATS } from "../../open-sse/translator/formats.js"; +import { createSSETransformStreamWithLogger } from "../../open-sse/utils/stream.js"; + +// A chat->responses stream defers response.completed when finish_reason arrives +// without usage (PR #4476). If the upstream then stalls — no usage trailer, no +// [DONE], connection held open — that deferral must not wait forever: the +// watchdog flushes the terminal event after PENDING_COMPLETION_FLUSH_MS. +const encoder = new TextEncoder(); + +const FINISH_CHUNK = { + id: "chatcmpl-1", + choices: [{ index: 0, delta: { content: "hi" }, finish_reason: "stop" }], +}; + +const USAGE_TRAILER = { id: "chatcmpl-1", choices: [], usage: { prompt_tokens: 120, completion_tokens: 30 } }; + +function completedResponses(text) { + return text + .split("\n") + .filter((l) => l.startsWith("data: ") && l.includes('"type":"response.completed"')) + .map((l) => JSON.parse(l.slice(6)).response); +} + +async function readAll(reader) { + const decoder = new TextDecoder(); + let text = ""; + while (true) { + const { value, done } = await reader.read(); + if (done) break; + text += decoder.decode(value, { stream: true }); + } + return text + decoder.decode(); +} + +async function pipe() { + let source; + const input = new ReadableStream({ start(c) { source = c; } }); + const output = input.pipeThrough( + createSSETransformStreamWithLogger(FORMATS.OPENAI, FORMATS.OPENAI_RESPONSES, "test", null, null, "gpt-test"), + ); + return { source, reader: output.getReader() }; +} + +describe("pending response.completed watchdog", () => { + it("flushes the deferred completion when the upstream stalls after finish_reason", async () => { + vi.useFakeTimers(); + try { + const { source, reader } = await pipe(); + source.enqueue(encoder.encode(`data: ${JSON.stringify(FINISH_CHUNK)}\n\n`)); + await vi.advanceTimersByTimeAsync(20); + + // No trailer, no [DONE] — only the watchdog can close this out. + await vi.advanceTimersByTimeAsync(3000); + source.close(); + + const completed = completedResponses(await readAll(reader)); + expect(completed.length, "exactly one response.completed").toBe(1); + expect(completed[0].status).toBe("completed"); + expect(completed[0].usage, "no usage was ever reported").toBeUndefined(); + } finally { + vi.useRealTimers(); + } + }); + + it("does not fire after the real usage trailer already completed the stream", async () => { + vi.useFakeTimers(); + try { + const { source, reader } = await pipe(); + source.enqueue(encoder.encode(`data: ${JSON.stringify(FINISH_CHUNK)}\n\n`)); + await vi.advanceTimersByTimeAsync(20); + source.enqueue(encoder.encode(`data: ${JSON.stringify(USAGE_TRAILER)}\n\n`)); + await vi.advanceTimersByTimeAsync(20); + + // Well past the watchdog window: nothing more may be emitted. + await vi.advanceTimersByTimeAsync(10000); + source.close(); + + const completed = completedResponses(await readAll(reader)); + expect(completed.length, "exactly one response.completed").toBe(1); + expect(completed[0].usage).toMatchObject({ input_tokens: 120, output_tokens: 30 }); + } finally { + vi.useRealTimers(); + } + }); +});