Skip to content

Commit 41b953d

Browse files
authored
fix(core): request/response events' stream discrimination. (#525)
1 parent 2f8aece commit 41b953d

5 files changed

Lines changed: 70 additions & 58 deletions

File tree

packages/core/src/events/AgenticaRequestEvent.ts

Lines changed: 21 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,25 @@ import type OpenAI from "openai";
33
import type { AgenticaEventBase } from "./AgenticaEventBase";
44
import type { AgenticaEventSource } from "./AgenticaEventSource";
55

6-
export interface AgenticaRequestEvent extends AgenticaEventBase<"request"> {
7-
source: AgenticaEventSource;
8-
body: OpenAI.ChatCompletionCreateParamsStreaming | OpenAI.ChatCompletionCreateParamsNonStreaming;
9-
options?: OpenAI.RequestOptions | undefined;
6+
export type AgenticaRequestEvent =
7+
| AgenticaRequestEvent.Streaming
8+
| AgenticaRequestEvent.NonStreaming;
9+
export namespace AgenticaRequestEvent {
10+
export type Streaming = Base<
11+
true,
12+
OpenAI.ChatCompletionCreateParamsStreaming
13+
>;
14+
15+
export type NonStreaming = Base<
16+
false,
17+
OpenAI.ChatCompletionCreateParamsNonStreaming
18+
>;
19+
20+
interface Base<Stream extends boolean, Body extends object>
21+
extends AgenticaEventBase<"request"> {
22+
source: AgenticaEventSource;
23+
stream: Stream;
24+
body: Body;
25+
options?: OpenAI.RequestOptions | undefined;
26+
}
1027
}

packages/core/src/events/AgenticaResponseEvent.ts

Lines changed: 23 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -3,42 +3,30 @@ import type OpenAI from "openai";
33
import type { AgenticaEventBase } from "./AgenticaEventBase";
44
import type { AgenticaEventSource } from "./AgenticaEventSource";
55

6-
export interface AgenticaResponseEvent extends AgenticaEventBase<"response"> {
7-
request_id: string;
8-
9-
/**
10-
* The source agent of the response.
11-
*/
12-
source: AgenticaEventSource;
13-
14-
/**
15-
* Request body.
16-
*/
17-
body: OpenAI.ChatCompletionCreateParamsStreaming;
18-
19-
/**
20-
* The response data.
21-
*/
22-
response: AgenticaResponseEvent.Response;
6+
export type AgenticaResponseEvent =
7+
| AgenticaResponseEvent.Streaming
8+
| AgenticaResponseEvent.NonStreaming;
9+
export namespace AgenticaResponseEvent {
10+
export type Streaming = Base<
11+
true,
12+
OpenAI.ChatCompletionCreateParamsStreaming,
13+
AsyncGenerator<OpenAI.ChatCompletionChunk, undefined, undefined>
14+
>;
2315

24-
/**
25-
* Options for the request.
26-
*/
27-
options?: OpenAI.RequestOptions | undefined;
16+
export type NonStreaming = Base<
17+
false,
18+
OpenAI.ChatCompletionCreateParamsNonStreaming,
19+
OpenAI.ChatCompletion
20+
>;
2821

29-
/**
30-
* Wait the completion.
31-
*/
32-
join: () => Promise<OpenAI.ChatCompletion>;
33-
}
34-
export namespace AgenticaResponseEvent {
35-
export type Response = StreamResponse | NonStreamResponse;
36-
export interface StreamResponse {
37-
stream: true;
38-
data: AsyncGenerator<OpenAI.ChatCompletionChunk, undefined, undefined>;
39-
}
40-
export interface NonStreamResponse {
41-
stream: false;
42-
data: OpenAI.ChatCompletion;
22+
interface Base<Stream extends boolean, Body extends object, Completion extends object>
23+
extends AgenticaEventBase<"response"> {
24+
source: AgenticaEventSource;
25+
request_id: string;
26+
stream: Stream;
27+
body: Body;
28+
completion: Completion;
29+
options?: OpenAI.RequestOptions | undefined;
30+
join: () => Promise<OpenAI.ChatCompletion>;
4331
}
4432
}

packages/core/src/factory/events.ts

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -308,9 +308,12 @@ export function createDescribeEvent(props: {
308308
/* -----------------------------------------------------------
309309
API REQUESTS
310310
----------------------------------------------------------- */
311-
export function createRequestEvent(props: {
311+
export function createRequestEvent<Stream extends boolean>(props: {
312312
source: AgenticaEventSource;
313-
body: OpenAI.ChatCompletionCreateParamsStreaming | OpenAI.ChatCompletionCreateParamsNonStreaming;
313+
stream: Stream;
314+
body: Stream extends true
315+
? OpenAI.ChatCompletionCreateParamsStreaming
316+
: OpenAI.ChatCompletionCreateParamsNonStreaming;
314317
options?: OpenAI.RequestOptions | undefined;
315318
}): AgenticaRequestEvent {
316319
const id: string = v4();
@@ -320,17 +323,23 @@ export function createRequestEvent(props: {
320323
id,
321324
created_at,
322325
source: props.source,
323-
body: props.body,
326+
stream: props.stream as false,
327+
body: props.body as OpenAI.ChatCompletionCreateParamsNonStreaming,
324328
options: props.options,
325329
};
326330
}
327331

328-
export function createResponseEvent(props: {
332+
export function createResponseEvent<Stream extends boolean>(props: {
329333
request_id: string;
330334
source: AgenticaEventSource;
331-
body: OpenAI.ChatCompletionCreateParamsStreaming;
335+
stream: Stream;
336+
body: Stream extends true
337+
? OpenAI.ChatCompletionCreateParamsStreaming
338+
: OpenAI.ChatCompletionCreateParamsNonStreaming;
332339
options?: OpenAI.RequestOptions | undefined;
333-
response: AgenticaResponseEvent.Response;
340+
completion: Stream extends true
341+
? AsyncGenerator<OpenAI.ChatCompletionChunk, undefined, undefined>
342+
: OpenAI.ChatCompletion;
334343
join: () => Promise<OpenAI.ChatCompletion>;
335344
}): AgenticaResponseEvent {
336345
const id: string = v4();
@@ -341,9 +350,10 @@ export function createResponseEvent(props: {
341350
request_id: props.request_id,
342351
created_at,
343352
source: props.source,
344-
body: props.body,
353+
stream: props.stream as false,
354+
body: props.body as OpenAI.ChatCompletionCreateParamsNonStreaming,
355+
completion: props.completion as OpenAI.ChatCompletion,
345356
options: props.options,
346-
response: props.response,
347357
join: props.join,
348358
};
349359
}

packages/core/src/utils/request.ts

Lines changed: 6 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ export function getChatCompletionFunction(props: {
2929
: { stream: false };
3030
const event: AgenticaRequestEvent = createRequestEvent({
3131
source,
32+
stream: streamOptions.stream,
3233
body: {
3334
...body,
3435
model: props.vendor.model,
@@ -75,11 +76,9 @@ export function getChatCompletionFunction(props: {
7576
type: "response",
7677
request_id: event.id,
7778
source,
78-
response: {
79-
stream: false,
80-
data: completion,
81-
},
82-
body: event.body as OpenAI.ChatCompletionCreateParamsStreaming,
79+
stream: false,
80+
body: event.body as OpenAI.ChatCompletionCreateParamsNonStreaming,
81+
completion,
8382
options: event.options,
8483
join: async () => completion,
8584
created_at: new Date().toISOString(),
@@ -122,11 +121,9 @@ export function getChatCompletionFunction(props: {
122121
type: "response",
123122
request_id: event.id,
124123
source,
125-
response: {
126-
stream: true,
127-
data: streamDefaultReaderToAsyncGenerator(streamForStream.getReader(), props.abortSignal),
128-
},
124+
stream: true,
129125
body: event.body as OpenAI.ChatCompletionCreateParamsStreaming,
126+
completion: streamDefaultReaderToAsyncGenerator(streamForStream.getReader(), props.abortSignal),
130127
options: event.options,
131128
join: async () => {
132129
const chunks = await StreamUtil.readAll(streamForJoin, props.abortSignal);

test/src/features/test_base_streaming.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,12 +44,12 @@ export async function test_base_streaming(): Promise<void | false> {
4444
});
4545

4646
agent.on("response", async (event: AgenticaResponseEvent) => {
47-
if (event.response.stream === false) {
47+
if (event.stream === false) {
4848
throw new Error("Response is not a stream");
4949
}
5050
responseEventFired = true;
5151
// Test the stream
52-
for await (const value of event.response.data) {
52+
for await (const value of event.completion) {
5353
if (value.choices !== undefined && value.choices[0]?.delta?.content !== undefined) {
5454
streamContentPieces.push(value.choices[0].delta.content as string);
5555
}

0 commit comments

Comments
 (0)