Anthropic SSE 実装: バイトから AssistantMessage へ
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 で増分を受け取る。
役割
- SSE バイトストリームのデコード:
iterateSseMessagesはReadableStream<Uint8Array>をServerSentEventの async generator に変換する。\r\n/\n/ chunk 跨ぎの結合をサポート。packages/ai/src/providers/anthropic.ts:321-369参照。 - 一行デコード:
decodeSseLineは SSE 仕様に従いevent:/data:/:コメント / 空行 flush を処理する。packages/ai/src/providers/anthropic.ts:266-290参照。 - イベントから型付きオブジェクトへ:
iterateAnthropicEventsは SSE data をparseJsonWithRepairでRawMessageStreamEventに解析する。Anthropic でないイベントをフィルタし、message_start/message_stopのペアを検証する。packages/ai/src/providers/anthropic.ts:380-419参照。 - ストリーミング主入口:
streamAnthropicは client を作り、buildParams を呼び、client.messages.create({stream:true})を呼び、イベントをAssistantMessageEventStreamへ push する。packages/ai/src/providers/anthropic.ts:421-510参照。 - Simple 版:
streamSimpleAnthropicはSimpleStreamOptionsを完全なAnthropicOptionsに翻訳し、adaptive thinking(Opus 4.6 / Sonnet 4.6)と budget-based thinking を区別する。packages/ai/src/providers/anthropic.ts:717-756参照。 - パラメータ構築:
buildParamsは messages / system / tools / thinking を組み立て、OAuth token の場合は Claude Code 身分を強制注入する。packages/ai/src/providers/anthropic.ts:868-925参照。 - メッセージ変換:
convertMessagesは内部のMessage[]を Anthropic のMessageParam[]に変換し、text/image/tool ブロック、cache_control、surrogate のクリーンアップを処理する。packages/ai/src/providers/anthropic.ts:978-1070参照。
設計動機
なぜ @anthropic-ai/sdk の client.messages.stream() をそのまま使わないのか? pi-ai はイベントの形状を自分で制御したいからだ。すべての provider は最終的に同じ AssistantMessageEvent を push し、消費側は for await で一つのコードですべての provider を処理できる。SDK 自前の stream helper が返すのは SDK 自身のイベント型で、さらに変換を挟むとかえって遠回りになる。自分で SSE を走らせれば、解析失敗時に sse.raw をエラー情報に同梱でき、CDN による切断やプロキシによるレスポンス改変を排查しやすい。
なぜ convertMessages と buildParams を分けるのか? OAuth token と通常の API key で system prompt の構築が違うからだ。buildParams が Claude Code 身分を注入するかどうかを決め、convertMessages は単一メッセージのフォーマットだけを管轄する。両者は各々の層を管轄し、cache 戦略を変えてもメッセージ変換は触らず、メッセージ変換を変えても OAuth ロジックは触らない。
主要ファイル
packages/ai/src/providers/anthropic.ts:421-510—streamAnthropic本体。async IIFE の中で client 生成 + ストリーミング消費。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))の主消費ループ。packages/ai/src/providers/anthropic.ts:321-369—iterateSseMessagesがバイトストリームを行に切る。packages/ai/src/providers/anthropic.ts:266-290—decodeSseLineの一行解析。packages/ai/src/providers/anthropic.ts:380-419—iterateAnthropicEventsが SSE → 型付きイベント + ペア検証。packages/ai/src/providers/anthropic.ts:717-756—streamSimpleAnthropic。simple → full options の翻訳。packages/ai/src/providers/anthropic.ts:762-810—createClient。beta header の結合 + OAuth 検出。packages/ai/src/providers/anthropic.ts:868-925—buildParams。OAuth で Claude Code 身分を注入。packages/ai/src/providers/anthropic.ts:978-1070—convertMessages。メッセージブロックの変換。
SSE デコードの核心。空行が flush を引き、完全なイベントを返す:
// 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 イベントイテレータはペア検証失敗時にエラーを投げ、切断されたストリームを黙って飲み込まない:
// 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 身分を強制付与する:
// 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 } : {}),
},
];データフロー
バイトストリームからイベントストリームまでの全リンク:
境界と失敗
- ストリームの切断:
sawMessageStart && !sawMessageEndでAnthropic stream ended before message_stopを投げる。packages/ai/src/providers/anthropic.ts:416-418参照。 - JSON 解析失敗:
parseJsonWithRepairがまず修復を試みてから投げる。エラー情報にはsse.dataとsse.rawを含む。packages/ai/src/providers/anthropic.ts:408-413参照。 - OAuth 注入:
isOAuthTokenはsk-ant-oatプレフィックスを検出し、Claude Code 身分の system prompt を強制する。packages/ai/src/providers/anthropic.ts:758-760とpackages/ai/src/providers/anthropic.ts:883-897参照。 - cache_control 戦略:
getCacheControl(model, cacheRetention)が system / messages のどのブロックにつけるかを決める。packages/ai/src/providers/anthropic.ts:54-110参照。 - Adaptive thinking の検出:
supportsAdaptiveThinking(model.id)が Opus 4.6 / Sonnet 4.6(effort を使う)と旧モデル(budget tokens を使う)を区別する。packages/ai/src/providers/anthropic.ts:681-715参照。 - Temperature と thinking の排他:
thinkingEnabledの時はparams.temperatureをスキップする。packages/ai/src/providers/anthropic.ts:910-912参照。
まとめ
anthropic.ts は自前で SSE を走らせ、バイトストリーム → ServerSentEvent → RawMessageStreamEvent → AssistantMessageEvent の四段変換を繋ぎ、OAuth 身分・cache_control・thinking 戦略を buildParams / convertMessages で処理する。イベントストリームのデータ構造は EventStream 非同期イテレータ、provider がどうレジストリに登録されるかは Provider 抽象と組み込み を参照。