async def anthropic_adapter(
raw_stream: AsyncIterator[Any],
) -> AsyncIterator[StreamEvent]:
"""Convert raw Anthropic stream chunks to StreamEvent.
- Known chunk types -> typed StreamEvent
- Unknown chunk types -> ExtensionEvent (never raise, never drop)
"""
async for chunk in raw_stream:
event_type = _get(chunk, "type", None)
if event_type == "message_start":
msg = _get(chunk, "message", chunk)
yield MessageStartEvent(
message_id=_get(msg, "id", ""),
model=_get(msg, "model", ""),
)
elif event_type == "content_block_delta":
yield _handle_content_block_delta(chunk, event_type)
elif event_type == "content_block_start":
yield _handle_content_block_start(chunk)
elif event_type == "message_delta":
delta = _get(chunk, "delta", {})
yield MessageEndEvent(
message_id="",
stop_reason=_get(delta, "stop_reason", None),
)
elif event_type == "message_stop":
yield MessageEndEvent(message_id="", stop_reason="end_turn")
elif event_type == "usage" or (
isinstance(chunk, dict) and "usage" in chunk
):
usage = _get(chunk, "usage", chunk)
yield UsageEvent(
input_tokens=_get(usage, "input_tokens", 0),
output_tokens=_get(usage, "output_tokens", 0),
)
else:
yield make_extension_event(f"anthropic.{event_type}", chunk)