Skip to content
Open
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
67 changes: 8 additions & 59 deletions src/lib/AssistantStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,66 +94,15 @@ export class AssistantStream
#currentRunStepSnapshot: Runs.RunStep | undefined;

[Symbol.asyncIterator](): AsyncIterator<AssistantStreamEvent> {
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<IteratorResult<AssistantStreamEvent>> => {
if (!pushQueue.length) {
if (done) {
return { value: undefined, done: true };
}
return new Promise<AssistantStreamEvent | undefined>((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<AssistantStreamEvent>(
(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 {
Expand Down
64 changes: 7 additions & 57 deletions src/lib/ChatCompletionStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -661,64 +661,14 @@ export class ChatCompletionStream<ParsedT = null>
}

[Symbol.asyncIterator](this: ChatCompletionStream<ParsedT>): AsyncIterator<ChatCompletionChunk> {
const pushQueue: ChatCompletionChunk[] = [];
const readQueue: {
resolve: (chunk: ChatCompletionChunk | undefined) => void;
reject: (err: unknown) => void;
}[] = [];
let done = false;

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', (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<IteratorResult<ChatCompletionChunk>> => {
if (!pushQueue.length) {
if (done) {
return { value: undefined, done: true };
}
return new Promise<ChatCompletionChunk | undefined>((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<ChatCompletionChunk>(
(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 {
Expand Down
102 changes: 26 additions & 76 deletions src/lib/ChatCompletionStreamingRunner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -69,89 +69,39 @@ export class ChatCompletionStreamingRunner<ParsedT = null>
}

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<ChatCompletionReadableStreamItem>(
(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<ChatCompletionReadableStreamItem> => ({
next: async (): Promise<IteratorResult<ChatCompletionReadableStreamItem>> => {
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<ChatCompletionReadableStreamItem | undefined>((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();
}

Expand Down
70 changes: 56 additions & 14 deletions src/lib/EventStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -180,17 +180,50 @@ export class EventStream<EventTypes extends BaseEvents> {
event: Event,
): AsyncIterableIterator<EventParameters<EventTypes, Event>> {
type Parameters = EventParameters<EventTypes, Event>;
type Result = IteratorResult<Parameters>;
return this._createIterator<Parameters>(
(push) => {
const onEvent = (...args: Parameters) => push(args);
this.on(event, onEvent as EventListener<EventTypes, Event>);
return () => this.off(event, onEvent as EventListener<EventTypes, Event>);
},
{
// 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<T>(
attach: (push: (value: T) => void) => () => void,
{
rejectOnError = true,
rejectOnAbort = true,
onReturn,
}: { rejectOnError?: boolean; rejectOnAbort?: boolean; onReturn?: () => void } = {},
): AsyncIterableIterator<T> {
type Result = IteratorResult<T>;
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 = () => {
Expand All @@ -204,18 +237,18 @@ export class EventStream<EventTypes extends BaseEvents> {
readQueue.shift()!.reject(failure);
};
const cleanup = () => {
this.off(event, onEvent as EventListener<EventTypes, Event>);
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) => {
Expand All @@ -232,16 +265,17 @@ export class EventStream<EventTypes extends BaseEvents> {
};

if (!ended) {
this.on(event, onEvent as EventListener<EventTypes, Event>);
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<Result> => {
if (pushQueue.length) {
return Promise.resolve({ value: pushQueue.shift()!, done: false });
}

if (failure && !failureDelivered) {
failureDelivered = true;
Expand All @@ -259,6 +293,14 @@ export class EventStream<EventTypes extends BaseEvents> {
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]() {
Expand Down
Loading