Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
92 changes: 89 additions & 3 deletions open-sse/utils/stream.js
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,43 @@ const STREAM_MODE = {
PASSTHROUGH: "passthrough" // No translation, normalize output, extract usage
};

const OPENAI_RESPONSES_TERMINAL_EVENTS = new Set([
"response.completed",
"response.failed",
"error"
]);

function getOpenAIResponsesEventName(eventName, chunk) {
if (eventName) return eventName;
if (chunk && typeof chunk.type === "string") return chunk.type;
return null;
}

function isOpenAIResponsesTerminalEvent(eventName, chunk) {
const type = getOpenAIResponsesEventName(eventName, chunk);
if (OPENAI_RESPONSES_TERMINAL_EVENTS.has(type)) return true;
const status = chunk?.response?.status;
return status === "completed" || status === "failed";
}

function formatIncompleteOpenAIResponsesStreamFailure() {
return formatSSE({
event: "response.failed",
data: {
type: "response.failed",
response: {
id: `resp_${Date.now()}`,
status: "failed",
error: {
type: "stream_error",
code: "stream_disconnected",
message: "stream closed before response.completed"
}
}
}
}, FORMATS.OPENAI_RESPONSES);
}

/**
* Create unified SSE transform stream
* @param {object} options
Expand Down Expand Up @@ -62,6 +99,9 @@ export function createSSEStream(options = {}) {
let sseLineCount = 0;
let sseEmittedCount = 0;
const eventTypeCounts = {};
let currentOpenAIResponsesEvent = null;
let openAIResponsesTerminalSeen = false;
let openAIResponsesDoneSent = false;

return new TransformStream({
transform(chunk, controller) {
Expand All @@ -83,6 +123,10 @@ export function createSSEStream(options = {}) {
}
}

if (mode === STREAM_MODE.TRANSLATE && targetFormat === FORMATS.OPENAI_RESPONSES && trimmed.startsWith("event:")) {
currentOpenAIResponsesEvent = trimmed.slice(6).trim();
}

// Passthrough mode: normalize and forward
if (mode === STREAM_MODE.PASSTHROUGH) {
let output;
Expand Down Expand Up @@ -174,12 +218,31 @@ export function createSSEStream(options = {}) {
const parsed = parseSSELine(trimmed, targetFormat);
if (!parsed) continue;

const isOpenAIResponsesStream = targetFormat === FORMATS.OPENAI_RESPONSES;
const keepsOpenAIResponsesFormat = isOpenAIResponsesStream && sourceFormat === FORMATS.OPENAI_RESPONSES;
const openAIResponsesEventName = isOpenAIResponsesStream
? getOpenAIResponsesEventName(currentOpenAIResponsesEvent, parsed)
: null;

if (isOpenAIResponsesStream && isOpenAIResponsesTerminalEvent(openAIResponsesEventName, parsed)) {
openAIResponsesTerminalSeen = true;
}

// For Ollama: done=true is the final chunk with finish_reason/usage, must translate
// For other formats: done=true is the [DONE] sentinel, skip
if (parsed && parsed.done && targetFormat !== FORMATS.OLLAMA) {
if (keepsOpenAIResponsesFormat && !openAIResponsesTerminalSeen) {
const failedOutput = formatIncompleteOpenAIResponsesStreamFailure();
reqLogger?.appendConvertedChunk?.(failedOutput);
controller.enqueue(sharedEncoder.encode(failedOutput));
openAIResponsesTerminalSeen = true;
sseEmittedCount++;
}

const output = "data: [DONE]\n\n";
reqLogger?.appendConvertedChunk?.(output);
controller.enqueue(sharedEncoder.encode(output));
if (keepsOpenAIResponsesFormat) openAIResponsesDoneSent = true;
continue;
}

Expand Down Expand Up @@ -224,6 +287,17 @@ export function createSSEStream(options = {}) {
const extracted = extractUsage(parsed);
if (extracted) state.usage = extracted; // Keep original usage for logging

if (keepsOpenAIResponsesFormat && openAIResponsesEventName) {
const output = formatSSE({ event: openAIResponsesEventName, data: parsed }, sourceFormat);
reqLogger?.appendConvertedChunk?.(output);
controller.enqueue(sharedEncoder.encode(output));
currentOpenAIResponsesEvent = null;
sseEmittedCount++;
continue;
}

currentOpenAIResponsesEvent = null;

// Translate: targetFormat -> openai -> sourceFormat
const translated = translateResponse(targetFormat, sourceFormat, parsed, state);

Expand Down Expand Up @@ -322,6 +396,7 @@ export function createSSEStream(options = {}) {

if (translated?.length > 0) {
for (const item of translated) {
if (item === null || item === undefined) continue;
const output = formatSSE(item, sourceFormat);
reqLogger?.appendConvertedChunk?.(output);
controller.enqueue(sharedEncoder.encode(output));
Expand All @@ -341,15 +416,26 @@ export function createSSEStream(options = {}) {

if (flushed?.length > 0) {
for (const item of flushed) {
if (item === null || item === undefined) continue;
const output = formatSSE(item, sourceFormat);
reqLogger?.appendConvertedChunk?.(output);
controller.enqueue(sharedEncoder.encode(output));
}
}

const doneOutput = "data: [DONE]\n\n";
reqLogger?.appendConvertedChunk?.(doneOutput);
controller.enqueue(sharedEncoder.encode(doneOutput));
const keepsOpenAIResponsesFormat = targetFormat === FORMATS.OPENAI_RESPONSES && sourceFormat === FORMATS.OPENAI_RESPONSES;
if (keepsOpenAIResponsesFormat && !openAIResponsesTerminalSeen) {
const failedOutput = formatIncompleteOpenAIResponsesStreamFailure();
reqLogger?.appendConvertedChunk?.(failedOutput);
controller.enqueue(sharedEncoder.encode(failedOutput));
openAIResponsesTerminalSeen = true;
}

if (!keepsOpenAIResponsesFormat || !openAIResponsesDoneSent) {
const doneOutput = "data: [DONE]\n\n";
reqLogger?.appendConvertedChunk?.(doneOutput);
controller.enqueue(sharedEncoder.encode(doneOutput));
}

if (!hasValidUsage(state?.usage) && totalContentLength > 0) {
state.usage = estimateUsage(body, totalContentLength, sourceFormat);
Expand Down
83 changes: 83 additions & 0 deletions tests/unit/openai-responses-terminal-event.test.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
import { describe, expect, it } from "vitest";

import { FORMATS } from "../../open-sse/translator/formats.js";
import { createSSETransformStreamWithLogger } from "../../open-sse/utils/stream.js";

async function runTransform(input) {
const encoder = new TextEncoder();
const stream = new ReadableStream({
start(controller) {
controller.enqueue(encoder.encode(input));
controller.close();
},
});

const output = stream.pipeThrough(
createSSETransformStreamWithLogger(
FORMATS.OPENAI_RESPONSES,
FORMATS.OPENAI_RESPONSES,
"codex",
null,
null,
"gpt-5.5",
),
);

const reader = output.getReader();
const decoder = new TextDecoder();
let text = "";

while (true) {
const { value, done } = await reader.read();
if (done) break;
text += decoder.decode(value, { stream: true });
}

text += decoder.decode();
return text;
}

describe("OpenAI Responses streaming termination", () => {
it("emits a response.failed event when a Responses stream closes before a terminal event", async () => {
const output = await runTransform([
`event: response.created`,
`data: ${JSON.stringify({ type: "response.created", response: { id: "resp_test", status: "in_progress" } })}`,
"",
`event: response.output_text.delta`,
`data: ${JSON.stringify({ type: "response.output_text.delta", delta: "partial" })}`,
"",
].join("\n"));

expect(output).toContain("event: response.failed");
expect(output).toContain('"type":"response.failed"');
expect(output).not.toContain("data: null");
expect(output).toContain("data: [DONE]");
});

it("does not add response.failed when a Responses stream already completed", async () => {
const output = await runTransform([
`event: response.completed`,
`data: ${JSON.stringify({ type: "response.completed", response: { id: "resp_test", status: "completed" } })}`,
"",
].join("\n"));

expect(output).toContain("event: response.completed");
expect(output).not.toContain("event: response.failed");
expect(output).not.toContain("data: null");
expect(output).toContain("data: [DONE]");
});

it("emits response.failed before DONE when a Responses stream sends DONE without a terminal event", async () => {
const output = await runTransform([
`event: response.created`,
`data: ${JSON.stringify({ type: "response.created", response: { id: "resp_test", status: "in_progress" } })}`,
"",
"data: [DONE]",
"",
].join("\n"));

expect(output.indexOf("event: response.failed")).toBeLessThan(output.indexOf("data: [DONE]"));
expect(output.match(/data: \[DONE\]/g)).toHaveLength(1);
expect(output).not.toContain("data: null");
});
});