Skip to content

Anthropic SSE 実装: バイトから AssistantMessage へ

源码版本v0.73.1

providers/anthropic.ts は 9 社の組み込み provider の中で最も重い一家だ。SSE 増分デコード、Anthropic イベントから AssistantMessageEvent への変換、OAuth/Claude Code 身分の注入、cache_control 戦略、thinking と tool call のブロック増分結合をすべて自前で実装する。Anthropic SDK のストリーミング解析には頼らず、response.body のバイトストリームを直接読んで SSE プロトコルを自分で走らせる。最終産物は AssistantMessageEventStream に push されたイベント列で、消費者は for await で増分を受け取る。

役割

  1. SSE バイトストリームのデコード: iterateSseMessagesReadableStream<Uint8Array>ServerSentEvent の async generator に変換する。\r\n / \n / chunk 跨ぎの結合をサポート。packages/ai/src/providers/anthropic.ts:321-369 参照。
  2. 一行デコード: decodeSseLine は SSE 仕様に従い event: / data: / : コメント / 空行 flush を処理する。packages/ai/src/providers/anthropic.ts:266-290 参照。
  3. イベントから型付きオブジェクトへ: iterateAnthropicEvents は SSE data を parseJsonWithRepairRawMessageStreamEvent に解析する。Anthropic でないイベントをフィルタし、message_start / message_stop のペアを検証する。packages/ai/src/providers/anthropic.ts:380-419 参照。
  4. ストリーミング主入口: streamAnthropic は client を作り、buildParams を呼び、client.messages.create({stream:true}) を呼び、イベントを AssistantMessageEventStream へ push する。packages/ai/src/providers/anthropic.ts:421-510 参照。
  5. Simple 版: streamSimpleAnthropicSimpleStreamOptions を完全な AnthropicOptions に翻訳し、adaptive thinking(Opus 4.6 / Sonnet 4.6)と budget-based thinking を区別する。packages/ai/src/providers/anthropic.ts:717-756 参照。
  6. パラメータ構築: buildParams は messages / system / tools / thinking を組み立て、OAuth token の場合は Claude Code 身分を強制注入する。packages/ai/src/providers/anthropic.ts:868-925 参照。
  7. メッセージ変換: convertMessages は内部の Message[] を Anthropic の MessageParam[] に変換し、text/image/tool ブロック、cache_control、surrogate のクリーンアップを処理する。packages/ai/src/providers/anthropic.ts:978-1070 参照。

設計動機

なぜ @anthropic-ai/sdkclient.messages.stream() をそのまま使わないのか? pi-ai はイベントの形状を自分で制御したいからだ。すべての provider は最終的に同じ AssistantMessageEvent を push し、消費側は for await で一つのコードですべての provider を処理できる。SDK 自前の stream helper が返すのは SDK 自身のイベント型で、さらに変換を挟むとかえって遠回りになる。自分で SSE を走らせれば、解析失敗時に sse.raw をエラー情報に同梱でき、CDN による切断やプロキシによるレスポンス改変を排查しやすい。

なぜ convertMessagesbuildParams を分けるのか? OAuth token と通常の API key で system prompt の構築が違うからだ。buildParams が Claude Code 身分を注入するかどうかを決め、convertMessages は単一メッセージのフォーマットだけを管轄する。両者は各々の層を管轄し、cache 戦略を変えてもメッセージ変換は触らず、メッセージ変換を変えても OAuth ロジックは触らない。

主要ファイル

SSE デコードの核心。空行が flush を引き、完全なイベントを返す:

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;
}

Anthropic イベントイテレータはペア検証失敗時にエラーを投げ、切断されたストリームを黙って飲み込まない:

typescript
// packages/ai/src/providers/anthropic.ts:416-418
	if (sawMessageStart && !sawMessageEnd) {
		throw new Error("Anthropic stream ended before message_stop");
	}

OAuth token は特殊な system prompt を使い、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 } : {}),
			},
		];

データフロー

バイトストリームからイベントストリームまでの全リンク:

境界と失敗

まとめ

anthropic.ts は自前で SSE を走らせ、バイトストリーム → ServerSentEventRawMessageStreamEventAssistantMessageEvent の四段変換を繋ぎ、OAuth 身分・cache_control・thinking 戦略を buildParams / convertMessages で処理する。イベントストリームのデータ構造は EventStream 非同期イテレータ、provider がどうレジストリに登録されるかは Provider 抽象と組み込み を参照。