Skip to content

Implementación SSE de Anthropic: de bytes a AssistantMessage

源码版本v0.73.1

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

  1. Decodificación de byte stream SSE: iterateSseMessages convierte ReadableStream<Uint8Array> en un async generator de ServerSentEvent, soportando \r\n / \n / empalme entre chunks. Ver packages/ai/src/providers/anthropic.ts:321-369.
  2. Decodificación de línea: decodeSseLine trata según la spec SSE event: / data: / comentarios : / flush en línea vacía. Ver packages/ai/src/providers/anthropic.ts:266-290.
  3. Eventos a objetos tipados: iterateAnthropicEvents parsea el data del SSE con parseJsonWithRepair a RawMessageStreamEvent, filtra eventos no Anthropic y verifica que message_start / message_stop estén apareados. Ver packages/ai/src/providers/anthropic.ts:380-419.
  4. Entrada principal de streaming: streamAnthropic crea el client, construye los params, invoca client.messages.create({stream:true}) y empuja los eventos al AssistantMessageEventStream. Ver packages/ai/src/providers/anthropic.ts:421-510.
  5. Versión Simple: streamSimpleAnthropic traduce SimpleStreamOptions a AnthropicOptions completos, distinguiendo adaptive thinking (Opus 4.6 / Sonnet 4.6) y thinking basado en budget. Ver packages/ai/src/providers/anthropic.ts:717-756.
  6. Construcción de parámetros: buildParams arma messages / system / tools / thinking; el token OAuth fuerza la inyección de identidad Claude Code. Ver packages/ai/src/providers/anthropic.ts:868-925.
  7. Conversión de mensajes: convertMessages convierte Message[] interno a MessageParam[] de Anthropic, gestionando bloques text/image/tool, cache_control y limpieza de surrogados. Ver packages/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

Núcleo del decodificador SSE: la línea vacía dispara el flush y devuelve el evento completo:

typescript
// 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:

typescript
// 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:

typescript
// 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

Resumen

anthropic.ts ejecuta el SSE a mano y encadena las cuatro conversiones byte stream → ServerSentEventRawMessageStreamEventAssistantMessageEvent; 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.