From 0a78b1122206868cb9a3492eed0895b9fcdd01f0 Mon Sep 17 00:00:00 2001 From: Aiden Cline <63023139+rekram1-node@users.noreply.github.com> Date: Tue, 25 Aug 2026 00:03:08 -0500 Subject: [PATCH] fix(ai): ignore unknown Anthropic stream variants (#44817) --- .../ai/src/protocols/anthropic-messages.ts | 86 ++++++++++++--- .../test/provider/anthropic-messages.test.ts | 102 ++++++++++++++++++ 2 files changed, 176 insertions(+), 12 deletions(-) diff --git a/packages/ai/src/protocols/anthropic-messages.ts b/packages/ai/src/protocols/anthropic-messages.ts index b4ab6cbdefa..07410e08be0 100644 --- a/packages/ai/src/protocols/anthropic-messages.ts +++ b/packages/ai/src/protocols/anthropic-messages.ts @@ -1,5 +1,5 @@ import { Buffer } from "node:buffer" -import { Effect, Schema } from "effect" +import { Effect, Option, Schema } from "effect" import { Tool } from "@opencode-ai/schema/tool" import { Route } from "../route/client.js" import { Auth } from "../route/auth.js" @@ -361,6 +361,8 @@ const AnthropicStreamBlock = Schema.Struct({ tool_use_id: Schema.optional(Schema.String), content: Schema.optional(Schema.Unknown), }) +type AnthropicStreamBlock = Schema.Schema.Type +const decodeAnthropicStreamBlock = Schema.decodeUnknownOption(AnthropicStreamBlock) const AnthropicStreamDelta = Schema.Struct({ type: Schema.optional(Schema.String), @@ -371,13 +373,15 @@ const AnthropicStreamDelta = Schema.Struct({ stop_reason: optionalNull(Schema.String), stop_sequence: optionalNull(Schema.String), }) +type AnthropicStreamDelta = Schema.Schema.Type +const decodeAnthropicStreamDelta = Schema.decodeUnknownOption(AnthropicStreamDelta) 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.Unknown), + delta: Schema.optional(Schema.Unknown), usage: Schema.optional(AnthropicUsage), // `type` and `message` are both required per Anthropic's spec, but // OpenAI-compatible proxies and gateway translations occasionally drop one @@ -1106,7 +1110,7 @@ const SERVER_TOOL_RESULT_NAMES: Record = const isServerToolResultType = (type: string): type is AnthropicServerToolResultType => type in SERVER_TOOL_RESULT_NAMES -const serverToolResultEvent = (block: NonNullable): LLMEvent | undefined => { +const serverToolResultEvent = (block: AnthropicStreamBlock): LLMEvent | undefined => { if (!block.type || !isServerToolResultType(block.type)) return undefined const errorPayload = typeof block.content === "object" && block.content !== null && "type" in block.content @@ -1133,7 +1137,10 @@ const onMessageStart = (state: ParserState, event: AnthropicEvent): StepResult = return [usage ? { ...state, usage: mergeUsage(state.usage, usage) } : state, NO_EVENTS] } -const onContentBlockStart = (state: ParserState, event: AnthropicEvent): StepResult => { +const onContentBlockStart = ( + state: ParserState, + event: AnthropicEvent & { readonly content_block: AnthropicStreamBlock }, +): StepResult => { const block = event.content_block if (!block) return [state, NO_EVENTS] @@ -1224,11 +1231,12 @@ const onContentBlockStart = (state: ParserState, event: AnthropicEvent): StepRes const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(function* ( state: ParserState, - event: AnthropicEvent, + event: AnthropicEvent & { readonly delta: AnthropicStreamDelta }, ) { const delta = event.delta if (delta?.type === "text_delta" && delta.text) { + if (!state.lifecycle.text.has(`text-${event.index ?? 0}`)) return [state, NO_EVENTS] satisfies StepResult const events: LLMEvent[] = [] return [ { ...state, lifecycle: Lifecycle.textDelta(state.lifecycle, events, `text-${event.index ?? 0}`, delta.text) }, @@ -1237,6 +1245,7 @@ const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(f } if (delta?.type === "thinking_delta" && delta.thinking) { + if (!state.lifecycle.reasoning.has(`reasoning-${event.index ?? 0}`)) return [state, NO_EVENTS] satisfies StepResult const events: LLMEvent[] = [] return [ { @@ -1249,6 +1258,7 @@ const onContentBlockDelta = Effect.fn("AnthropicMessages.onContentBlockDelta")(f if (delta?.type === "signature_delta" && delta.signature) { const index = event.index ?? 0 + if (!state.lifecycle.reasoning.has(`reasoning-${index}`)) return [state, NO_EVENTS] satisfies StepResult return [ { ...state, @@ -1301,7 +1311,10 @@ const onContentBlockStop = Effect.fn("AnthropicMessages.onContentBlockStop")(fun return [{ ...state, lifecycle, tools: result.tools, reasoningSignatures }, events] satisfies StepResult }) -const onMessageDelta = (state: ParserState, event: AnthropicEvent): StepResult => { +const onMessageDelta = ( + state: ParserState, + event: AnthropicEvent & { readonly delta?: AnthropicStreamDelta }, +): StepResult => { const usage = mergeUsage(state.usage, mapUsage(event.usage)) return [ { @@ -1356,11 +1369,49 @@ const onError = (event: AnthropicEvent) => }), ) +const isKnownStreamBlockType = (type: string) => + type === "text" || + type === "thinking" || + type === "redacted_thinking" || + type === "tool_use" || + type === "server_tool_use" || + isServerToolResultType(type) + +const isKnownStreamDeltaType = (type: string) => + type === "text_delta" || type === "thinking_delta" || type === "signature_delta" || type === "input_json_delta" + +const invalidStreamEvent = (event: AnthropicEvent) => + Effect.fail( + ProviderShared.eventError( + ADAPTER, + "Invalid anthropic/anthropic-messages stream event", + ProviderShared.encodeJson(event), + ), + ) + const step = (state: ParserState, event: AnthropicEvent) => { + if (!SSE_EVENTS.has(event.type)) return Effect.succeed([state, NO_EVENTS]) + if ( + event.type !== "content_block_start" && + event.content_block !== undefined && + Option.isNone(decodeAnthropicStreamBlock(event.content_block)) + ) + return invalidStreamEvent(event) + if ( + event.type !== "content_block_delta" && + event.delta !== undefined && + Option.isNone(decodeAnthropicStreamDelta(event.delta)) + ) + return invalidStreamEvent(event) if (event.type === "message_start") return Effect.succeed(onMessageStart(state, event)) if (event.type === "content_block_start") { - const block = event.content_block - if (block && (block.type === "tool_use" || block.type === "server_tool_use")) { + if (!ProviderShared.isRecord(event.content_block) || typeof event.content_block.type !== "string") + return invalidStreamEvent(event) + if (!isKnownStreamBlockType(event.content_block.type)) return Effect.succeed([state, NO_EVENTS]) + const decoded = decodeAnthropicStreamBlock(event.content_block) + if (Option.isNone(decoded)) return invalidStreamEvent(event) + const block = decoded.value + if (block.type === "tool_use" || block.type === "server_tool_use") { if (event.index === undefined) return Effect.fail(ProviderShared.eventError(ADAPTER, `Anthropic ${block.type} missing index`)) if (!block.id) @@ -1368,11 +1419,22 @@ const step = (state: ParserState, event: AnthropicEvent) => { ProviderShared.eventError(ADAPTER, `Anthropic tool_use missing id at index ${event.index}`), ) } - return Effect.succeed(onContentBlockStart(state, event)) + return Effect.succeed(onContentBlockStart(state, { ...event, content_block: block })) + } + if (event.type === "content_block_delta") { + if (!ProviderShared.isRecord(event.delta)) return invalidStreamEvent(event) + if (typeof event.delta.type === "string" && !isKnownStreamDeltaType(event.delta.type)) + return Effect.succeed([state, NO_EVENTS]) + const decoded = decodeAnthropicStreamDelta(event.delta) + if (Option.isNone(decoded)) return invalidStreamEvent(event) + return onContentBlockDelta(state, { ...event, delta: decoded.value }) } - if (event.type === "content_block_delta") return onContentBlockDelta(state, event) if (event.type === "content_block_stop") return onContentBlockStop(state, event) - if (event.type === "message_delta") return Effect.succeed(onMessageDelta(state, event)) + if (event.type === "message_delta") { + const decoded = decodeAnthropicStreamDelta(event.delta) + if (Option.isNone(decoded)) return invalidStreamEvent(event) + return Effect.succeed(onMessageDelta(state, { ...event, delta: decoded.value })) + } if (event.type === "message_stop") return onMessageStop(state) if (event.type === "error") return onError(event) return Effect.succeed([state, NO_EVENTS]) diff --git a/packages/ai/test/provider/anthropic-messages.test.ts b/packages/ai/test/provider/anthropic-messages.test.ts index 6a10144c903..ca63e5ad0ee 100644 --- a/packages/ai/test/provider/anthropic-messages.test.ts +++ b/packages/ai/test/provider/anthropic-messages.test.ts @@ -770,6 +770,108 @@ describe("Anthropic Messages route", () => { }), ) + it.effect("ignores unknown content block and delta variants", () => + Effect.gen(function* () { + const response = yield* LLMClient.generate(request).pipe( + Effect.provide( + fixedResponse( + sseEvents( + { type: "message_start", message: { usage: { input_tokens: 5 } } }, + { type: "future_event", content_block: 42, delta: 42 }, + { type: "content_block_start", index: 0, content_block: { type: "future_block", text: 42 } }, + { type: "content_block_delta", index: 0, delta: { text: "ignored" } }, + { type: "content_block_delta", index: 0, delta: { type: "future_delta", text: 42 } }, + { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "hidden" } }, + { type: "content_block_delta", index: 0, delta: { type: "thinking_delta", thinking: "hidden" } }, + { type: "content_block_delta", index: 0, delta: { type: "signature_delta", signature: "hidden" } }, + { 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.message.content).toEqual([{ type: "text", text: "Hello" }]) + expect(response.finishReason).toEqual({ normalized: "stop", raw: "end_turn" }) + }), + ) + + it.effect("rejects malformed recognized content block variants", () => + Effect.gen(function* () { + const error = yield* LLMClient.generate(request).pipe( + Effect.provide( + fixedResponse( + sseEvents( + { type: "message_start", message: { usage: { input_tokens: 5 } } }, + { type: "content_block_start", index: 0, content_block: { type: "text", text: 42 } }, + ), + ), + ), + Effect.flip, + ) + + expect(error.reason).toMatchObject({ + _tag: "InvalidProviderOutput", + message: "Invalid anthropic/anthropic-messages stream event", + }) + }), + ) + + it.effect("rejects malformed recognized content delta variants", () => + Effect.gen(function* () { + const error = yield* LLMClient.generate(request).pipe( + Effect.provide( + fixedResponse( + sseEvents( + { type: "message_start", message: { usage: { input_tokens: 5 } } }, + { type: "content_block_start", index: 0, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: 42 } }, + ), + ), + ), + Effect.flip, + ) + + expect(error.reason).toMatchObject({ + _tag: "InvalidProviderOutput", + message: "Invalid anthropic/anthropic-messages stream event", + }) + }), + ) + + it.effect("rejects malformed payloads on unrelated stream events", () => + Effect.gen(function* () { + const events = [ + { type: "message_start", message: { usage: { input_tokens: 1 } }, delta: 42 }, + { type: "content_block_start", index: 0 }, + { type: "content_block_delta", index: 0 }, + { type: "content_block_stop", index: 0, content_block: { type: "text", text: 42 } }, + { type: "message_delta" }, + { type: "message_delta", delta: { stop_reason: 42 } }, + { type: "message_stop", delta: { text: 42 } }, + { type: "error", error: { type: "overloaded_error", message: "busy" }, content_block: 42 }, + ] + + yield* Effect.forEach(events, (event) => + Effect.gen(function* () { + const error = yield* LLMClient.generate(request).pipe( + Effect.provide(fixedResponse(sseEvents(event))), + Effect.flip, + ) + + expect(error.reason).toMatchObject({ + _tag: "InvalidProviderOutput", + message: "Invalid anthropic/anthropic-messages stream event", + }) + }), + ) + }), + ) + it.effect("rejects malformed recognized SSE events", () => Effect.gen(function* () { const error = yield* LLMClient.generate(request).pipe(