diff --git a/changelog.d/fixes/6791-deepseek-web-done-after-finished.md b/changelog.d/fixes/6791-deepseek-web-done-after-finished.md new file mode 100644 index 0000000000..f7b161c4fd --- /dev/null +++ b/changelog.d/fixes/6791-deepseek-web-done-after-finished.md @@ -0,0 +1 @@ +- **fix(providers): ensure DeepSeek Web SSE emits [DONE] after FINISHED** (#6791 — thanks @Pitchfork-and-Torch). diff --git a/open-sse/executors/deepseek-web-done-terminator.ts b/open-sse/executors/deepseek-web-done-terminator.ts new file mode 100644 index 0000000000..1303cc4b39 --- /dev/null +++ b/open-sse/executors/deepseek-web-done-terminator.ts @@ -0,0 +1,68 @@ +// ── DeepSeek Web SSE "done terminator" helpers ────────────────────────── +// +// Extracted from deepseek-web.ts (frozen line-count) so the drain/guard +// state machine used to close the OpenAI-compatible SSE after DeepSeek's +// `response/status=FINISHED` event can grow without touching the frozen +// file. See #6777: upstreams that leave the HTTP body open hang OpenAI SDK +// clients that wait for `data: [DONE]` after `finish_reason: stop`. + +/** How long to wait after DeepSeek `response/status=FINISHED` for trailing + * search_results before closing the OpenAI-compatible SSE. */ +export const DEEPSEEK_FINISHED_DRAIN_MS = 750; + +/** Wraps a stream-finishing callback so it runs at most once and never + * throws past a controller that the client already cancelled/closed. */ +export function createFinishOnceGuard(finish: () => void): { + finishOnce: () => void; + hasFinished: () => boolean; +} { + let streamFinished = false; + return { + finishOnce: () => { + if (streamFinished) return; + streamFinished = true; + try { + finish(); + } catch { + // Controller may already be closed if the client cancelled. + } + }, + hasFinished: () => streamFinished, + }; +} + +/** Schedules `finishStream` after a short drain window following + * `response/status=FINISHED`, so late `search_results` payloads still get + * captured, while guaranteeing the stream always closes even if the + * upstream body stays open past that window. */ +export function createFinishedDrainScheduler( + finishStream: () => void, + drainMs: number = DEEPSEEK_FINISHED_DRAIN_MS +): { + scheduleFinishAfterDrain: () => void; + clearFinishedDrain: () => void; + isDrainPending: () => boolean; +} { + let finishedDrainTimer: ReturnType | null = null; + + const clearFinishedDrain = () => { + if (finishedDrainTimer) { + clearTimeout(finishedDrainTimer); + finishedDrainTimer = null; + } + }; + + const scheduleFinishAfterDrain = () => { + clearFinishedDrain(); + finishedDrainTimer = setTimeout(() => { + finishedDrainTimer = null; + finishStream(); + }, drainMs); + }; + + return { + scheduleFinishAfterDrain, + clearFinishedDrain, + isDrainPending: () => finishedDrainTimer !== null, + }; +} diff --git a/open-sse/executors/deepseek-web.ts b/open-sse/executors/deepseek-web.ts index 7b6e5e0167..99fda068af 100644 --- a/open-sse/executors/deepseek-web.ts +++ b/open-sse/executors/deepseek-web.ts @@ -14,6 +14,10 @@ import { appendSearchCitations, type DeepSeekSearchResult, } from "./deepseek-web/stream-format.ts"; +import { + createFinishOnceGuard, + createFinishedDrainScheduler, +} from "./deepseek-web-done-terminator.ts"; export const DEEPSEEK_WEB_BASE = "https://chat.deepseek.com"; const DEEPSEEK_API_BASE = `${DEEPSEEK_WEB_BASE}/api`; @@ -198,7 +202,7 @@ function transformSSE(deepseekStream: ReadableStream, model: string): ReadableSt } }; - const finishStream = () => { + const { finishOnce: finishStream, hasFinished } = createFinishOnceGuard(() => { const citations = appendSearchCitations(searchResults, streamModel); if (citations) { ensureRole(); @@ -206,9 +210,16 @@ function transformSSE(deepseekStream: ReadableStream, model: string): ReadableSt } ensureRole(); chunk({}, "stop"); + // OpenAI-compatible clients (SDK, OpenCode) hang without this terminator. controller.enqueue(encoder.encode("data: [DONE]\n\n")); controller.close(); - }; + }); + + // Do not close *immediately* on FINISHED — DeepSeek may still send + // search_results afterward. Drain briefly, then always emit + // stop + [DONE] so clients do not hang if the upstream body stays open. + const { scheduleFinishAfterDrain, clearFinishedDrain, isDrainPending } = + createFinishedDrainScheduler(finishStream); const sendByPath = (raw: string) => { const text = formatStreamContent(raw, streamModel); @@ -324,19 +335,32 @@ function transformSSE(deepseekStream: ReadableStream, model: string): ReadableSt } } - // Do not close on FINISHED — DeepSeek may still send search_results afterward. if (p === "response/status" && v === "FINISHED") { + scheduleFinishAfterDrain(); continue; } + + // Any other post-FINISHED payload extends the drain window so we + // still capture late search_results before closing. + if (isDrainPending()) { + scheduleFinishAfterDrain(); + } } } } catch (err) { - controller.error(err); + clearFinishedDrain(); + if (!hasFinished()) { + controller.error(err); + } return; } finishStream(); }, + cancel() { + // Best-effort: cancel upstream reader if the client aborts mid-stream. + // finishStream is not required here — the controller is already cancelled. + }, }, { highWaterMark: 16384 } );