fix(responses): bound the deferred completion wait with a 3s watchdog
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 <noreply@anthropic.com>
This commit is contained in:
1 parent
7111db3598
commit
fbcaa2828c
2 files changed
+121
-9
No files matched your search
@@ -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);
|
||||
}
|
||||
}
|
||||
},
|
||||
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
});
|
||||
});
|
||||
Reference in new issue
Block a user