Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
37 changes: 22 additions & 15 deletions src/lib/ChatCompletionStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -667,6 +667,22 @@ export class ChatCompletionStream<ParsedT = null>
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();
Expand All @@ -685,25 +701,16 @@ export class ChatCompletionStream<ParsedT = null>
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<IteratorResult<ChatCompletionChunk>> => {
if (!pushQueue.length) {
if (failure !== null && !failureDelivered) {
failureDelivered = true;
return Promise.reject(failure);
}
if (done) {
return { value: undefined, done: true };
}
Expand Down
37 changes: 22 additions & 15 deletions src/lib/responses/ResponseStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -233,6 +233,22 @@ export class ResponseStream<ParsedT = null>
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();
Expand All @@ -251,25 +267,16 @@ export class ResponseStream<ParsedT = null>
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<IteratorResult<ResponseStreamEvent>> => {
if (!pushQueue.length) {
if (failure !== null && !failureDelivered) {
failureDelivered = true;
return Promise.reject(failure);
}
if (done) {
return { value: undefined, done: true };
}
Expand Down
47 changes: 47 additions & 0 deletions tests/lib/ChatCompletionStream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<OpenAI.Chat.ChatCompletionChunk>;
while (!(result = await iterator.next()).done) {
collected.push(result.value);
}
} catch (err) {
caught = err;
}

expect(collected).toHaveLength(chunks.length);
expect(caught).toBeInstanceOf(Error);
});
});
49 changes: 49 additions & 0 deletions tests/lib/ResponseStream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<ResponseStreamEvent>;
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[]) {
Expand Down