diff --git a/packages/ai/src/protocols/open-responses.ts b/packages/ai/src/protocols/open-responses.ts index 1dd449319eb..0b7c7b0714b 100644 --- a/packages/ai/src/protocols/open-responses.ts +++ b/packages/ai/src/protocols/open-responses.ts @@ -759,7 +759,7 @@ const TERMINAL_TYPES = new Set(["error", "response.completed", "response.incompl export const terminal = (event: Event) => TERMINAL_TYPES.has(event.type) const onOutputTextDelta = (state: ParserState, event: Event, id: string): StepResult => { - if (!event.delta) return [state, NO_EVENTS] + if (!event.delta || !state.messageItems.has(id)) return [state, NO_EVENTS] const events: LLMEvent[] = [] const phase = state.messagePhases[id] const metadata = providerMetadata(state, { itemId: id, ...(phase === undefined ? {} : { phase }) }) @@ -777,10 +777,9 @@ const onOutputTextDone = (state: ParserState, event: Event, id: string): StepRes } export const onReasoningDelta = (state: ParserState, event: Event, itemID: string): StepResult => { - if (!event.delta) return [state, NO_EVENTS] + if (!event.delta || !state.reasoningItems[itemID]) return [state, NO_EVENTS] const events: LLMEvent[] = [] - const id = - event.summary_index !== undefined || state.reasoningItems[itemID] ? `${itemID}:${event.summary_index ?? 0}` : itemID + const id = `${itemID}:${event.summary_index ?? 0}` return [ { ...state, @@ -858,27 +857,9 @@ const onOutputItemAdded = (state: ParserState, event: Event): StepResult => { const onReasoningSummaryPartAdded = (state: ParserState, event: Event): StepResult => { if (!event.item_id || event.summary_index === undefined) return [state, NO_EVENTS] - const item = state.reasoningItems[event.item_id] ?? { encryptedContent: undefined, summaryParts: {} } - if (event.summary_index === 0) { - if (state.reasoningItems[event.item_id]) return [state, NO_EVENTS] - const events: LLMEvent[] = [] - return [ - { - ...state, - lifecycle: Lifecycle.reasoningStart( - state.lifecycle, - events, - `${event.item_id}:0`, - providerMetadata(state, { itemId: event.item_id, reasoningEncryptedContent: null }), - ), - reasoningItems: { - ...state.reasoningItems, - [event.item_id]: { ...item, summaryParts: { 0: "active" } }, - }, - }, - events, - ] - } + const item = state.reasoningItems[event.item_id] + if (!item) return [state, NO_EVENTS] + if (event.summary_index === 0) return [state, NO_EVENTS] const events: LLMEvent[] = [] const closed = Object.entries(item.summaryParts) @@ -957,7 +938,7 @@ const onFunctionCallArgumentsDelta = Effect.fn("OpenResponses.onFunctionCallArgu state: ParserState, event: Event, ) { - if (!event.item_id || !event.delta) return [state, NO_EVENTS] satisfies StepResult + if (!event.item_id || !event.delta || !state.tools[event.item_id]) return [state, NO_EVENTS] satisfies StepResult const result = ToolStream.appendExisting( state.id, state.tools, @@ -1165,7 +1146,10 @@ export const step = (state: ParserState, event: Event) => { return ProviderShared.eventError(state.id, `${event.type} message is missing id`) return Effect.succeed(onOutputItemAdded(state, event)) } - if (event.type === "response.function_call_arguments.delta") return onFunctionCallArgumentsDelta(state, event) + if (event.type === "response.function_call_arguments.delta") + return event.item_id + ? onFunctionCallArgumentsDelta(state, event) + : ProviderShared.eventError(state.id, `${event.type} is missing item_id`) if (event.type === "response.output_item.done") { if (event.item?.type === "message" && !event.item.id) return ProviderShared.eventError(state.id, `${event.type} message is missing id`) diff --git a/packages/ai/test/provider/google-vertex.test.ts b/packages/ai/test/provider/google-vertex.test.ts index e7ed8a0b09b..1b5da5d7aac 100644 --- a/packages/ai/test/provider/google-vertex.test.ts +++ b/packages/ai/test/provider/google-vertex.test.ts @@ -237,6 +237,7 @@ describe("Google Vertex providers", () => { }) return input.respond( sseEvents( + { type: "response.output_item.added", item: { type: "message", id: "msg_1" } }, { type: "response.output_text.delta", item_id: "msg_1", delta: "Hello." }, { type: "response.completed", response: { id: "resp_1" } }, ), diff --git a/packages/ai/test/provider/openai-responses.test.ts b/packages/ai/test/provider/openai-responses.test.ts index 08735d828f7..cd020775e26 100644 --- a/packages/ai/test/provider/openai-responses.test.ts +++ b/packages/ai/test/provider/openai-responses.test.ts @@ -321,6 +321,10 @@ describe("OpenAI Responses route", () => { }), messages: Stream.fromArray([ ProviderShared.encodeJson({ type: "response.created", response: { id: "resp_ws" } }), + ProviderShared.encodeJson({ + type: "response.output_item.added", + item: { type: "message", id: "msg_1" }, + }), ProviderShared.encodeJson({ type: "response.output_text.delta", item_id: "msg_1", delta: "Hi" }), ProviderShared.encodeJson({ type: "response.completed", response: { id: "resp_ws" } }), ]), @@ -1665,6 +1669,7 @@ describe("OpenAI Responses route", () => { it.effect("parses text and usage stream fixtures", () => Effect.gen(function* () { const body = sseEvents( + { type: "response.output_item.added", item: { type: "message", id: "msg_1" } }, { type: "response.output_text.delta", item_id: "msg_1", delta: "Hello" }, { type: "response.output_text.delta", item_id: "msg_1", delta: "!" }, { @@ -1926,6 +1931,67 @@ describe("OpenAI Responses route", () => { }), ) + it.effect("ignores deltas without a matching output item", () => + Effect.gen(function* () { + const response = yield* LLMClient.generate(request).pipe( + Effect.provide( + fixedResponse( + sseEvents( + { type: "response.output_text.delta", item_id: "msg_missing", delta: "orphaned text" }, + { type: "response.refusal.delta", item_id: "refusal_missing", delta: "orphaned refusal" }, + { + type: "response.reasoning_summary_text.delta", + item_id: "rs_missing", + summary_index: 0, + delta: "orphaned reasoning", + }, + { + type: "response.reasoning_summary_part.added", + item_id: "rs_still_missing", + summary_index: 0, + }, + { + type: "response.reasoning_summary_text.delta", + item_id: "rs_still_missing", + summary_index: 0, + delta: "still orphaned reasoning", + }, + { + type: "response.function_call_arguments.delta", + item_id: "fc_missing", + delta: '{"orphaned":true}', + }, + { type: "response.completed", response: { id: "resp_1" } }, + ), + ), + ), + ) + + expect(response.text).toBe("") + expect(response.message.content).toEqual([]) + expect(response.events.some(LLMEvent.is.toolCall)).toBeFalse() + }), + ) + + it.effect("rejects function argument deltas without the spec-required item id", () => + Effect.gen(function* () { + const error = yield* LLMClient.generate(request).pipe( + Effect.provide( + fixedResponse( + sseEvents( + { type: "response.function_call_arguments.delta", delta: "{}" }, + { type: "response.completed", response: { id: "resp_1" } }, + ), + ), + ), + Effect.flip, + ) + + expect(error.reason._tag).toBe("InvalidProviderOutput") + expect(error.message).toContain("response.function_call_arguments.delta is missing item_id") + }), + ) + it.effect("rejects reasoning events without the spec-required item id", () => Effect.gen(function* () { const events = [ @@ -1981,8 +2047,10 @@ describe("OpenAI Responses route", () => { Effect.provide( fixedResponse( sseEvents( + { type: "response.output_item.added", item: { type: "message", id: "msg_1" } }, { type: "response.output_text.delta", item_id: "msg_1", delta: "First" }, - { type: "response.output_text.done", item_id: "msg_1" }, + { type: "response.output_item.done", item: { type: "message", id: "msg_1" } }, + { type: "response.output_item.added", item: { type: "message", id: "msg_2" } }, { type: "response.output_text.delta", item_id: "msg_2", delta: "Second" }, { type: "response.output_item.done", item: { type: "message", id: "msg_2" } }, { type: "response.completed", response: { id: "resp_1" } }, @@ -1993,10 +2061,10 @@ describe("OpenAI Responses route", () => { expect(response.events.filter((event) => event.type.startsWith("text-"))).toEqual([ { type: "text-start", id: "msg_1", providerMetadata: { openai: { itemId: "msg_1" } } }, - { type: "text-delta", id: "msg_1", text: "First" }, - { type: "text-end", id: "msg_1", providerMetadata: undefined }, + { type: "text-delta", id: "msg_1", text: "First", providerMetadata: undefined }, + { type: "text-end", id: "msg_1", providerMetadata: { openai: { itemId: "msg_1" } } }, { type: "text-start", id: "msg_2", providerMetadata: { openai: { itemId: "msg_2" } } }, - { type: "text-delta", id: "msg_2", text: "Second" }, + { type: "text-delta", id: "msg_2", text: "Second", providerMetadata: undefined }, { type: "text-end", id: "msg_2", providerMetadata: { openai: { itemId: "msg_2" } } }, ]) }), @@ -2005,7 +2073,9 @@ describe("OpenAI Responses route", () => { it.effect("parses reasoning summary stream fixtures", () => Effect.gen(function* () { const body = sseEvents( + { type: "response.output_item.added", item: { type: "reasoning", id: "rs_1" } }, { type: "response.reasoning_summary_text.delta", item_id: "rs_1", delta: "thinking" }, + { type: "response.output_item.added", item: { type: "message", id: "msg_1" } }, { type: "response.output_text.delta", item_id: "msg_1", delta: "Hello" }, { type: "response.reasoning_summary_text.done", item_id: "rs_1" }, { type: "response.completed", response: { id: "resp_1" } }, @@ -2017,18 +2087,22 @@ describe("OpenAI Responses route", () => { expect(response.text).toBe("Hello") expect(response.events).toMatchObject([ { type: "step-start", index: 0 }, - { type: "reasoning-start", id: "rs_1" }, - { type: "reasoning-delta", id: "rs_1", text: "thinking" }, + { type: "reasoning-start", id: "rs_1:0" }, + { type: "reasoning-delta", id: "rs_1:0", text: "thinking" }, { type: "text-start", id: "msg_1" }, { type: "text-delta", id: "msg_1", text: "Hello" }, - { type: "reasoning-end", id: "rs_1" }, + { type: "reasoning-end", id: "rs_1:0" }, { type: "text-end", id: "msg_1" }, { type: "step-finish", index: 0, reason: { normalized: "stop", raw: undefined } }, { type: "finish", reason: { normalized: "stop", raw: undefined } }, ]) expect(response.events.filter((event) => event.type === "finish")).toHaveLength(1) expect(response.message.content).toEqual([ - { type: "reasoning", text: "thinking" }, + { + type: "reasoning", + text: "thinking", + providerMetadata: { openai: { itemId: "rs_1", reasoningEncryptedContent: null } }, + }, { type: "text", text: "Hello", providerMetadata: { openai: { itemId: "msg_1" } } }, ]) }), @@ -2040,6 +2114,7 @@ describe("OpenAI Responses route", () => { Effect.provide( fixedResponse( sseEvents( + { type: "response.output_item.added", item: { type: "reasoning", id: "rs_1" } }, { type: "response.reasoning_summary_text.delta", item_id: "rs_1", delta: "thinking" }, { type: "response.output_item.done", @@ -2059,7 +2134,7 @@ describe("OpenAI Responses route", () => { expect(response.events).toContainEqual( expect.objectContaining({ type: "reasoning-end", - id: "rs_1", + id: "rs_1:0", providerMetadata: { openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" } }, }), ) @@ -2200,6 +2275,7 @@ describe("OpenAI Responses route", () => { }) return input.respond( sseEvents( + { type: "response.output_item.added", item: { type: "message", id: "msg_1" } }, { type: "response.output_text.delta", item_id: "msg_1", delta: "Parser now round-trips reasoning." }, { type: "response.completed", response: { id: "resp_1" } }, ), diff --git a/packages/ai/test/provider/xai-responses.test.ts b/packages/ai/test/provider/xai-responses.test.ts index 4f9d14048be..5defc903623 100644 --- a/packages/ai/test/provider/xai-responses.test.ts +++ b/packages/ai/test/provider/xai-responses.test.ts @@ -30,6 +30,10 @@ describe("xAI Responses route", () => { Effect.provide( fixedResponse( sseEvents( + { + type: "response.output_item.added", + item: { type: "reasoning", id: "reasoning_1" }, + }, { type: "response.reasoning_text.delta", item_id: "reasoning_1", delta: "Considering." }, { type: "response.reasoning_text.done", item_id: "reasoning_1" }, {