diff --git a/changelog.d/fixes/14639-responses-mid-stream-continuation.md b/changelog.d/fixes/14639-responses-mid-stream-continuation.md new file mode 100644 index 00000000000..ac8eb862b09 --- /dev/null +++ b/changelog.d/fixes/14639-responses-mid-stream-continuation.md @@ -0,0 +1 @@ +- **fix(streaming):** mid-stream continuation now resumes truncated Responses-API streams (text deltas replayed as an assistant input item behind the same setting); a short complete Responses turn ending in `response.completed` is never re-requested ([#14639](https://github.com/diegosouzapw/OmniRoute/pull/14639)) — thanks @maxmad64bis diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index d285f1f6567..121fa7942b5 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -3429,7 +3429,7 @@ export async function handleChatCore({ // Mid-stream continuation (Fase 4.4): re-request with the partial text as an // assistant prefill. Gated by its own setting and only for OpenAI-compatible - // bodies (makeContinuationBody returns null otherwise). + // request bodies, chat or Responses (makeContinuationBody returns null otherwise). const continueStream = continueMidStreamEnabled ? (assistantSoFar: string) => { const continuationBody = makeContinuationBody( diff --git a/open-sse/services/streamRecovery.ts b/open-sse/services/streamRecovery.ts index 3bd95e5123b..90b59d3ed5b 100644 --- a/open-sse/services/streamRecovery.ts +++ b/open-sse/services/streamRecovery.ts @@ -174,6 +174,13 @@ export function isRetryableStreamError(error: unknown): boolean { // Anthropic `event: message_stop`. Presence means the stream ended cleanly. const OPENAI_DONE_MARKER = "[DONE]"; const ANTHROPIC_STOP_MARKER = "message_stop"; +// A complete Responses-API turn ends with `response.completed`, never `[DONE]` — +// without this the holdback mistakes a short complete Responses stream for a +// truncation and replays it. `failed`/`incomplete` end the turn the same way: +// they are never resumed or retried, only flushed as-is. +const RESPONSES_COMPLETED_MARKER = "response.completed"; +const RESPONSES_FAILED_MARKER = "response.failed"; +const RESPONSES_INCOMPLETE_MARKER = "response.incomplete"; /** * Heuristic check for a terminal SSE marker in the buffered opening window. Used to @@ -184,7 +191,13 @@ const ANTHROPIC_STOP_MARKER = "message_stop"; export function hasTerminalMarker(bytes: Uint8Array): boolean { if (!bytes || bytes.byteLength === 0) return false; const text = new TextDecoder().decode(bytes); - return text.includes(OPENAI_DONE_MARKER) || text.includes(ANTHROPIC_STOP_MARKER); + return ( + text.includes(OPENAI_DONE_MARKER) || + text.includes(ANTHROPIC_STOP_MARKER) || + text.includes(RESPONSES_COMPLETED_MARKER) || + text.includes(RESPONSES_FAILED_MARKER) || + text.includes(RESPONSES_INCOMPLETE_MARKER) + ); } // ──────────────── Mid-stream continuation primitives (Fase 4.4) ──────────────── @@ -227,84 +240,236 @@ export interface OpenAiSseScan { finishReason: string | null; /** True if at least one OpenAI-shaped `choices[].delta` was parsed (format gate). */ parsedOpenAi: boolean; + /** + * True if at least one Responses-API event below was parsed (format gate, + * symmetric with `parsedOpenAi`): text/resasoning deltas, function-call + * signals, or a terminal lifecycle event. + */ + parsedResponses: boolean; } +/** Responses-API `type` values the scanner understands. Anything else is ignored. */ +const RESPONSES_TEXT_DELTA = "response.output_text.delta"; +const RESPONSES_TEXT_DONE = "response.output_text.done"; +const RESPONSES_REASONING_DELTA = "response.reasoning_summary_text.delta"; +const RESPONSES_FN_ARGS_DELTA = "response.function_call_arguments.delta"; +const RESPONSES_FN_ARGS_DONE = "response.function_call_arguments.done"; +const RESPONSES_ITEM_ADDED = "response.output_item.added"; +const RESPONSES_ITEM_DONE = "response.output_item.done"; +const RESPONSES_COMPLETED = "response.completed"; +const RESPONSES_FAILED = "response.failed"; +const RESPONSES_INCOMPLETE = "response.incomplete"; + /** * Scan a slice of OpenAI-compatible SSE for the assistant text, tool-call presence, and * a terminal marker. Non-OpenAI bodies (e.g. Anthropic `content_block_delta` events) parse * to `parsedOpenAi:false` with empty text, so the caller falls back to current behavior. */ export function scanOpenAiSseText(sse: string): OpenAiSseScan { - let text = ""; - let reasoningText = ""; - let sawToolCall = false; - let toolCallFinished = false; - let terminal = false; - let finishReason: string | null = null; - let parsedOpenAi = false; + const fold: SseScanFold = { + text: "", + reasoningText: "", + sawToolCall: false, + toolCallFinished: false, + terminal: false, + finishReason: null, + parsedOpenAi: false, + parsedResponses: false, + }; if (typeof sse !== "string" || sse.length === 0) { return { - text, - reasoningText, - sawToolCall, + text: fold.text, + reasoningText: fold.reasoningText, + sawToolCall: fold.sawToolCall, sawToolCallInFlight: false, - terminal, - finishReason, - parsedOpenAi, + terminal: fold.terminal, + finishReason: fold.finishReason, + parsedOpenAi: fold.parsedOpenAi, + parsedResponses: fold.parsedResponses, }; } for (const line of sse.split("\n")) { - const trimmed = line.trimStart(); - if (!trimmed.startsWith("data:")) continue; - const payload = trimmed.slice(5).trim(); - if (!payload) continue; - if (payload === "[DONE]") { - terminal = true; - continue; - } - let json: unknown; - try { - json = JSON.parse(payload); - } catch { - continue; - } - const choices = (json as { choices?: unknown })?.choices; - if (!Array.isArray(choices)) continue; - for (const choice of choices) { - const delta = (choice as { delta?: unknown })?.delta; - if (delta && typeof delta === "object") { - parsedOpenAi = true; - const content = (delta as { content?: unknown }).content; - if (typeof content === "string") text += content; - const reasoning = (delta as { reasoning_content?: unknown }).reasoning_content; - if (typeof reasoning === "string") reasoningText += reasoning; - const toolCalls = (delta as { tool_calls?: unknown }).tool_calls; - if (Array.isArray(toolCalls) && toolCalls.length > 0) sawToolCall = true; - } - const rawFinishReason = (choice as { finish_reason?: unknown })?.finish_reason; - if (rawFinishReason === "tool_calls") { - // Ends this one choice, but the overall stream/turn stays continuable — - // never counts as the general terminal marker (see OpenAiSseScan.terminal). - toolCallFinished = true; - finishReason = "tool_calls"; - } else if (rawFinishReason != null) { - terminal = true; - if (typeof rawFinishReason === "string") finishReason = rawFinishReason; - } - } + scanSseLine(line, fold); } - const sawToolCallInFlight = sawToolCall && !toolCallFinished; + const sawToolCallInFlight = fold.sawToolCall && !fold.toolCallFinished; return { - text, - reasoningText, - sawToolCall, + text: fold.text, + reasoningText: fold.reasoningText, + sawToolCall: fold.sawToolCall, sawToolCallInFlight, - terminal, - finishReason, - parsedOpenAi, + terminal: fold.terminal, + finishReason: fold.finishReason, + parsedOpenAi: fold.parsedOpenAi, + parsedResponses: fold.parsedResponses, }; } +interface SseScanFold { + text: string; + reasoningText: string; + sawToolCall: boolean; + toolCallFinished: boolean; + terminal: boolean; + finishReason: string | null; + parsedOpenAi: boolean; + parsedResponses: boolean; +} + +/** Parse one raw SSE line into the running scan fold. */ +function scanSseLine(line: string, fold: SseScanFold): void { + const trimmed = line.trimStart(); + if (!trimmed.startsWith("data:")) return; + const payload = trimmed.slice(5).trim(); + if (!payload) return; + if (payload === "[DONE]") { + fold.terminal = true; + return; + } + let json: unknown; + try { + json = JSON.parse(payload); + } catch { + return; + } + const choices = (json as { choices?: unknown })?.choices; + if (Array.isArray(choices)) { + for (const choice of choices) scanChatChoice(choice, fold); + return; + } + // Responses-API events: same scan, allowlisted `type` values only. Gated on + // `data:` JSON exactly like the chat path above, so a bare `event:` line + // without a payload stays invisible here (the holdback substring in + // `hasTerminalMarker` still covers the pre-commit window). + scanResponsesEvent(json as Record, { + addText: (delta) => { + fold.text += delta; + }, + addReasoning: (delta) => { + fold.reasoningText += delta; + }, + markToolCall: () => { + fold.sawToolCall = true; + }, + markToolCallFinished: () => { + fold.sawToolCall = true; + fold.toolCallFinished = true; + }, + markTerminal: () => { + fold.terminal = true; + }, + markParsed: () => { + fold.parsedResponses = true; + }, + }); +} + +/** Fold one chat-completions `choices[]` entry into the running scan fold. */ +function scanChatChoice(choice: unknown, fold: SseScanFold): void { + const delta = (choice as { delta?: unknown })?.delta; + if (delta && typeof delta === "object") { + fold.parsedOpenAi = true; + const content = (delta as { content?: unknown }).content; + if (typeof content === "string") fold.text += content; + const reasoning = (delta as { reasoning_content?: unknown }).reasoning_content; + if (typeof reasoning === "string") fold.reasoningText += reasoning; + const toolCalls = (delta as { tool_calls?: unknown }).tool_calls; + if (Array.isArray(toolCalls) && toolCalls.length > 0) fold.sawToolCall = true; + } + scanChatFinishReason((choice as { finish_reason?: unknown })?.finish_reason, fold); +} + +/** Fold one chat-completions `finish_reason` value into the running scan fold. */ +function scanChatFinishReason(rawFinishReason: unknown, fold: SseScanFold): void { + if (rawFinishReason === "tool_calls") { + // Ends this one choice, but the overall stream/turn stays continuable — + // never counts as the general terminal marker (see OpenAiSseScan.terminal). + fold.toolCallFinished = true; + fold.finishReason = "tool_calls"; + } else if (rawFinishReason != null) { + fold.terminal = true; + if (typeof rawFinishReason === "string") fold.finishReason = rawFinishReason; + } +} + +interface ResponsesScanSink { + addText: (delta: string) => void; + addReasoning: (delta: string) => void; + markToolCall: () => void; + markToolCallFinished: () => void; + markTerminal: () => void; + markParsed: () => void; +} + +/** + * Fold one parsed Responses-API `data:` payload into the scan. Only the + * allowlisted `type` values below set any flag; every other event (lifecycle + * echoes such as `response.output_text.done`, unknown provider extensions) is + * ignored. `output_item.*` is gated on a `function_call` item — a plain text + * turn announces `message` items that must never trip the tool-call lock. + */ +function scanResponsesEvent(json: Record, sink: ResponsesScanSink): void { + const eventType = (json as { type?: unknown }).type; + if (typeof eventType !== "string") return; + if (eventType === RESPONSES_TEXT_DELTA || eventType === RESPONSES_REASONING_DELTA) { + scanResponsesTextDelta(json, eventType, sink); + return; + } + if (eventType === RESPONSES_TEXT_DONE) { + // Lifecycle echo only (downstream drops it); never text, never a tool call. + return; + } + if (eventType === RESPONSES_FN_ARGS_DELTA || eventType === RESPONSES_FN_ARGS_DONE) { + if (eventType === RESPONSES_FN_ARGS_DONE) sink.markToolCallFinished(); + else sink.markToolCall(); + sink.markParsed(); + return; + } + if (eventType === RESPONSES_ITEM_ADDED || eventType === RESPONSES_ITEM_DONE) { + scanResponsesOutputItem(json, eventType, sink); + return; + } + if ( + eventType === RESPONSES_COMPLETED || + eventType === RESPONSES_FAILED || + eventType === RESPONSES_INCOMPLETE + ) { + // End of transport in every status, including a `completed` carrying a + // failed response snapshot: never resumed, never retried. + sink.markTerminal(); + sink.markParsed(); + } +} + +/** Fold a text or reasoning delta event into the scan. */ +function scanResponsesTextDelta( + json: Record, + eventType: string, + sink: ResponsesScanSink +): void { + const delta = (json as { delta?: unknown }).delta; + if (typeof delta !== "string" || delta.length === 0) { + sink.markParsed(); + return; + } + if (eventType === RESPONSES_REASONING_DELTA) sink.addReasoning(delta); + else sink.addText(delta); + sink.markParsed(); +} + +/** Fold an `output_item.*` event into the scan when it carries a function call. */ +function scanResponsesOutputItem( + json: Record, + eventType: string, + sink: ResponsesScanSink +): void { + const item = (json as { item?: unknown }).item; + const itemType = item && typeof item === "object" ? (item as { type?: unknown }).type : undefined; + if (itemType !== "function_call") return; + if (eventType === RESPONSES_ITEM_DONE) sink.markToolCallFinished(); + else sink.markToolCall(); + sink.markParsed(); +} + export interface ContinuableBody { messages?: unknown; stream?: unknown; @@ -313,26 +478,48 @@ export interface ContinuableBody { /** * Build a re-request body that continues from `assistantSoFar` by appending it as an - * assistant turn. When `assistantSoFar` is empty (nothing usable was emitted yet — e.g. a - * clean stop that only produced reasoning), the messages are re-sent unchanged instead of - * appending an empty assistant turn: this simply re-asks for a real answer. Returns null - * only when the body has no `messages` array at all (nothing to continue from). + * assistant turn. Chat bodies (`messages`) win over Responses bodies (`input`) when + * both are present. When `assistantSoFar` is empty (nothing usable was emitted yet — + * e.g. a clean stop that only produced reasoning), the turns are re-sent unchanged + * instead of appending an empty assistant turn: this simply re-asks for a real answer. + * Returns null when the body carries neither a non-empty `messages` array nor a + * non-empty `input` array (nothing to continue from). */ export function makeContinuationBody( body: ContinuableBody, assistantSoFar: string -): (ContinuableBody & { messages: unknown[] }) | null { +): (ContinuableBody & { messages: unknown[] }) | (ContinuableBody & { input: unknown[] }) | null { if (!body || typeof body !== "object") return null; - if (!Array.isArray(body.messages) || body.messages.length === 0) return null; if (typeof assistantSoFar !== "string") return null; - return { - ...body, - messages: - assistantSoFar.length > 0 - ? [...body.messages, { role: "assistant", content: assistantSoFar }] - : [...body.messages], - stream: true, - }; + if (Array.isArray(body.messages) && body.messages.length > 0) { + return { + ...body, + messages: + assistantSoFar.length > 0 + ? [...body.messages, { role: "assistant", content: assistantSoFar }] + : [...body.messages], + stream: true, + }; + } + if (Array.isArray(body.input) && body.input.length > 0) { + return { + ...body, + input: + assistantSoFar.length > 0 + ? [ + ...body.input, + { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: assistantSoFar }], + status: "completed", + }, + ] + : [...body.input], + stream: true, + }; + } + return null; } /** @@ -467,6 +654,7 @@ export function createRecoverableStream( let emittedSawToolCall = false; // any tool_call delta seen, complete or not let emittedToolCallFinish = false; // any finish_reason "tool_calls" seen let emittedParsedOpenAi = false; + let emittedParsedResponses = false; // STREAM_RECOVERY_TOOLCALL_ORDER_FIX, resolved lazily at most once per stream and only // on a recovery decision, so the flag costs nothing on streams that end cleanly. let toolCallOrderFix: boolean | undefined; @@ -494,6 +682,7 @@ export function createRecoverableStream( if (scan.sawToolCall) emittedSawToolCall = true; if (scan.finishReason === "tool_calls") emittedToolCallFinish = true; if (scan.parsedOpenAi) emittedParsedOpenAi = true; + if (scan.parsedResponses) emittedParsedResponses = true; }; const flushHeld = (controller: ReadableStreamDefaultController) => { @@ -543,7 +732,7 @@ export function createRecoverableStream( const canContinue = () => continueEnabled && continuations < maxContinuations && - emittedParsedOpenAi && + (emittedParsedOpenAi || emittedParsedResponses) && !emittedToolCallInFlight && (emittedText.length > 0 ? !emittedTerminal : hallucinatedEmptyStop()) && !toolCallBlocksContinuation(); @@ -551,14 +740,39 @@ export function createRecoverableStream( // Report why a cut is not continued. Silent for non-OpenAI bodies (continuation never // applies to them) so the hook stays quiet on every Claude/Gemini-format stream end. const reportRefusal = () => { - if (!continueEnabled || !emittedParsedOpenAi || !options.onContinueOutcome) return; + if ( + !continueEnabled || + (!emittedParsedOpenAi && !emittedParsedResponses) || + !options.onContinueOutcome + ) + return; let reason: ContinuationRefusal = "not-continuable"; if (continuations >= maxContinuations) reason = "budget"; else if (emittedToolCallInFlight || toolCallBlocksContinuation()) reason = "tool-call"; options.onContinueOutcome({ attempt: continuations, outcome: "refused", reason }); }; + // True for the Responses path (see `isResponsesTurn` below): the client stream + // carries raw upstream Responses events here (translation to chat happens + // downstream), so the stitched suffix and the clean terminal must speak the + // same format. Mixed-format turns fall back to the chat envelope. + const isResponsesTurn = () => emittedParsedResponses && !emittedParsedOpenAi; + const emitCleanTerminal = (controller: ReadableStreamDefaultController) => { + if (isResponsesTurn()) { + // A minimal `response.completed` snapshot: the downstream Responses-to-chat + // translator treats it exactly like a real upstream terminal — usage + // extraction, snapshot tool-call synthesis, then `computeFinishReason` on + // the turn state it already accumulated, so `tool_calls` vs `stop` needs + // no logic here. A second `completed` after `finishReasonSent` resolves to + // no chunk downstream, so re-stating completion is idempotent. + controller.enqueue( + encoder.encode( + 'data: {"type":"response.completed","response":{"status":"completed","output":[]}}\n\n' + ) + ); + return; + } controller.enqueue( encoder.encode('data: {"choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}\n\n') ); @@ -629,12 +843,24 @@ export function createRecoverableStream( } const suffix = overlapResult; if (suffix) { - emit( - controller, - encoder.encode( - `data: ${JSON.stringify({ choices: [{ index: 0, delta: { content: suffix } }] })}\n\n` - ) - ); + if (isResponsesTurn()) { + // Same envelope as the upstream text deltas so the downstream + // translator folds the suffix into a chat `content` chunk with no + // prior announcement events required. + emit( + controller, + encoder.encode( + `data: ${JSON.stringify({ type: "response.output_text.delta", output_index: 0, content_index: 0, delta: suffix })}\n\n` + ) + ); + } else { + emit( + controller, + encoder.encode( + `data: ${JSON.stringify({ choices: [{ index: 0, delta: { content: suffix } }] })}\n\n` + ) + ); + } report({ attempt: continuations, outcome: "suffix", suffixChars: suffix.length }); } // A clean finish, or a tool call we cannot safely stitch, ends the recovered stream. diff --git a/tests/unit/stream-continuation-responses.test.ts b/tests/unit/stream-continuation-responses.test.ts new file mode 100644 index 00000000000..1513f130901 --- /dev/null +++ b/tests/unit/stream-continuation-responses.test.ts @@ -0,0 +1,426 @@ +import { after, afterEach, test } from "node:test"; +import assert from "node:assert/strict"; + +import { + createRecoverableStream, + hasTerminalMarker, + makeContinuationBody, + scanOpenAiSseText, + TruncatedStreamError, +} from "../../open-sse/services/streamRecovery.ts"; +import { openaiResponsesToOpenAIResponse } from "../../open-sse/translator/response/openai-responses.ts"; +import { resetDbInstance } from "../../src/lib/db/core.ts"; + +const enc = new TextEncoder(); +const dec = new TextDecoder(); + +const ORDER_FIX_FLAG = "STREAM_RECOVERY_TOOLCALL_ORDER_FIX"; +const ORIGINAL_ORDER_FIX_FLAG = process.env[ORDER_FIX_FLAG]; + +function setOrderFix(on: boolean) { + if (on) process.env[ORDER_FIX_FLAG] = "true"; + else delete process.env[ORDER_FIX_FLAG]; +} + +afterEach(() => { + if (ORIGINAL_ORDER_FIX_FLAG === undefined) delete process.env[ORDER_FIX_FLAG]; + else process.env[ORDER_FIX_FLAG] = ORIGINAL_ORDER_FIX_FLAG; +}); + +after(() => { + resetDbInstance(); +}); + +/** Clock that jumps past the holdback window so the first chunk commits immediately. */ +function jumpingClock(): () => number { + let t = 0; + return () => (t += 1000); +} + +/** Emit each chunk on its own read, then cut with a retryable truncation. */ +function streamFrom(chunks: string[]): ReadableStream { + let i = 0; + return new ReadableStream({ + pull(controller) { + if (i < chunks.length) { + controller.enqueue(enc.encode(chunks[i++])); + return; + } + controller.error(new TruncatedStreamError()); + }, + }); +} + +function completedStream(sse: string): ReadableStream { + return new ReadableStream({ + start(controller) { + controller.enqueue(enc.encode(sse)); + controller.close(); + }, + }); +} + +async function collectText(stream: ReadableStream): Promise { + const reader = stream.getReader(); + try { + let out = ""; + for (;;) { + let read: ReadableStreamReadResult; + try { + read = await reader.read(); + } catch { + // A refused continuation surfaces the original truncation error. + break; + } + if (read.done) break; + if (read.value) out += dec.decode(read.value, { stream: true }); + } + return out; + } finally { + reader.releaseLock(); + } +} + +const textDelta = (delta: string) => + `data: ${JSON.stringify({ type: "response.output_text.delta", output_index: 0, content_index: 0, delta })}\n\n`; +const textDone = + 'data: {"type":"response.output_text.done","output_index":0,"content_index":0}\n\n'; +const completed = (status = "completed") => + `data: ${JSON.stringify({ type: "response.completed", response: { status, output: [] } })}\n\n`; +const summaryDelta = (delta: string) => + `data: ${JSON.stringify({ type: "response.reasoning_summary_text.delta", item_id: "rs_1", output_index: 0, summary_index: 0, delta })}\n\n`; +const fnArgsDelta = (itemId: string, outputIndex: number, delta: string) => + `data: ${JSON.stringify({ type: "response.function_call_arguments.delta", item_id: itemId, output_index: outputIndex, delta })}\n\n`; +const fnArgsDone = (itemId: string, outputIndex: number) => + `data: ${JSON.stringify({ type: "response.function_call_arguments.done", item_id: itemId, output_index: outputIndex })}\n\n`; +const fnAdded = (callId: string, name: string) => + `data: ${JSON.stringify({ type: "response.output_item.added", output_index: 0, item: { type: "function_call", id: "item_1", call_id: callId, name, arguments: "{}" } })}\n\n`; +const msgAdded = () => + 'data: {"type":"response.output_item.added","output_index":0,"item":{"type":"message","id":"msg_1","status":"in_progress","role":"assistant","content":[]}}\n\n'; + +// ── scan: text deltas ───────────────────────────────────────────────────────── + +test("scan accumulates Responses text deltas without flagging the chat format", () => { + const r = scanOpenAiSseText(textDelta("Hello") + textDelta(" world")); + assert.equal(r.text, "Hello world"); + assert.equal(r.parsedResponses, true); + assert.equal(r.parsedOpenAi, false); + assert.equal(r.sawToolCall, false); + assert.equal(r.terminal, false); +}); + +test("scan ignores the named output_text.done lifecycle echo", () => { + const r = scanOpenAiSseText(textDelta("hi") + textDone); + assert.equal(r.text, "hi"); + assert.equal(r.parsedResponses, true); + assert.equal(r.sawToolCall, false); + assert.equal(r.terminal, false); +}); + +test("scan keeps reasoning summary text out of the visible text field", () => { + const r = scanOpenAiSseText(summaryDelta("thinking...") + summaryDelta(" more")); + assert.equal(r.reasoningText, "thinking... more"); + assert.equal(r.text, ""); + assert.equal(r.parsedResponses, true); +}); + +test("scan treats response.completed as terminal whatever the status", () => { + const ok = scanOpenAiSseText(textDelta("hi") + completed("completed")); + assert.equal(ok.text, "hi"); + assert.equal(ok.terminal, true); + assert.equal(ok.parsedResponses, true); + const failedStatus = scanOpenAiSseText(textDelta("hi") + completed("failed")); + assert.equal(failedStatus.terminal, true); +}); + +test("scan treats response.failed and response.incomplete as terminal", () => { + const failed = scanOpenAiSseText( + 'data: {"type":"response.failed","response":{"status":"failed"}}\n\n' + ); + assert.equal(failed.terminal, true); + assert.equal(failed.parsedResponses, true); + const incomplete = scanOpenAiSseText( + 'data: {"type":"response.incomplete","response":{"status":"incomplete"}}\n\n' + ); + assert.equal(incomplete.terminal, true); + assert.equal(incomplete.parsedResponses, true); +}); + +// ── scan: tool-call signals ─────────────────────────────────────────────────── + +test("scan flags Responses function-call deltas as an in-flight tool call", () => { + const r = scanOpenAiSseText(fnArgsDelta("item_1", 0, '{"q":1}')); + assert.equal(r.sawToolCall, true); + assert.equal(r.sawToolCallInFlight, true); + assert.equal(r.parsedResponses, true); +}); + +test("scan closes the in-flight tool call on the matching done event", () => { + const r = scanOpenAiSseText(fnArgsDelta("item_1", 0, '{"q":1}') + fnArgsDone("item_1", 0)); + assert.equal(r.sawToolCall, true); + assert.equal(r.sawToolCallInFlight, false); + assert.equal(r.finishReason, null); +}); + +test("scan only flags output_item events carrying a function_call item", () => { + const message = scanOpenAiSseText(msgAdded()); + assert.equal(message.sawToolCall, false); + const fn = scanOpenAiSseText(fnAdded("call_1", "lookup")); + assert.equal(fn.sawToolCall, true); + assert.equal(fn.sawToolCallInFlight, true); + assert.equal(fn.parsedResponses, true); +}); + +test("scan ignores bare event lines without a data payload", () => { + const r = scanOpenAiSseText("event: response.completed\n\n"); + assert.equal(r.terminal, false); + assert.equal(r.parsedResponses, false); +}); + +// ── hasTerminalMarker ───────────────────────────────────────────────────────── + +test("hasTerminalMarker recognizes Responses terminal events without [DONE]", () => { + const withCompleted = enc.encode(textDelta("hi") + completed()); + assert.equal(hasTerminalMarker(withCompleted), true); + const failedOnly = enc.encode( + 'data: {"type":"response.failed","response":{"status":"failed"}}\n\n' + ); + assert.equal(hasTerminalMarker(failedOnly), true); + const incompleteOnly = enc.encode( + 'data: {"type":"response.incomplete","response":{"status":"incomplete"}}\n\n' + ); + assert.equal(hasTerminalMarker(incompleteOnly), true); +}); + +test("hasTerminalMarker stays false on a truncated Responses stream", () => { + const truncated = enc.encode(textDelta("partial answer")); + assert.equal(hasTerminalMarker(truncated), false); +}); + +test("hasTerminalMarker covers a bare event line without a data payload", () => { + assert.equal(hasTerminalMarker(enc.encode("event: response.completed\n\n")), true); +}); + +// ── makeContinuationBody ────────────────────────────────────────────────────── + +test("makeContinuationBody replays received text as an assistant input item", () => { + const body = { model: "m", input: [{ type: "message", role: "user", content: "hi" }] }; + const out = makeContinuationBody(body as never, "partial answer") as { + input: unknown[]; + stream: unknown; + }; + assert.equal(out.stream, true); + assert.equal(out.input.length, 2); + assert.deepEqual(out.input[1], { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "partial answer" }], + status: "completed", + }); +}); + +test("makeContinuationBody keeps messages priority over input", () => { + const body = { + model: "m", + messages: [{ role: "user", content: "hi" }], + input: [{ type: "message", role: "user", content: "hi" }], + }; + const out = makeContinuationBody(body as never, "partial") as { + messages: unknown[]; + } & Record; + assert.ok(Array.isArray(out.messages)); + assert.equal(out.messages.length, 2); + assert.ok(!("input" in out) || out.input === body.input); +}); + +test("makeContinuationBody rejects empty or missing input like empty messages", () => { + assert.equal(makeContinuationBody({ model: "m", input: [] } as never, "t"), null); + assert.equal(makeContinuationBody({ model: "m" } as never, "t"), null); +}); + +test("makeContinuationBody re-sends input unchanged on an empty prefill", () => { + const body = { model: "m", input: [{ type: "message", role: "user", content: "hi" }] }; + const out = makeContinuationBody(body as never, "") as { input: unknown[] }; + assert.equal(out.input.length, 1); +}); + +// ── trap 1: a short complete Responses stream is never resumed ──────────────── + +test("a short complete Responses stream is never resumed after serving", async () => { + const sse = textDelta("Hello there") + completed(); + let continuations = 0; + const stream = createRecoverableStream(completedStream(sse), async () => null, { + finalize: () => {}, + now: jumpingClock(), + continueStream: async () => { + continuations += 1; + return null; + }, + }); + const out = await collectText(stream); + assert.ok(out.includes("Hello there")); + assert.equal(continuations, 0, "a served complete response must not be re-requested"); +}); + +// ── cut after content resumes with a Responses envelope ─────────────────────── + +test("a cut after content resumes and stitches a Responses suffix", async () => { + const initial = streamFrom([textDelta("Hello there world")]); + let prefill = ""; + const stream = createRecoverableStream(initial, async () => null, { + finalize: () => {}, + now: jumpingClock(), + continueStream: async (soFar: string) => { + prefill = soFar; + return completedStream(textDelta("there world, nice to meet you!") + completed()); + }, + }); + const out = await collectText(stream); + assert.equal(prefill, "Hello there world"); + const scan = scanOpenAiSseText(out); + assert.equal(scan.parsedResponses, true); + assert.ok(out.includes("response.output_text.delta")); + assert.ok(out.includes("nice to meet you!")); + assert.ok(out.includes("response.completed")); +}); + +test("an in-flight Responses tool call refuses the cut without a re-request", async () => { + const initial = streamFrom([textDelta("partial") + fnArgsDelta("item_9", 0, '{"q":1}')]); + let continuations = 0; + const outcomes: unknown[] = []; + const stream = createRecoverableStream(initial, async () => null, { + finalize: () => {}, + now: jumpingClock(), + continueStream: async () => { + continuations += 1; + return null; + }, + onContinueOutcome: (event) => { + outcomes.push(event); + }, + }); + await collectText(stream); + assert.equal(continuations, 0); + assert.ok( + outcomes.some( + (event) => + (event as { outcome?: string }).outcome === "refused" && + (event as { reason?: string }).reason === "tool-call" + ) + ); +}); + +test("a finished tool call stays continuable when the order fix is off", async () => { + setOrderFix(false); + const initial = streamFrom([ + fnAdded("call_1", "lookup") + fnArgsDone("item_1", 0) + textDelta("trailing prose"), + ]); + let continuations = 0; + const stream = createRecoverableStream(initial, async () => null, { + finalize: () => {}, + now: jumpingClock(), + continueStream: async () => { + continuations += 1; + return completedStream(textDelta("trailing prose plus more") + completed()); + }, + }); + const out = await collectText(stream); + assert.equal(continuations, 1); + assert.ok(out.includes("plus more")); +}); + +test("a finished tool call blocks continuation when the order fix is on", async () => { + setOrderFix(true); + const initial = streamFrom([ + fnAdded("call_1", "lookup") + fnArgsDone("item_1", 0) + textDelta("trailing prose"), + ]); + let continuations = 0; + const stream = createRecoverableStream(initial, async () => null, { + finalize: () => {}, + now: jumpingClock(), + continueStream: async () => { + continuations += 1; + return null; + }, + }); + await collectText(stream); + assert.equal(continuations, 0); +}); + +// ── stitched envelope drives the downstream translator to a stop chunk ──────── + +test("the stitched Responses envelope closes through the downstream translator", () => { + const state = {}; + const suffixChunk = openaiResponsesToOpenAIResponse( + { + type: "response.output_text.delta", + output_index: 0, + content_index: 0, + delta: " world", + }, + state + ) as { choices?: Array<{ delta?: { content?: string } }> } | null; + assert.ok(suffixChunk); + assert.equal(suffixChunk?.choices?.[0]?.delta?.content, " world"); + const finalChunk = openaiResponsesToOpenAIResponse( + { type: "response.completed", response: { status: "completed", output: [] } }, + state + ) as { choices?: Array<{ finish_reason?: string }> } | null; + assert.ok(finalChunk); + assert.equal(finalChunk?.choices?.[0]?.finish_reason, "stop"); +}); + +test("a completed snapshot carrying a function call closes as tool_calls", () => { + const state = {}; + openaiResponsesToOpenAIResponse( + { + type: "response.output_item.added", + output_index: 0, + item: { + type: "function_call", + id: "item_1", + call_id: "call_1", + name: "lookup", + arguments: "{}", + }, + }, + state + ); + const finalChunk = openaiResponsesToOpenAIResponse( + { + type: "response.completed", + response: { + status: "completed", + output: [ + { + type: "function_call", + id: "item_1", + call_id: "call_1", + name: "lookup", + arguments: "{}", + }, + ], + }, + }, + state + ); + const chunks = Array.isArray(finalChunk) ? finalChunk : [finalChunk]; + const reasons = chunks.flatMap( + (chunk) => (chunk as { choices?: Array<{ finish_reason?: string }> })?.choices ?? [] + ); + assert.ok(reasons.some((choice) => choice.finish_reason === "tool_calls")); +}); + +test("a synthetic completed after finish_reason was sent resolves to no chunk", () => { + const state = {}; + const first = openaiResponsesToOpenAIResponse( + { type: "response.completed", response: { status: "completed", output: [] } }, + state + ); + assert.ok(first); + const second = openaiResponsesToOpenAIResponse( + { type: "response.completed", response: { status: "completed", output: [] } }, + state + ); + assert.equal(second, null); +});