From f53e713352dfd0c099f3573f2e5075a96864932f Mon Sep 17 00:00:00 2001 From: Aiden Cline Date: Sun, 23 Aug 2026 20:46:19 -0500 Subject: [PATCH] fix(ai): ignore unknown provider stream parts --- .../ai/src/protocols/anthropic-messages.ts | 80 ++++++++++++++----- packages/ai/src/protocols/gemini.ts | 18 ++++- .../test/provider/anthropic-messages.test.ts | 36 +++++++++ packages/ai/test/provider/gemini.test.ts | 34 ++++++++ 4 files changed, 147 insertions(+), 21 deletions(-) diff --git a/packages/ai/src/protocols/anthropic-messages.ts b/packages/ai/src/protocols/anthropic-messages.ts index 6c22068278a..6673094a108 100644 --- a/packages/ai/src/protocols/anthropic-messages.ts +++ b/packages/ai/src/protocols/anthropic-messages.ts @@ -372,12 +372,33 @@ const AnthropicStreamDelta = Schema.Struct({ stop_sequence: optionalNull(Schema.String), }) +const AnthropicStreamBlockTypes = new Set([ + "tool_use", + "server_tool_use", + "text", + "thinking", + "redacted_thinking", + "web_search_tool_result", + "code_execution_tool_result", + "web_fetch_tool_result", +]) +const AnthropicStreamDeltaTypes = new Set(["text_delta", "thinking_delta", "signature_delta", "input_json_delta"]) + +const AnthropicOpaqueStreamBlock = Schema.declare( + (input): input is Schema.Schema.Type => + ProviderShared.isRecord(input) && typeof input.type === "string" && !AnthropicStreamBlockTypes.has(input.type), +) +const AnthropicOpaqueStreamDelta = Schema.declare( + (input): input is Schema.Schema.Type => + ProviderShared.isRecord(input) && typeof input.type === "string" && !AnthropicStreamDeltaTypes.has(input.type), +) + const AnthropicEvent = Schema.Struct({ type: Schema.String, index: Schema.optional(Schema.Number), message: Schema.optional(Schema.Struct({ usage: Schema.optional(AnthropicUsage) })), - content_block: Schema.optional(AnthropicStreamBlock), - delta: Schema.optional(AnthropicStreamDelta), + content_block: Schema.optional(Schema.Union([AnthropicStreamBlock, AnthropicOpaqueStreamBlock])), + delta: Schema.optional(Schema.Union([AnthropicStreamDelta, AnthropicOpaqueStreamDelta])), usage: Schema.optional(AnthropicUsage), // `type` and `message` are both required per Anthropic's spec, but // OpenAI-compatible proxies and gateway translations occasionally drop one @@ -391,6 +412,7 @@ type AnthropicEvent = Schema.Schema.Type interface ParserState { readonly tools: ToolStream.State + readonly contentBlocks: Readonly> readonly reasoningSignatures: Readonly> readonly usage?: Usage readonly pendingFinish?: { @@ -1101,6 +1123,10 @@ const onMessageStart = (state: ParserState, event: AnthropicEvent): StepResult = const onContentBlockStart = (state: ParserState, event: AnthropicEvent): StepResult => { const block = event.content_block if (!block) return [state, NO_EVENTS] + const index = event.index ?? 0 + const nextState = AnthropicStreamBlockTypes.has(block.type) + ? { ...state, contentBlocks: { ...state.contentBlocks, [index]: true as const } } + : state if (block.type === "tool_use" || block.type === "server_tool_use") { if (event.index === undefined || !block.id) return [state, NO_EVENTS] @@ -1108,7 +1134,7 @@ const onContentBlockStart = (state: ParserState, event: AnthropicEvent): StepRes const lifecycle = Lifecycle.stepStart(state.lifecycle, events) return [ { - ...state, + ...nextState, lifecycle, tools: ToolStream.start(state.tools, event.index, { id: block.id, @@ -1133,23 +1159,26 @@ const onContentBlockStart = (state: ParserState, event: AnthropicEvent): StepRes if (block.type === "text" && block.text !== undefined) { const events: LLMEvent[] = [] - const id = `text-${event.index ?? 0}` + const id = `text-${index}` const lifecycle = Lifecycle.textStart(state.lifecycle, events, id) return [ - { ...state, lifecycle: block.text ? Lifecycle.textDelta(lifecycle, events, id, block.text) : lifecycle }, + { + ...nextState, + lifecycle: block.text ? Lifecycle.textDelta(lifecycle, events, id, block.text) : lifecycle, + }, events, ] } if (block.type === "thinking" && block.thinking !== undefined) { const events: LLMEvent[] = [] - const id = `reasoning-${event.index ?? 0}` + const id = `reasoning-${index}` const providerMetadata = block.signature === undefined ? undefined : anthropicMetadata({ signature: block.signature }) const lifecycle = Lifecycle.reasoningStart(state.lifecycle, events, id, providerMetadata) return [ { - ...state, + ...nextState, lifecycle: block.thinking ? Lifecycle.reasoningDelta(lifecycle, events, id, block.thinking, providerMetadata) : lifecycle, @@ -1169,11 +1198,11 @@ const onContentBlockStart = (state: ParserState, event: AnthropicEvent): StepRes const events: LLMEvent[] = [] return [ { - ...state, + ...nextState, lifecycle: Lifecycle.reasoningStart( state.lifecycle, events, - `reasoning-${event.index ?? 0}`, + `reasoning-${index}`, anthropicMetadata({ redactedData: block.data }), ), }, @@ -1182,9 +1211,9 @@ const onContentBlockStart = (state: ParserState, event: AnthropicEvent): StepRes } const result = serverToolResultEvent(block) - if (!result) return [state, NO_EVENTS] + if (!result) return [nextState, NO_EVENTS] const events: LLMEvent[] = [] - return [{ ...state, lifecycle: Lifecycle.stepStart(state.lifecycle, events) }, [...events, result]] + return [{ ...nextState, lifecycle: Lifecycle.stepStart(state.lifecycle, events) }, [...events, result]] } const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(function* ( @@ -1194,19 +1223,23 @@ const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(f const delta = event.delta if (delta?.type === "text_delta" && delta.text) { + const index = event.index ?? 0 + if (!state.contentBlocks[index]) return [state, NO_EVENTS] satisfies StepResult const events: LLMEvent[] = [] return [ - { ...state, lifecycle: Lifecycle.textDelta(state.lifecycle, events, `text-${event.index ?? 0}`, delta.text) }, + { ...state, lifecycle: Lifecycle.textDelta(state.lifecycle, events, `text-${index}`, delta.text) }, events, ] satisfies StepResult } if (delta?.type === "thinking_delta" && delta.thinking) { + const index = event.index ?? 0 + if (!state.contentBlocks[index]) return [state, NO_EVENTS] satisfies StepResult const events: LLMEvent[] = [] return [ { ...state, - lifecycle: Lifecycle.reasoningDelta(state.lifecycle, events, `reasoning-${event.index ?? 0}`, delta.thinking), + lifecycle: Lifecycle.reasoningDelta(state.lifecycle, events, `reasoning-${index}`, delta.thinking), }, events, ] satisfies StepResult @@ -1214,6 +1247,7 @@ const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(f if (delta?.type === "signature_delta" && delta.signature) { const index = event.index ?? 0 + if (!state.contentBlocks[index]) return [state, NO_EVENTS] satisfies StepResult return [ { ...state, @@ -1251,19 +1285,24 @@ const onContentBlockStop = Effect.fn("AnthropicMessages.onContentBlockStop")(fun const result = yield* ToolStream.finish(ADAPTER, state.tools, event.index) const events: LLMEvent[] = [] const resultEvents = result.events ?? [] + const accepted = state.contentBlocks[event.index] const signature = state.reasoningSignatures[event.index] const lifecycle = resultEvents.length ? Lifecycle.stepStart(state.lifecycle, events) - : Lifecycle.reasoningEnd( - Lifecycle.textEnd(state.lifecycle, events, `text-${event.index}`), - events, - `reasoning-${event.index}`, - signature === undefined ? undefined : anthropicMetadata({ signature }), - ) + : accepted + ? Lifecycle.reasoningEnd( + Lifecycle.textEnd(state.lifecycle, events, `text-${event.index}`), + events, + `reasoning-${event.index}`, + signature === undefined ? undefined : anthropicMetadata({ signature }), + ) + : state.lifecycle events.push(...resultEvents) + const contentBlocks = { ...state.contentBlocks } + delete contentBlocks[event.index] const reasoningSignatures = { ...state.reasoningSignatures } delete reasoningSignatures[event.index] - return [{ ...state, lifecycle, tools: result.tools, reasoningSignatures }, events] satisfies StepResult + return [{ ...state, contentBlocks, lifecycle, tools: result.tools, reasoningSignatures }, events] satisfies StepResult }) const onMessageDelta = (state: ParserState, event: AnthropicEvent): StepResult => { @@ -1361,6 +1400,7 @@ export const protocol = Protocol.make({ event: Protocol.jsonEvent(AnthropicEvent), initial: () => ({ tools: ToolStream.empty(), + contentBlocks: {}, reasoningSignatures: {}, lifecycle: Lifecycle.initial(), }), diff --git a/packages/ai/src/protocols/gemini.ts b/packages/ai/src/protocols/gemini.ts index 3ccf8eafcae..43df1b88e8b 100644 --- a/packages/ai/src/protocols/gemini.ts +++ b/packages/ai/src/protocols/gemini.ts @@ -126,6 +126,21 @@ const GeminiContentPart = Schema.Union([ GeminiFunctionResponsePart, ]) +const GeminiOpaqueResponsePart = Schema.declare( + (input): input is Schema.Schema.Type => + input !== null && + (!ProviderShared.isRecord(input) || + (!("text" in input) && + !("functionCall" in input) && + !("inlineData" in input) && + !("functionResponse" in input))), +) + +const GeminiResponseContent = Schema.Struct({ + role: optionalNull(Schema.Literals(["user", "model"])), + parts: optionalNull(Schema.Array(Schema.Union([GeminiContentPart, GeminiOpaqueResponsePart]))), +}) + const GeminiContent = Schema.Struct({ role: optionalNull(Schema.Literals(["user", "model"])), parts: optionalNull(Schema.Array(GeminiContentPart)), @@ -200,7 +215,7 @@ const GeminiUsage = Schema.Struct({ type GeminiUsage = Schema.Schema.Type const GeminiCandidate = Schema.Struct({ - content: optionalNull(GeminiContent), + content: optionalNull(GeminiResponseContent), finishReason: optionalNull(Schema.String), }) @@ -599,6 +614,7 @@ const step = (state: ParserState, event: GeminiEvent) => { const seenCallIds = new Set(nextState.seenCallIds) for (const part of candidate.content.parts ?? []) { + if (!ProviderShared.isRecord(part)) continue const signature = "thoughtSignature" in part && part.thoughtSignature ? part.thoughtSignature : undefined // Gemini attaches replay signatures to thought parts, visible text, or function calls; // each block kind must retain the signature attached to its own parts. diff --git a/packages/ai/test/provider/anthropic-messages.test.ts b/packages/ai/test/provider/anthropic-messages.test.ts index a2447aea900..d9ae32dd687 100644 --- a/packages/ai/test/provider/anthropic-messages.test.ts +++ b/packages/ai/test/provider/anthropic-messages.test.ts @@ -638,6 +638,42 @@ describe("Anthropic Messages route", () => { }), ) + it.effect("ignores unknown blocks and deltas without accepting their sequence", () => + Effect.gen(function* () { + const response = yield* LLMClient.generate(request).pipe( + Effect.provide( + fixedResponse( + sseEvents( + { type: "message_start", message: { usage: { input_tokens: 1 } } }, + { + type: "content_block_start", + index: 0, + content_block: { type: "future_block", text: 42, thinking: null, id: { value: 1 } }, + }, + { + type: "content_block_delta", + index: 0, + delta: { type: "future_delta", text: 42, partial_json: null }, + }, + { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "ignored" } }, + { type: "content_block_stop", index: 0 }, + { type: "content_block_start", index: 1, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 1, delta: { type: "text_delta", text: "Hello" } }, + { type: "content_block_stop", index: 1 }, + { type: "message_delta", delta: { stop_reason: "end_turn" }, usage: { output_tokens: 1 } }, + { type: "message_stop" }, + ), + ), + ), + ) + + expect(response.text).toBe("Hello") + expect(response.reasoning ?? "").toBe("") + expect(response.toolCalls).toEqual([]) + expect(response.finishReason).toEqual({ normalized: "stop", raw: "end_turn" }) + }), + ) + it.effect("requires message_stop before completing a streamed message", () => Effect.gen(function* () { const error = yield* LLMClient.generate(request).pipe( diff --git a/packages/ai/test/provider/gemini.test.ts b/packages/ai/test/provider/gemini.test.ts index 8aeaf400e20..7765cc5c7cc 100644 --- a/packages/ai/test/provider/gemini.test.ts +++ b/packages/ai/test/provider/gemini.test.ts @@ -906,6 +906,40 @@ describe("Gemini route", () => { }), ) + it.effect("ignores unknown response parts and primitive elements", () => + Effect.gen(function* () { + const response = yield* LLMClient.generate(request).pipe( + Effect.provide( + fixedResponse( + sseEvents({ + candidates: [ + { + content: { parts: [42, "future", { futurePart: { value: 1 } }, { text: "Hello" }] }, + finishReason: "STOP", + }, + ], + }), + ), + ), + ) + + expect(response.text).toBe("Hello") + expect(response.finishReason).toEqual({ normalized: "stop", raw: "STOP" }) + }), + ) + + it.effect("rejects non-array response parts", () => + Effect.gen(function* () { + const error = yield* LLMClient.generate(request).pipe( + Effect.provide(fixedResponse(sseEvents({ candidates: [{ content: { parts: {} } }] }))), + Effect.flip, + ) + + expect(error.reason).toMatchObject({ _tag: "InvalidProviderOutput" }) + expect(error.message).toContain("Invalid google/gemini stream event") + }), + ) + it.effect("preserves thoughtSignature for reasoning and tool-call continuation", () => Effect.gen(function* () { const body = sseEvents({