Revert "fix(ai): own Anthropic SSE parsing to avoid SDK JSON.parse hard-failures"

This reverts commit 4b926a30a2.
This commit is contained in:
Mario Zechner
2026-04-21 23:19:25 +02:00
parent 81de426f96
commit fc9220d2de
12 changed files with 124 additions and 430 deletions

View File

@@ -4,7 +4,6 @@ import type {
ContentBlockParam,
MessageCreateParamsStreaming,
MessageParam,
RawMessageStreamEvent,
} from "@anthropic-ai/sdk/resources/messages.js";
import { getEnvApiKey } from "../env-api-keys.js";
import { calculateCost } from "../models.js";
@@ -28,7 +27,7 @@ import type {
} from "../types.js";
import { AssistantMessageEventStream } from "../utils/event-stream.js";
import { headersToRecord } from "../utils/headers.js";
import { parseJsonWithRepair, parseStreamingJson } from "../utils/json-parse.js";
import { parseStreamingJson } from "../utils/json-parse.js";
import { sanitizeSurrogates } from "../utils/sanitize-unicode.js";
import { buildCopilotDynamicHeaders, hasCopilotVisionInput } from "./github-copilot-headers.js";
@@ -214,176 +213,6 @@ function mergeHeaders(...headerSources: (Record<string, string> | undefined)[]):
return merged;
}
interface ServerSentEvent {
event: string | null;
data: string;
raw: string[];
}
interface SseDecoderState {
event: string | null;
data: string[];
raw: string[];
}
function flushSseEvent(state: SseDecoderState): ServerSentEvent | null {
if (!state.event && state.data.length === 0) {
return null;
}
const event: ServerSentEvent = {
event: state.event,
data: state.data.join("\n"),
raw: [...state.raw],
};
state.event = null;
state.data = [];
state.raw = [];
return event;
}
function decodeSseLine(line: string, state: SseDecoderState): ServerSentEvent | null {
if (line === "") {
return flushSseEvent(state);
}
state.raw.push(line);
if (line.startsWith(":")) {
return null;
}
const delimiterIndex = line.indexOf(":");
const fieldName = delimiterIndex === -1 ? line : line.slice(0, delimiterIndex);
let value = delimiterIndex === -1 ? "" : line.slice(delimiterIndex + 1);
if (value.startsWith(" ")) {
value = value.slice(1);
}
if (fieldName === "event") {
state.event = value;
} else if (fieldName === "data") {
state.data.push(value);
}
return null;
}
function nextLineBreakIndex(text: string): number {
const carriageReturnIndex = text.indexOf("\r");
const newlineIndex = text.indexOf("\n");
if (carriageReturnIndex === -1) {
return newlineIndex;
}
if (newlineIndex === -1) {
return carriageReturnIndex;
}
return Math.min(carriageReturnIndex, newlineIndex);
}
function consumeLine(text: string): { line: string; rest: string } | null {
const lineBreakIndex = nextLineBreakIndex(text);
if (lineBreakIndex === -1) {
return null;
}
let nextIndex = lineBreakIndex + 1;
if (text[lineBreakIndex] === "\r" && text[nextIndex] === "\n") {
nextIndex += 1;
}
return {
line: text.slice(0, lineBreakIndex),
rest: text.slice(nextIndex),
};
}
async function* iterateSseMessages(
body: ReadableStream<Uint8Array>,
signal?: AbortSignal,
): AsyncGenerator<ServerSentEvent> {
const reader = body.getReader();
const decoder = new TextDecoder();
const state: SseDecoderState = { event: null, data: [], raw: [] };
let buffer = "";
try {
while (true) {
if (signal?.aborted) {
throw new Error("Request was aborted");
}
const { value, done } = await reader.read();
if (done) {
break;
}
buffer += decoder.decode(value, { stream: true });
let consumed = consumeLine(buffer);
while (consumed) {
buffer = consumed.rest;
const event = decodeSseLine(consumed.line, state);
if (event) {
yield event;
}
consumed = consumeLine(buffer);
}
}
buffer += decoder.decode();
let consumed = consumeLine(buffer);
while (consumed) {
buffer = consumed.rest;
const event = decodeSseLine(consumed.line, state);
if (event) {
yield event;
}
consumed = consumeLine(buffer);
}
if (buffer.length > 0) {
const event = decodeSseLine(buffer, state);
if (event) {
yield event;
}
}
const trailingEvent = flushSseEvent(state);
if (trailingEvent) {
yield trailingEvent;
}
} finally {
reader.releaseLock();
}
}
async function* iterateAnthropicEvents(
response: Response,
signal?: AbortSignal,
): AsyncGenerator<RawMessageStreamEvent> {
if (!response.body) {
throw new Error("Attempted to iterate over an Anthropic response with no body");
}
for await (const sse of iterateSseMessages(response.body, signal)) {
if (!sse.event || sse.event === "ping") {
continue;
}
if (sse.event === "error") {
throw new Error(sse.data);
}
try {
yield parseJsonWithRepair<RawMessageStreamEvent>(sse.data);
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
throw new Error(
`Could not parse Anthropic SSE event ${sse.event}: ${message}; data=${sse.data}; raw=${sse.raw.join("\\n")}`,
);
}
}
}
export const streamAnthropic: StreamFunction<"anthropic-messages", AnthropicOptions> = (
model: Model<"anthropic-messages">,
context: Context,
@@ -444,16 +273,16 @@ export const streamAnthropic: StreamFunction<"anthropic-messages", AnthropicOpti
if (nextParams !== undefined) {
params = nextParams as MessageCreateParamsStreaming;
}
const response = await client.messages
.create({ ...params, stream: true }, { signal: options?.signal })
.asResponse();
const { data: anthropicStream, response } = await client.messages
.stream({ ...params, stream: true }, { signal: options?.signal })
.withResponse();
await options?.onResponse?.({ status: response.status, headers: headersToRecord(response.headers) }, model);
stream.push({ type: "start", partial: output });
type Block = (ThinkingContent | TextContent | (ToolCall & { partialJson: string })) & { index: number };
const blocks = output.content as Block[];
for await (const event of iterateAnthropicEvents(response, options?.signal)) {
for await (const event of anthropicStream) {
if (event.type === "message_start") {
output.responseId = event.message.id;
// Capture initial token usage from message_start event