streamAnthropic function
AssistantMessageEventStream
streamAnthropic(
- Model model,
- Context context, [
- AnthropicOptions? options,
- Client? client,
Streams an assistant message from an Anthropic messages endpoint.
Ported from pi's stream in anthropic-messages.ts. The endpoint is
{model.baseUrl}/v1/messages (default base URL
https://api.anthropic.com on the model descriptor).
Errors-as-events invariant (non-negotiable): this function never
throws. Network failures, non-200 responses, provider error SSE events,
malformed SSE, and aborts all terminate the returned stream with an
ErrorEvent carrying StopReason.error or StopReason.aborted.
client overrides the HTTP client (used by tests with
http.testing.MockClient); when omitted, an owned client is created and
closed when the stream finishes.
Implementation
AssistantMessageEventStream streamAnthropic(
Model model,
Context context, [
AnthropicOptions? options,
http.Client? client,
]) {
final eventStream = AssistantMessageEventStream();
final cancelToken = options?.cancelToken;
// No injected client: the shared keep-alive client (never closed per call).
final httpClient = client ?? sharedProviderHttpClient();
// Blocks accumulate in the shared state holder; each event carries a
// fresh immutable snapshot of them (pi mutates one `output` object).
final state = ProviderStreamState(model);
final session = _AnthropicStreamSession(model, eventStream, state);
unawaited(
runProviderStream(
eventStream,
state,
cancelToken,
httpClient,
ownsClient: false, // shared or injected — never closed per call
body: () => session.run(context, options, httpClient, cancelToken),
),
);
return eventStream;
}