From 6d110e6d10a9a760d75fe6d7d5414c7796437e75 Mon Sep 17 00:00:00 2001 From: pouya shahrdami <131306412+pouyashahrdami@users.noreply.github.com> Date: Mon, 3 Aug 2026 00:01:19 -0400 Subject: [PATCH 1/2] fix(streaming): surface mid-stream errors when iterating buffered streams The `[Symbol.asyncIterator]()` implementations in ChatCompletionStream and ResponseStream buffer events in `pushQueue` when the producer is ahead of the consumer. The error/abort handlers only rejected readers waiting at that moment and never stored the error, so an error that fired while events were buffered (no reader waiting) was lost: the consumer drained the buffer and `next()` returned done without throwing. Capture the terminal error and re-throw it once the buffer drains, mirroring the existing EventStream.events() pattern. Behavior when a reader is already waiting is unchanged. Adds a regression test to each stream. Closes #2046 --- src/lib/ChatCompletionStream.ts | 37 +++++++++++-------- src/lib/responses/ResponseStream.ts | 37 +++++++++++-------- tests/lib/ChatCompletionStream.test.ts | 47 ++++++++++++++++++++++++ tests/lib/ResponseStream.test.ts | 49 ++++++++++++++++++++++++++ 4 files changed, 140 insertions(+), 30 deletions(-) diff --git a/src/lib/ChatCompletionStream.ts b/src/lib/ChatCompletionStream.ts index 61fc127e7e..88ae1f98ac 100644 --- a/src/lib/ChatCompletionStream.ts +++ b/src/lib/ChatCompletionStream.ts @@ -667,6 +667,22 @@ export class ChatCompletionStream reject: (err: unknown) => void; }[] = []; let done = false; + // Capture a terminal error so it can still be surfaced after any buffered + // chunks drain, even when no reader was waiting at the time it fired. + let failure: unknown = null; + let failureDelivered = false; + + const rejectQueuedReaders = (err: unknown) => { + done = true; + failure = err; + if (readQueue.length) { + failureDelivered = true; + } + for (const reader of readQueue) { + reader.reject(err); + } + readQueue.length = 0; + }; this.on('chunk', (chunk) => { const reader = readQueue.shift(); @@ -685,25 +701,16 @@ export class ChatCompletionStream readQueue.length = 0; }); - this.on('abort', (err) => { - done = true; - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }); - - this.on('error', (err) => { - done = true; - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }); + this.on('abort', rejectQueuedReaders); + this.on('error', rejectQueuedReaders); return { next: async (): Promise> => { if (!pushQueue.length) { + if (failure !== null && !failureDelivered) { + failureDelivered = true; + return Promise.reject(failure); + } if (done) { return { value: undefined, done: true }; } diff --git a/src/lib/responses/ResponseStream.ts b/src/lib/responses/ResponseStream.ts index a857ef10ae..130f163489 100644 --- a/src/lib/responses/ResponseStream.ts +++ b/src/lib/responses/ResponseStream.ts @@ -233,6 +233,22 @@ export class ResponseStream reject: (err: unknown) => void; }[] = []; let done = false; + // Capture a terminal error so it can still be surfaced after any buffered + // events drain, even when no reader was waiting at the time it fired. + let failure: unknown = null; + let failureDelivered = false; + + const rejectQueuedReaders = (err: unknown) => { + done = true; + failure = err; + if (readQueue.length) { + failureDelivered = true; + } + for (const reader of readQueue) { + reader.reject(err); + } + readQueue.length = 0; + }; this.on('event', (event) => { const reader = readQueue.shift(); @@ -251,25 +267,16 @@ export class ResponseStream readQueue.length = 0; }); - this.on('abort', (err) => { - done = true; - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }); - - this.on('error', (err) => { - done = true; - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }); + this.on('abort', rejectQueuedReaders); + this.on('error', rejectQueuedReaders); return { next: async (): Promise> => { if (!pushQueue.length) { + if (failure !== null && !failureDelivered) { + failureDelivered = true; + return Promise.reject(failure); + } if (done) { return { value: undefined, done: true }; } diff --git a/tests/lib/ChatCompletionStream.test.ts b/tests/lib/ChatCompletionStream.test.ts index d343efd037..7af52005be 100644 --- a/tests/lib/ChatCompletionStream.test.ts +++ b/tests/lib/ChatCompletionStream.test.ts @@ -752,4 +752,51 @@ describe('.stream()', () => { `); expect(capturedLogProbs?.length).toEqual(choice?.logprobs?.refusal?.length); }); + + it('surfaces a mid-stream error when chunks are buffered before consumption', async () => { + const chunks = [ + { + id: 'chatcmpl-test', + object: 'chat.completion.chunk', + created: 1, + model: 'gpt-4', + choices: [{ index: 0, delta: { role: 'assistant', content: 'hel' }, finish_reason: null }], + }, + { + id: 'chatcmpl-test', + object: 'chat.completion.chunk', + created: 1, + model: 'gpt-4', + choices: [{ index: 0, delta: { content: 'lo' }, finish_reason: null }], + }, + ] as unknown as OpenAI.Chat.ChatCompletionChunk[]; + // Yield valid chunks, then throw to error the stream after they have been + // delivered (mimics a connection drop mid-response). + const readable = new Stream(async function* () { + for (const chunk of chunks) yield chunk; + throw new Error('network boom'); + }, new AbortController()).toReadableStream(); + + const stream = ChatCompletionStream.fromReadableStream(readable); + // Grab the iterator (registering its listeners) but do not consume yet, so + // the valid chunks and the error land while no reader is waiting: they + // buffer in the iterator's internal queue instead of rejecting a pending + // reader. + const iterator = stream[Symbol.asyncIterator](); + await new Promise((resolve) => setTimeout(resolve, 10)); + + const collected: OpenAI.Chat.ChatCompletionChunk[] = []; + let caught: unknown = null; + try { + let result: IteratorResult; + while (!(result = await iterator.next()).done) { + collected.push(result.value); + } + } catch (err) { + caught = err; + } + + expect(collected).toHaveLength(chunks.length); + expect(caught).toBeInstanceOf(Error); + }); }); diff --git a/tests/lib/ResponseStream.test.ts b/tests/lib/ResponseStream.test.ts index 4d971273eb..31a915cc24 100644 --- a/tests/lib/ResponseStream.test.ts +++ b/tests/lib/ResponseStream.test.ts @@ -253,6 +253,55 @@ describe('.stream()', () => { } expect(final.output_text).toBe('The answer is 42'); }); + + it('surfaces a mid-stream error when events are buffered before consumption', async () => { + // Two valid events, then a malformed delta that references a missing output + // index so accumulation throws mid-stream (the stream itself closes cleanly, + // so the two earlier events are delivered). + const validEvents: ResponseStreamEvent[] = [ + { type: 'response.created', sequence_number: 0, response: makeResponse() }, + { + type: 'response.output_item.added', + sequence_number: 1, + output_index: 0, + item: { id: 'msg_1', type: 'message', role: 'assistant', status: 'in_progress', content: [] }, + }, + ]; + const malformedEvent = { + type: 'response.output_text.delta', + sequence_number: 2, + item_id: 'msg_1', + output_index: 99, + content_index: 0, + delta: 'boom', + logprobs: [], + } as unknown as ResponseStreamEvent; + + const stream = ResponseStream.fromReadableStream( + readableStreamFromEvents([...validEvents, malformedEvent]), + ); + // Grab the iterator (registering its listeners) but do not consume yet, so + // the valid events and the error land while no reader is waiting: they + // buffer in the iterator's internal queue instead of rejecting a pending + // reader. + const iterator = stream[Symbol.asyncIterator](); + // Let the producer drain the readable and hit the error before we read. + await new Promise((resolve) => setTimeout(resolve, 10)); + + const collected: ResponseStreamEvent[] = []; + let caught: unknown = null; + try { + let result: IteratorResult; + while (!(result = await iterator.next()).done) { + collected.push(result.value); + } + } catch (err) { + caught = err; + } + + expect(collected).toHaveLength(validEvents.length); + expect(caught).toBeInstanceOf(Error); + }); }); function readableStreamFromEvents(events: ResponseStreamEvent[]) { From 5f986015405d3f7cdc498c0df5f13ccbbd5944f7 Mon Sep 17 00:00:00 2001 From: pouya shahrdami <131306412+pouyashahrdami@users.noreply.github.com> Date: Mon, 3 Aug 2026 19:49:20 -0400 Subject: [PATCH 2/2] fix(streaming): share one buffered-iterator adapter so terminal errors survive in every stream bridge Address review: the retained-failure fix previously covered only ChatCompletionStream and ResponseStream, while ChatCompletionStreamingRunner.toReadableStream() and AssistantStream[Symbol.asyncIterator]() kept their own copies of the buffering logic and still dropped terminal errors once events were buffered with no reader waiting. - extract EventStream#_createIterator, a shared buffered async-iterator adapter (the events() logic, generalized); events() now delegates to it - rebuild all four duplicated adapters on the shared helper: chunk/event iterators keep their break-aborts semantics via onReturn, the runner bridge keeps eager listener registration, AssistantStream keeps structuredClone-at-push - remove listeners and drain queues on end/return (previously leaked); mark the terminal promise handled when the consumer explicitly ends iteration so the self-inflicted abort no longer surfaces as an unhandled rejection - add regression tests for the runner bridge and AssistantStream paths, plus edge-case tests for abort retention, exactly-once failure delivery, break-aborts, post-end iteration, and events('error') - make the buffered-error tests deterministic by awaiting the stream's terminal signal and asserting the exact error class and message --- src/lib/AssistantStream.ts | 67 ++------- src/lib/ChatCompletionStream.ts | 71 +--------- src/lib/ChatCompletionStreamingRunner.ts | 102 ++++---------- src/lib/EventStream.ts | 70 ++++++++-- src/lib/responses/ResponseStream.ts | 71 +--------- tests/lib/ChatCompletionStream.test.ts | 140 ++++++++++++++++++- tests/lib/EventStream.test.ts | 30 +++- tests/lib/ResponseStream.test.ts | 10 +- tests/streaming/assistants/assistant.test.ts | 59 +++++++- 9 files changed, 334 insertions(+), 286 deletions(-) diff --git a/src/lib/AssistantStream.ts b/src/lib/AssistantStream.ts index f6814289ac..559e9c40c2 100644 --- a/src/lib/AssistantStream.ts +++ b/src/lib/AssistantStream.ts @@ -94,66 +94,15 @@ export class AssistantStream #currentRunStepSnapshot: Runs.RunStep | undefined; [Symbol.asyncIterator](): AsyncIterator { - const pushQueue: AssistantStreamEvent[] = []; - const readQueue: { - resolve: (chunk: AssistantStreamEvent | undefined) => void; - reject: (err: unknown) => void; - }[] = []; - let done = false; - - //Catch all for passing along all events - this.on('event', (event) => { - const eventCopy = structuredClone(event); - const reader = readQueue.shift(); - if (reader) { - reader.resolve(eventCopy); - } else { - pushQueue.push(eventCopy); - } - }); - - this.on('end', () => { - done = true; - for (const reader of readQueue) { - reader.resolve(undefined); - } - readQueue.length = 0; - }); - - this.on('abort', (err) => { - done = true; - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }); - - this.on('error', (err) => { - done = true; - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }); - - return { - next: async (): Promise> => { - if (!pushQueue.length) { - if (done) { - return { value: undefined, done: true }; - } - return new Promise((resolve, reject) => - readQueue.push({ resolve, reject }), - ).then((chunk) => (chunk ? { value: chunk, done: false } : { value: undefined, done: true })); - } - const chunk = pushQueue.shift()!; - return { value: chunk, done: false }; + return this._createIterator( + (push) => { + //Catch all for passing along all events + const onEvent = (event: AssistantStreamEvent) => push(structuredClone(event)); + this.on('event', onEvent); + return () => this.off('event', onEvent); }, - return: async () => { - this.abort(); - return { value: undefined, done: true }; - }, - }; + { onReturn: () => this.abort() }, + ); } static fromReadableStream(stream: ReadableStream): AssistantStream { diff --git a/src/lib/ChatCompletionStream.ts b/src/lib/ChatCompletionStream.ts index 88ae1f98ac..5eff19282e 100644 --- a/src/lib/ChatCompletionStream.ts +++ b/src/lib/ChatCompletionStream.ts @@ -661,71 +661,14 @@ export class ChatCompletionStream } [Symbol.asyncIterator](this: ChatCompletionStream): AsyncIterator { - const pushQueue: ChatCompletionChunk[] = []; - const readQueue: { - resolve: (chunk: ChatCompletionChunk | undefined) => void; - reject: (err: unknown) => void; - }[] = []; - let done = false; - // Capture a terminal error so it can still be surfaced after any buffered - // chunks drain, even when no reader was waiting at the time it fired. - let failure: unknown = null; - let failureDelivered = false; - - const rejectQueuedReaders = (err: unknown) => { - done = true; - failure = err; - if (readQueue.length) { - failureDelivered = true; - } - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }; - - this.on('chunk', (chunk) => { - const reader = readQueue.shift(); - if (reader) { - reader.resolve(chunk); - } else { - pushQueue.push(chunk); - } - }); - - this.on('end', () => { - done = true; - for (const reader of readQueue) { - reader.resolve(undefined); - } - readQueue.length = 0; - }); - - this.on('abort', rejectQueuedReaders); - this.on('error', rejectQueuedReaders); - - return { - next: async (): Promise> => { - if (!pushQueue.length) { - if (failure !== null && !failureDelivered) { - failureDelivered = true; - return Promise.reject(failure); - } - if (done) { - return { value: undefined, done: true }; - } - return new Promise((resolve, reject) => - readQueue.push({ resolve, reject }), - ).then((chunk) => (chunk ? { value: chunk, done: false } : { value: undefined, done: true })); - } - const chunk = pushQueue.shift()!; - return { value: chunk, done: false }; + return this._createIterator( + (push) => { + const onChunk = (chunk: ChatCompletionChunk) => push(chunk); + this.on('chunk', onChunk); + return () => this.off('chunk', onChunk); }, - return: async () => { - this.abort(); - return { value: undefined, done: true }; - }, - }; + { onReturn: () => this.abort() }, + ); } toReadableStream(): ReadableStream { diff --git a/src/lib/ChatCompletionStreamingRunner.ts b/src/lib/ChatCompletionStreamingRunner.ts index a3f33880e3..007284108e 100644 --- a/src/lib/ChatCompletionStreamingRunner.ts +++ b/src/lib/ChatCompletionStreamingRunner.ts @@ -69,89 +69,39 @@ export class ChatCompletionStreamingRunner } override toReadableStream(): ReadableStream { - const pushQueue: ChatCompletionReadableStreamItem[] = []; - const readQueue: { - resolve: (event: ChatCompletionReadableStreamItem | undefined) => void; - reject: (err: unknown) => void; - }[] = []; - let done = false; let lastChunk: ChatCompletionChunk | undefined; let toolCallIds: string[] | undefined; - const pushEvent = (event: ChatCompletionReadableStreamItem) => { - const reader = readQueue.shift(); - if (reader) { - reader.resolve(event); - } else { - pushQueue.push(event); - } - }; - - this.on('chunk', (chunk) => { - lastChunk = chunk; - pushEvent(chunk); - }); - this.on('message', (message: ChatCompletionMessageParam) => { - if (isAssistantMessage(message)) { - toolCallIds = message.tool_calls?.map((toolCall) => toolCall.id); - return; - } - - if (isToolMessage(message)) { - if (!lastChunk) { - throw new OpenAIError('cannot serialize a tool message before receiving any chunks'); - } - pushEvent(makeChatCompletionReadableStreamMessageChunk(lastChunk, message, toolCallIds)); - } - }); - - this.on('end', () => { - done = true; - for (const reader of readQueue) { - reader.resolve(undefined); - } - readQueue.length = 0; - }); - - this.on('abort', (err) => { - done = true; - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }); - - this.on('error', (err) => { - done = true; - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }); + const iterator = this._createIterator( + (push) => { + const onChunk = (chunk: ChatCompletionChunk) => { + lastChunk = chunk; + push(chunk); + }; + const onMessage = (message: ChatCompletionMessageParam) => { + if (isAssistantMessage(message)) { + toolCallIds = message.tool_calls?.map((toolCall) => toolCall.id); + return; + } - const iterator = (): AsyncIterator => ({ - next: async (): Promise> => { - if (!pushQueue.length) { - if (done) { - return { value: undefined, done: true }; + if (isToolMessage(message)) { + if (!lastChunk) { + throw new OpenAIError('cannot serialize a tool message before receiving any chunks'); + } + push(makeChatCompletionReadableStreamMessageChunk(lastChunk, message, toolCallIds)); } - return new Promise((resolve, reject) => - readQueue.push({ resolve, reject }), - ).then((event) => (event ? { value: event, done: false } : { value: undefined, done: true })); - } - const event = pushQueue.shift(); - if (!event) { - return { value: undefined, done: true }; - } - return { value: event, done: false }; + }; + this.on('chunk', onChunk); + this.on('message', onMessage); + return () => { + this.off('chunk', onChunk); + this.off('message', onMessage); + }; }, - return: async () => { - this.abort(); - return { value: undefined, done: true }; - }, - }); + { onReturn: () => this.abort() }, + ); - const stream = new Stream(iterator, this.controller); + const stream = new Stream(() => iterator, this.controller); return stream.toReadableStream(); } diff --git a/src/lib/EventStream.ts b/src/lib/EventStream.ts index 77fe2ffd26..f31f93255e 100644 --- a/src/lib/EventStream.ts +++ b/src/lib/EventStream.ts @@ -180,17 +180,50 @@ export class EventStream { event: Event, ): AsyncIterableIterator> { type Parameters = EventParameters; - type Result = IteratorResult; + return this._createIterator( + (push) => { + const onEvent = (...args: Parameters) => push(args); + this.on(event, onEvent as EventListener); + return () => this.off(event, onEvent as EventListener); + }, + { + // When iterating the 'error' or 'abort' event itself, yield it as a + // value instead of rejecting the iterator. + rejectOnError: event !== 'error', + rejectOnAbort: event !== 'abort', + }, + ); + } + + /** + * Shared buffered async-iterator adapter over this stream's events. + * + * `attach` registers the producer listener(s) with the given `push` and + * returns a cleanup function that removes them. Termination is handled + * here: the iterator ends when the stream ends, listeners are removed on + * end/return, and a terminal error is retained until buffered values have + * drained so it is surfaced even when no reader was waiting when it fired. + */ + protected _createIterator( + attach: (push: (value: T) => void) => () => void, + { + rejectOnError = true, + rejectOnAbort = true, + onReturn, + }: { rejectOnError?: boolean; rejectOnAbort?: boolean; onReturn?: () => void } = {}, + ): AsyncIterableIterator { + type Result = IteratorResult; type Reader = { resolve: (result: Result) => void; reject: (error: OpenAIError) => void; }; - const pushQueue: Parameters[] = []; + const pushQueue: T[] = []; const readQueue: Reader[] = []; let ended = this.ended; let failure: OpenAIError | undefined; let failureDelivered = false; + let detach: () => void = () => {}; const doneResult = (): Result => ({ value: undefined as never, done: true }); const finishReaders = () => { @@ -204,18 +237,18 @@ export class EventStream { readQueue.shift()!.reject(failure); }; const cleanup = () => { - this.off(event, onEvent as EventListener); + detach(); this.off('end', onEnd); - if (event !== 'error') this.off('error', onFailure); - if (event !== 'abort') this.off('abort', onFailure); + if (rejectOnError) this.off('error', onFailure); + if (rejectOnAbort) this.off('abort', onFailure); }; - const onEvent = (...args: Parameters) => { + const push = (value: T) => { if (ended) return; const reader = readQueue.shift(); if (reader) { - reader.resolve({ value: args, done: false }); + reader.resolve({ value, done: false }); } else { - pushQueue.push(args); + pushQueue.push(value); } }; const onFailure = (error: OpenAIError) => { @@ -232,16 +265,17 @@ export class EventStream { }; if (!ended) { - this.on(event, onEvent as EventListener); + detach = attach(push); this.on('end', onEnd); - if (event !== 'error') this.on('error', onFailure); - if (event !== 'abort') this.on('abort', onFailure); + if (rejectOnError) this.on('error', onFailure); + if (rejectOnAbort) this.on('abort', onFailure); } return { - next: () => { - const value = pushQueue.shift(); - if (value) return Promise.resolve({ value, done: false }); + next: (): Promise => { + if (pushQueue.length) { + return Promise.resolve({ value: pushQueue.shift()!, done: false }); + } if (failure && !failureDelivered) { failureDelivered = true; @@ -259,6 +293,14 @@ export class EventStream { pushQueue.length = 0; cleanup(); finishReaders(); + if (onReturn) { + // The consumer explicitly ended iteration, so any failure the + // onReturn callback triggers (e.g. aborting the stream) is + // self-inflicted; mark the stream's terminal promise as handled so + // it does not surface as an unhandled rejection. + this.done().catch(() => {}); + onReturn(); + } return Promise.resolve(doneResult()); }, [Symbol.asyncIterator]() { diff --git a/src/lib/responses/ResponseStream.ts b/src/lib/responses/ResponseStream.ts index 130f163489..fb1623a043 100644 --- a/src/lib/responses/ResponseStream.ts +++ b/src/lib/responses/ResponseStream.ts @@ -227,71 +227,14 @@ export class ResponseStream } [Symbol.asyncIterator](this: ResponseStream): AsyncIterator { - const pushQueue: ResponseStreamEvent[] = []; - const readQueue: { - resolve: (event: ResponseStreamEvent | undefined) => void; - reject: (err: unknown) => void; - }[] = []; - let done = false; - // Capture a terminal error so it can still be surfaced after any buffered - // events drain, even when no reader was waiting at the time it fired. - let failure: unknown = null; - let failureDelivered = false; - - const rejectQueuedReaders = (err: unknown) => { - done = true; - failure = err; - if (readQueue.length) { - failureDelivered = true; - } - for (const reader of readQueue) { - reader.reject(err); - } - readQueue.length = 0; - }; - - this.on('event', (event) => { - const reader = readQueue.shift(); - if (reader) { - reader.resolve(event); - } else { - pushQueue.push(event); - } - }); - - this.on('end', () => { - done = true; - for (const reader of readQueue) { - reader.resolve(undefined); - } - readQueue.length = 0; - }); - - this.on('abort', rejectQueuedReaders); - this.on('error', rejectQueuedReaders); - - return { - next: async (): Promise> => { - if (!pushQueue.length) { - if (failure !== null && !failureDelivered) { - failureDelivered = true; - return Promise.reject(failure); - } - if (done) { - return { value: undefined, done: true }; - } - return new Promise((resolve, reject) => - readQueue.push({ resolve, reject }), - ).then((event) => (event ? { value: event, done: false } : { value: undefined, done: true })); - } - const event = pushQueue.shift()!; - return { value: event, done: false }; - }, - return: async () => { - this.abort(); - return { value: undefined, done: true }; + return this._createIterator( + (push) => { + const onEvent = (event: ResponseStreamEvent) => push(event); + this.on('event', onEvent); + return () => this.off('event', onEvent); }, - }; + { onReturn: () => this.abort() }, + ); } /** diff --git a/tests/lib/ChatCompletionStream.test.ts b/tests/lib/ChatCompletionStream.test.ts index 7af52005be..389a2f1ba1 100644 --- a/tests/lib/ChatCompletionStream.test.ts +++ b/tests/lib/ChatCompletionStream.test.ts @@ -1,5 +1,5 @@ import { vi } from 'vitest'; -import OpenAI from 'openai'; +import OpenAI, { OpenAIError } from 'openai'; import { zodResponseFormat } from 'openai/helpers/zod'; import { ChatCompletionStream } from 'openai/lib/ChatCompletionStream'; import { ChatCompletionStreamingRunner } from 'openai/lib/ChatCompletionStreamingRunner'; @@ -783,7 +783,9 @@ describe('.stream()', () => { // buffer in the iterator's internal queue instead of rejecting a pending // reader. const iterator = stream[Symbol.asyncIterator](); - await new Promise((resolve) => setTimeout(resolve, 10)); + // Wait for the stream's terminal signal so the chunks and the error have + // definitely been emitted before we start reading. + await stream.done().catch(() => {}); const collected: OpenAI.Chat.ChatCompletionChunk[] = []; let caught: unknown = null; @@ -797,6 +799,138 @@ describe('.stream()', () => { } expect(collected).toHaveLength(chunks.length); - expect(caught).toBeInstanceOf(Error); + expect(caught).toBeInstanceOf(OpenAIError); + expect((caught as OpenAIError).message).toBe('network boom'); + }); + + it('rejects a pending read exactly once when the stream errors while a reader is waiting', async () => { + const chunks = [ + { + id: 'chatcmpl-test', + object: 'chat.completion.chunk', + created: 1, + model: 'gpt-4', + choices: [{ index: 0, delta: { role: 'assistant', content: 'hel' }, finish_reason: null }], + }, + { + id: 'chatcmpl-test', + object: 'chat.completion.chunk', + created: 1, + model: 'gpt-4', + choices: [{ index: 0, delta: { content: 'lo' }, finish_reason: null }], + }, + ] as unknown as OpenAI.Chat.ChatCompletionChunk[]; + const readable = new Stream(async function* () { + for (const chunk of chunks) yield chunk; + throw new Error('network boom'); + }, new AbortController()).toReadableStream(); + + const stream = ChatCompletionStream.fromReadableStream(readable); + // Consume eagerly so each read is awaiting when its chunk (and finally the + // error) arrives, exercising the pending-reader path rather than the + // buffered path. + const iterator = stream[Symbol.asyncIterator](); + + await expect(iterator.next()).resolves.toMatchObject({ done: false }); + await expect(iterator.next()).resolves.toMatchObject({ done: false }); + const caught = await iterator.next().then( + () => null, + (err) => err, + ); + expect(caught).toBeInstanceOf(OpenAIError); + expect((caught as OpenAIError).message).toBe('network boom'); + // The failure is delivered exactly once; iteration then ends cleanly. + await expect(iterator.next()).resolves.toEqual({ value: undefined, done: true }); + }); + + it('aborts the stream when the consumer breaks out of iteration', async () => { + const readable = new Stream(async function* (): AsyncGenerator { + yield { + id: 'chatcmpl-test', + object: 'chat.completion.chunk', + created: 1, + model: 'gpt-4', + choices: [{ index: 0, delta: { role: 'assistant', content: 'hel' }, finish_reason: null }], + } as unknown as OpenAI.Chat.ChatCompletionChunk; + // Hang so the only way the consumer stops is by breaking out. + await new Promise(() => {}); + }, new AbortController()).toReadableStream(); + + const stream = ChatCompletionStream.fromReadableStream(readable); + for await (const chunk of stream) { + void chunk; + break; + } + + expect(stream.controller.signal.aborted).toBe(true); + }); + + it('returns done immediately when iterating after the stream has ended', async () => { + const chunks = [ + { + id: 'chatcmpl-test', + object: 'chat.completion.chunk', + created: 1, + model: 'gpt-4', + choices: [{ index: 0, delta: { role: 'assistant', content: 'hello' }, finish_reason: 'stop' }], + }, + ] as unknown as OpenAI.Chat.ChatCompletionChunk[]; + const readable = new Stream(async function* () { + for (const chunk of chunks) yield chunk; + }, new AbortController()).toReadableStream(); + + const stream = ChatCompletionStream.fromReadableStream(readable); + await stream.done(); + + const iterator = stream[Symbol.asyncIterator](); + await expect(iterator.next()).resolves.toEqual({ value: undefined, done: true }); + }); + + it('toReadableStream surfaces a mid-stream error when items are buffered before consumption', async () => { + const chunks = [ + { + id: 'chatcmpl-test', + object: 'chat.completion.chunk', + created: 1, + model: 'gpt-4', + choices: [{ index: 0, delta: { role: 'assistant', content: 'hel' }, finish_reason: null }], + }, + { + id: 'chatcmpl-test', + object: 'chat.completion.chunk', + created: 1, + model: 'gpt-4', + choices: [{ index: 0, delta: { content: 'lo' }, finish_reason: null }], + }, + ] as unknown as OpenAI.Chat.ChatCompletionChunk[]; + const readable = new Stream(async function* () { + for (const chunk of chunks) yield chunk; + throw new Error('network boom'); + }, new AbortController()).toReadableStream(); + + const runner = ChatCompletionStreamingRunner.fromReadableStream(readable); + // Bridge to a ReadableStream immediately (registering its listeners) but + // do not read from it until the runner has already errored, so the chunks + // and the error land while nothing is pulling: they buffer in the + // adapter's internal queue instead of rejecting a pending reader. + const proxied = Stream.fromReadableStream( + runner.toReadableStream(), + new AbortController(), + ); + await runner.done().catch(() => {}); + + const collected: OpenAI.Chat.ChatCompletionChunk[] = []; + let caught: unknown = null; + try { + for await (const chunk of proxied) { + collected.push(chunk); + } + } catch (err) { + caught = err; + } + + expect(collected).toHaveLength(chunks.length); + expect(caught).toBeInstanceOf(OpenAIError); + expect((caught as OpenAIError).message).toBe('network boom'); }); }); diff --git a/tests/lib/EventStream.test.ts b/tests/lib/EventStream.test.ts index d7d8a424af..fc3d7d5ee7 100644 --- a/tests/lib/EventStream.test.ts +++ b/tests/lib/EventStream.test.ts @@ -1,5 +1,5 @@ import { vi } from 'vitest'; -import { OpenAIError } from 'openai/error'; +import { APIUserAbortError, OpenAIError } from 'openai/error'; import { type BaseEvents, EventStream } from 'openai/lib/EventStream'; interface TestEvents extends BaseEvents { @@ -15,6 +15,10 @@ class TestStream extends EventStream { this._emit('error', error); } + emitAbort(error: APIUserAbortError) { + this._emit('abort', error); + } + end() { this._emit('end'); } @@ -67,6 +71,30 @@ describe('EventStream.events', () => { await expect(iterator.next()).resolves.toEqual({ value: undefined, done: true }); }); + test('drains queued events before rejecting on abort', async () => { + const stream = new TestStream(); + const iterator = stream.events('foo'); + const error = new APIUserAbortError(); + + stream.emitFoo('first', 1); + stream.emitAbort(error); + + await expect(iterator.next()).resolves.toEqual({ value: ['first', 1], done: false }); + await expect(iterator.next()).rejects.toBe(error); + await expect(iterator.next()).resolves.toEqual({ value: undefined, done: true }); + }); + + test("yields the 'error' event as a value instead of rejecting when iterating it", async () => { + const stream = new TestStream(); + const iterator = stream.events('error'); + const error = new OpenAIError('oops'); + + stream.emitError(error); + + await expect(iterator.next()).resolves.toEqual({ value: [error], done: false }); + await expect(iterator.next()).resolves.toEqual({ value: undefined, done: true }); + }); + test('does not suppress errors after iterator cleanup', async () => { const stream = new TestStream(); const iterator = stream.events('foo'); diff --git a/tests/lib/ResponseStream.test.ts b/tests/lib/ResponseStream.test.ts index 31a915cc24..716968b14f 100644 --- a/tests/lib/ResponseStream.test.ts +++ b/tests/lib/ResponseStream.test.ts @@ -1,5 +1,5 @@ import { vi } from 'vitest'; -import OpenAI, { APIUserAbortError } from 'openai'; +import OpenAI, { APIUserAbortError, OpenAIError } from 'openai'; import { ReadableStreamFrom } from 'openai/internal/shims'; import { ResponseStream } from 'openai/lib/responses/ResponseStream'; import type { Response, ResponseStreamEvent } from 'openai/resources/responses/responses'; @@ -285,8 +285,9 @@ describe('.stream()', () => { // buffer in the iterator's internal queue instead of rejecting a pending // reader. const iterator = stream[Symbol.asyncIterator](); - // Let the producer drain the readable and hit the error before we read. - await new Promise((resolve) => setTimeout(resolve, 10)); + // Wait for the stream's terminal signal so the events and the error have + // definitely been emitted before we start reading. + await stream.done().catch(() => {}); const collected: ResponseStreamEvent[] = []; let caught: unknown = null; @@ -300,7 +301,8 @@ describe('.stream()', () => { } expect(collected).toHaveLength(validEvents.length); - expect(caught).toBeInstanceOf(Error); + expect(caught).toBeInstanceOf(OpenAIError); + expect((caught as OpenAIError).message).toBe('missing output at index 99'); }); }); diff --git a/tests/streaming/assistants/assistant.test.ts b/tests/streaming/assistants/assistant.test.ts index 8d03ba64b0..72a6cca1da 100644 --- a/tests/streaming/assistants/assistant.test.ts +++ b/tests/streaming/assistants/assistant.test.ts @@ -1,6 +1,7 @@ -import OpenAI from 'openai'; +import OpenAI, { OpenAIError } from 'openai'; import { ReadableStreamFrom } from 'openai/internal/shims'; import { AssistantStream } from 'openai/lib/AssistantStream'; +import { AssistantStreamEvent } from 'openai/resources/beta/assistants'; import { Stream } from 'openai/streaming'; const openai = new OpenAI({ @@ -96,4 +97,60 @@ describe('assistant tests', () => { expect(deltas).toEqual(['E', 'ddy']); }); + + test('surfaces a mid-stream error when events are buffered before consumption', async () => { + const encoder = new TextEncoder(); + const events = [ + { + event: 'thread.message.created', + data: { + id: 'msg_1', + content: [], + }, + }, + { + event: 'thread.message.delta', + data: { + id: 'msg_1', + delta: { + content: [{ index: 0, type: 'text', text: { value: 'hi', annotations: [] } }], + }, + }, + }, + ]; + // Yield valid events, then throw to error the stream after they have been + // delivered (mimics a connection drop mid-run). + const input = ReadableStreamFrom( + (async function* () { + for (const event of events) { + yield encoder.encode(JSON.stringify(event) + '\n'); + } + throw new Error('assistant boom'); + })(), + ); + const assistantStream = AssistantStream.fromReadableStream(input); + // Grab the iterator (registering its listeners) but do not consume yet, so + // the valid events and the error land while no reader is waiting: they + // buffer in the iterator's internal queue instead of rejecting a pending + // reader. + const iterator = assistantStream[Symbol.asyncIterator](); + // Wait for the stream's terminal signal so the events and the error have + // definitely been emitted before we start reading. + await assistantStream.done().catch(() => {}); + + const collected: AssistantStreamEvent[] = []; + let caught: unknown = null; + try { + let result: IteratorResult; + while (!(result = await iterator.next()).done) { + collected.push(result.value); + } + } catch (err) { + caught = err; + } + + expect(collected.map((event) => event.event)).toEqual(['thread.message.created', 'thread.message.delta']); + expect(caught).toBeInstanceOf(OpenAIError); + expect((caught as OpenAIError).message).toBe('assistant boom'); + }); });