Implementación SSE de Anthropic: de bytes a AssistantMessage
providers/anthropic.ts es el provider más pesado de los 9 integrados: implementa su propio decodificador incremental SSE, convierte eventos Anthropic a AssistantMessageEvent, inyecta identidad OAuth/Claude Code, aplica la estrategia de cache_control y ensambla los bloques incrementales de thinking y tool calls. No depende del helper de streaming del SDK de Anthropic; en cambio lee directamente el byte stream de response.body y ejecuta el protocolo SSE a mano. El producto final es una secuencia de eventos empujada a AssistantMessageEventStream que el consumidor consume con for await.
Responsabilidades
- Decodificación de byte stream SSE:
iterateSseMessagesconvierteReadableStream<Uint8Array>en un async generator deServerSentEvent, soportando\r\n/\n/ empalme entre chunks. Verpackages/ai/src/providers/anthropic.ts:321-369. - Decodificación de línea:
decodeSseLinetrata según la spec SSEevent:/data:/ comentarios:/ flush en línea vacía. Verpackages/ai/src/providers/anthropic.ts:266-290. - Eventos a objetos tipados:
iterateAnthropicEventsparsea el data del SSE conparseJsonWithRepairaRawMessageStreamEvent, filtra eventos no Anthropic y verifica quemessage_start/message_stopestén apareados. Verpackages/ai/src/providers/anthropic.ts:380-419. - Entrada principal de streaming:
streamAnthropiccrea el client, construye los params, invocaclient.messages.create({stream:true})y empuja los eventos alAssistantMessageEventStream. Verpackages/ai/src/providers/anthropic.ts:421-510. - Versión Simple:
streamSimpleAnthropictraduceSimpleStreamOptionsaAnthropicOptionscompletos, distinguiendo adaptive thinking (Opus 4.6 / Sonnet 4.6) y thinking basado en budget. Verpackages/ai/src/providers/anthropic.ts:717-756. - Construcción de parámetros:
buildParamsarma messages / system / tools / thinking; el token OAuth fuerza la inyección de identidad Claude Code. Verpackages/ai/src/providers/anthropic.ts:868-925. - Conversión de mensajes:
convertMessagesconvierteMessage[]interno aMessageParam[]de Anthropic, gestionando bloques text/image/tool, cache_control y limpieza de surrogados. Verpackages/ai/src/providers/anthropic.ts:978-1070.
Motivación de diseño
¿Por qué no usar directamente client.messages.stream() del @anthropic-ai/sdk? Porque pi-ai quiere controlar la forma de los eventos: todos los providers terminan empujando el mismo AssistantMessageEvent, y el lado del consumidor usa un único bloque de código for await para todos los providers. El helper del SDK devuelve eventos con la forma propia del SDK; convertirlos después sería dar un rodeo. Ejecutar SSE a mano además permite, si el parseo falla, meter sse.raw en el mensaje de error, útil para diagnosticar truncado de CDN o proxies que reescriben la respuesta.
¿Por qué separar convertMessages y buildParams? Porque el token OAuth y la API key común construyen el system prompt de formas distintas: buildParams decide si inyectar o no la identidad Claude Code; convertMessages sólo se ocupa del formato de cada mensaje. Cada uno maneja su capa; cambiar la estrategia de cache no toca la conversión de mensajes, y viceversa.
Archivos clave
packages/ai/src/providers/anthropic.ts:421-510— Cuerpo destreamAnthropic, con una IIFE async que crea el client y consume el stream.packages/ai/src/providers/anthropic.ts:487-489—client.messages.create({ ...params, stream: true })+stream.push({ type: "start" }).packages/ai/src/providers/anthropic.ts:494-494—for await (const event of iterateAnthropicEvents(response, signal)), bucle principal de consumo.packages/ai/src/providers/anthropic.ts:321-369—iterateSseMessages, corta el byte stream en líneas.packages/ai/src/providers/anthropic.ts:266-290—decodeSseLine, parseo de línea.packages/ai/src/providers/anthropic.ts:380-419—iterateAnthropicEvents, SSE → eventos tipados + validación de apareo.packages/ai/src/providers/anthropic.ts:717-756—streamSimpleAnthropic, traduce simple → full options.packages/ai/src/providers/anthropic.ts:762-810—createClient, beta header + detección de OAuth.packages/ai/src/providers/anthropic.ts:868-925—buildParams, OAuth inyecta identidad Claude Code.packages/ai/src/providers/anthropic.ts:978-1070—convertMessages, conversión de bloques de mensaje.
Núcleo del decodificador SSE: la línea vacía dispara el flush y devuelve el evento completo:
// packages/ai/src/providers/anthropic.ts:266-290
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;
}El iterador de eventos Anthropic lanza si el apareo falla, para no tragarse silenciosamente un stream truncado:
// packages/ai/src/providers/anthropic.ts:416-418
if (sawMessageStart && !sawMessageEnd) {
throw new Error("Anthropic stream ended before message_stop");
}El token OAuth usa un system prompt especial que fuerza la identidad Claude Code:
// packages/ai/src/providers/anthropic.ts:883-890
if (isOAuthToken) {
params.system = [
{
type: "text",
text: "You are Claude Code, Anthropic's official CLI for Claude.",
...(cacheControl ? { cache_control: cacheControl } : {}),
},
];Flujo de datos
Cadena completa del byte stream al flujo de eventos:
Límites y fallos
- Stream truncado:
sawMessageStart && !sawMessageEndlanzaAnthropic stream ended before message_stop. Verpackages/ai/src/providers/anthropic.ts:416-418. - JSON inválido:
parseJsonWithRepairintenta reparar antes de lanzar; el mensaje de error incluyesse.dataysse.raw. Verpackages/ai/src/providers/anthropic.ts:408-413. - Inyección OAuth:
isOAuthTokendetecta el prefijosk-ant-oaty fuerza el system prompt de identidad Claude Code. Verpackages/ai/src/providers/anthropic.ts:758-760ypackages/ai/src/providers/anthropic.ts:883-897. - Estrategia cache_control:
getCacheControl(model, cacheRetention)decide en qué bloques de system / messages aplicarlo. Verpackages/ai/src/providers/anthropic.ts:54-110. - Detección adaptive thinking:
supportsAdaptiveThinking(model.id)distingue Opus 4.6 / Sonnet 4.6 (usan effort) de los modelos más antiguos (usan budget tokens). Verpackages/ai/src/providers/anthropic.ts:681-715. - Temperature y thinking son mutuamente excluyentes: cuando
thinkingEnabled, se saltaparams.temperature. Verpackages/ai/src/providers/anthropic.ts:910-912.
Resumen
anthropic.ts ejecuta el SSE a mano y encadena las cuatro conversiones byte stream → ServerSentEvent → RawMessageStreamEvent → AssistantMessageEvent; la identidad OAuth, cache_control y la estrategia de thinking se manejan en buildParams / convertMessages. La estructura de datos del stream de eventos en EventStream async iterator; cómo se registra el provider en abstracción de provider y built-ins.