stream method

  1. @override
Stream<ChatEvent> stream(
  1. ChatRequest request
)
override

Streams Anthropic message events.

Emits chat.delta for content_block_delta text events and a final chat.completed event when Anthropic sends message_stop.

Implementation

@override

/// Streams Anthropic message events.
///
/// Emits `chat.delta` for `content_block_delta` text events and a final
/// `chat.completed` event when Anthropic sends `message_stop`.
Stream<ChatEvent> stream(ChatRequest request) async* {
  String? system;
  final messages = <Map<String, dynamic>>[];
  final content = StringBuffer();

  for (final message in request.messages) {
    if (message.role == 'system') {
      system =
          system == null ? message.content : '$system\n${message.content}';
      continue;
    }
    messages.add({
      'role': message.role,
      'content': message.content,
    });
  }

  try {
    final chunks = httpClient.stream(
      AiHttpRequest(
        method: 'POST',
        uri: endpoint,
        headers: {
          'Content-Type': 'application/json',
          'Accept': 'text/event-stream',
          'x-api-key': apiKey,
          'anthropic-version': '2023-06-01',
        },
        body: jsonEncode({
          'model': request.model,
          'messages': messages,
          'max_tokens': request.maxTokens ?? 1024,
          'stream': true,
          if (request.temperature != null) 'temperature': request.temperature,
          if (system != null && system.isNotEmpty) 'system': system,
          ...request.metadata,
        }),
      ),
    );

    await for (final eventData in aiSseDataEvents(chunks)) {
      final payload = jsonDecode(eventData);
      if (payload is! Map) continue;
      final map = Map<String, dynamic>.from(payload);
      final type = map['type']?.toString();

      if (type == 'content_block_delta') {
        final delta = Map<String, dynamic>.from(
          map['delta'] as Map? ?? const {},
        );
        final text = delta['text']?.toString();
        if (text != null && text.isNotEmpty) {
          content.write(text);
          yield ChatEvent(
            type: 'chat.delta',
            payload: {
              'providerId': id,
              'model': request.model,
              'content': text,
              'raw': map,
            },
          );
        }
      } else if (type == 'message_stop') {
        yield ChatEvent(
          type: 'chat.completed',
          payload: {
            'providerId': id,
            'model': request.model,
            'content': content.toString(),
            'raw': map,
          },
        );
        return;
      }
    }
  } catch (error) {
    throw AiProviderException(
      providerId: id,
      message: 'Anthropic stream transport failed',
      cause: error,
    );
  }

  yield ChatEvent(
    type: 'chat.completed',
    payload: {
      'providerId': id,
      'model': request.model,
      'content': content.toString(),
    },
  );
}