Skip to content
Merged
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
53 changes: 39 additions & 14 deletions packages/core/src/aisdk.ts
Original file line number Diff line number Diff line change
Expand Up @@ -697,7 +697,7 @@ function metadataProviderOptions(input: ProviderMetadata | undefined): SharedV3P
}

function streamLanguage(language: LanguageModelV3, options: LanguageModelV3CallOptions, http?: HttpMiddleware) {
const state = { step: 0, toolNames: {} as Record<string, string> }
const state: StreamState = { step: 0, toolNames: {}, open: {} }
return Stream.concat(
Stream.make(LLMEvent.stepStart({ index: state.step })),
Stream.unwrap(
Expand All @@ -723,8 +723,16 @@ function streamLanguage(language: LanguageModelV3, options: LanguageModelV3CallO
)
}

type Fragment = "text" | "reasoning"

type StreamState = {
step: number
toolNames: Record<string, string>
open: Partial<Record<Fragment, string>>
}

function streamPartEvents(
state: { step: number; toolNames: Record<string, string> },
state: StreamState,
event: LanguageModelV3StreamPart,
): Effect.Effect<ReadonlyArray<LLMEvent>, AIError> {
switch (event.type) {
Expand All @@ -736,37 +744,31 @@ function streamPartEvents(
case "tool-approval-request":
return Effect.succeed([])
case "text-start":
return Effect.succeed([
LLMEvent.textStart({ id: event.id, providerMetadata: providerMetadata(event.providerMetadata) }),
])
return Effect.succeed(openFragment(state, "text", event.id, providerMetadata(event.providerMetadata)))
case "text-delta":
return Effect.succeed([
...openFragment(state, "text", event.id),
LLMEvent.textDelta({
id: event.id,
text: event.delta,
providerMetadata: providerMetadata(event.providerMetadata),
}),
])
case "text-end":
return Effect.succeed([
LLMEvent.textEnd({ id: event.id, providerMetadata: providerMetadata(event.providerMetadata) }),
])
return Effect.succeed(closeFragment(state, "text", event.id, providerMetadata(event.providerMetadata)))
case "reasoning-start":
return Effect.succeed([
LLMEvent.reasoningStart({ id: event.id, providerMetadata: providerMetadata(event.providerMetadata) }),
])
return Effect.succeed(openFragment(state, "reasoning", event.id, providerMetadata(event.providerMetadata)))
case "reasoning-delta":
return Effect.succeed([
...openFragment(state, "reasoning", event.id),
LLMEvent.reasoningDelta({
id: event.id,
text: event.delta,
providerMetadata: providerMetadata(event.providerMetadata),
}),
])
case "reasoning-end":
return Effect.succeed([
LLMEvent.reasoningEnd({ id: event.id, providerMetadata: providerMetadata(event.providerMetadata) }),
])
return Effect.succeed(closeFragment(state, "reasoning", event.id, providerMetadata(event.providerMetadata)))
case "tool-input-start":
state.toolNames[event.id] = event.toolName
return Effect.succeed([
Expand Down Expand Up @@ -843,6 +845,29 @@ function streamPartEvents(
}
}

// Session persists one open text and one open reasoning fragment at a time, while AI SDK providers may overlap,
// repeat, or omit fragment boundaries. Like the native protocol lifecycles, a start or delta for another fragment
// closes the open one, repeated starts are ignored, and ends for fragments that are not open are dropped.
function openFragment(state: StreamState, kind: Fragment, id: string, providerMetadata?: ProviderMetadata) {
const open = state.open[kind]
if (open === id) return []
state.open[kind] = id
const start =
kind === "text" ? LLMEvent.textStart({ id, providerMetadata }) : LLMEvent.reasoningStart({ id, providerMetadata })
if (open === undefined) return [start]
return [fragmentEnd(kind, open), start]
}

function closeFragment(state: StreamState, kind: Fragment, id: string, providerMetadata?: ProviderMetadata) {
if (state.open[kind] !== id) return []
state.open[kind] = undefined
return [fragmentEnd(kind, id, providerMetadata)]
}

function fragmentEnd(kind: Fragment, id: string, providerMetadata?: ProviderMetadata) {
return kind === "text" ? LLMEvent.textEnd({ id, providerMetadata }) : LLMEvent.reasoningEnd({ id, providerMetadata })
}

function usage(input: Extract<LanguageModelV3StreamPart, { type: "finish" }>["usage"]): UsageInput | undefined {
const output = {
inputTokens: input.inputTokens.total,
Expand Down
84 changes: 84 additions & 0 deletions packages/core/test/aisdk.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -384,6 +384,90 @@ it.effect("routes AI Gateway model options by upstream prefix", () =>
}),
)

it.effect("closes the open AI SDK reasoning part when the next one starts", () =>
Effect.gen(function* () {
// AI SDK OpenAI Responses can start summary part 1 before part 0 ends, then end both at item completion (#50662).
const aisdk = yield* AISDK.Service
yield* aisdk.hook.sdk((event) => {
event.sdk = {
languageModel: () =>
streamModel([
{ type: "reasoning-start", id: "rs_1:0", providerMetadata: { gateway: { generationId: "gen_1" } } },
{ type: "reasoning-start", id: "rs_1:1" },
{ type: "reasoning-delta", id: "rs_1:1", delta: "Second summary" },
{ type: "reasoning-end", id: "rs_1:0", providerMetadata: { gateway: { encrypted: "late" } } },
{ type: "reasoning-end", id: "rs_1:1", providerMetadata: { gateway: { encrypted: "final" } } },
{ type: "finish", finishReason: { unified: "stop", raw: "stop" }, usage },
]),
}
})

const resolved = yield* aisdk.model(model("@ai-sdk/gateway"))
const response = yield* LLMClient.generate(LLM.request({ model: resolved, prompt: "Think" })).pipe(
Effect.provide(client),
)

expect(response.events.filter((event) => event.type.startsWith("reasoning-"))).toEqual([
{ type: "reasoning-start", id: "rs_1:0", providerMetadata: { gateway: { generationId: "gen_1" } } },
{ type: "reasoning-end", id: "rs_1:0" },
{ type: "reasoning-start", id: "rs_1:1", providerMetadata: undefined },
{ type: "reasoning-delta", id: "rs_1:1", text: "Second summary", providerMetadata: undefined },
{ type: "reasoning-end", id: "rs_1:1", providerMetadata: { gateway: { encrypted: "final" } } },
])
}),
)

it.effect("normalizes repeated, reopened, and overlapping AI SDK fragment boundaries", () =>
Effect.gen(function* () {
const aisdk = yield* AISDK.Service
yield* aisdk.hook.sdk((event) => {
event.sdk = {
languageModel: () =>
streamModel([
// Older xAI Responses repeat the start for every summary part.
{ type: "reasoning-start", id: "rs_1" },
{ type: "reasoning-start", id: "rs_1" },
{ type: "reasoning-delta", id: "rs_1", delta: "First" },
{ type: "reasoning-end", id: "rs_1" },
// xAI Chat keeps streaming an ended reasoning id after an empty tool_calls chunk.
{ type: "reasoning-delta", id: "rs_1", delta: "Second" },
{ type: "reasoning-end", id: "rs_1" },
// xAI Responses ends every message item only when the stream flushes.
{ type: "text-start", id: "msg_1" },
{ type: "text-delta", id: "msg_1", delta: "One" },
{ type: "text-start", id: "msg_2" },
{ type: "text-delta", id: "msg_2", delta: "Two" },
{ type: "text-end", id: "msg_1" },
{ type: "text-end", id: "msg_2" },
{ type: "finish", finishReason: { unified: "stop", raw: "stop" }, usage },
]),
}
})

const resolved = yield* aisdk.model(model("@ai-sdk/gateway"))
const response = yield* LLMClient.generate(LLM.request({ model: resolved, prompt: "Think" })).pipe(
Effect.provide(client),
)

expect(
response.events.filter((event) => event.type.startsWith("reasoning-") || event.type.startsWith("text-")),
).toMatchObject([
{ type: "reasoning-start", id: "rs_1" },
{ type: "reasoning-delta", id: "rs_1", text: "First" },
{ type: "reasoning-end", id: "rs_1" },
{ type: "reasoning-start", id: "rs_1" },
{ type: "reasoning-delta", id: "rs_1", text: "Second" },
{ type: "reasoning-end", id: "rs_1" },
{ type: "text-start", id: "msg_1" },
{ type: "text-delta", id: "msg_1", text: "One" },
{ type: "text-end", id: "msg_1" },
{ type: "text-start", id: "msg_2" },
{ type: "text-delta", id: "msg_2", text: "Two" },
{ type: "text-end", id: "msg_2" },
])
}),
)

it.effect("projects replay metadata onto AI SDK prompt parts", () =>
Effect.gen(function* () {
const aisdk = yield* AISDK.Service
Expand Down
Loading