stream method

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

Streams OpenAI-compatible Server-Sent Events.

Emits chat.delta events as text chunks arrive and a final chat.completed event when the provider sends a finish reason or the stream ends.

Implementation

@override

/// Streams OpenAI-compatible Server-Sent Events.
///
/// Emits `chat.delta` events as text chunks arrive and a final
/// `chat.completed` event when the provider sends a finish reason or the
/// stream ends.
Stream<ChatEvent> stream(ChatRequest request) async* {
  final token = bearerToken ?? apiKey!;
  final content = StringBuffer();

  try {
    final chunks = httpClient.stream(
      AiHttpRequest(
        method: 'POST',
        uri: endpoint,
        headers: {
          'Content-Type': 'application/json',
          'Accept': 'text/event-stream',
          'Authorization': 'Bearer $token',
        },
        body: jsonEncode({
          'model': request.model,
          'messages':
              request.messages.map((message) => message.toMap()).toList(),
          'stream': true,
          if (request.temperature != null) 'temperature': request.temperature,
          if (request.maxTokens != null) 'max_tokens': request.maxTokens,
          ...request.metadata,
        }),
      ),
    );

    await for (final eventData in aiSseDataEvents(chunks)) {
      if (eventData.trim() == '[DONE]') {
        break;
      }

      final payload = jsonDecode(eventData);
      if (payload is! Map) continue;
      final map = Map<String, dynamic>.from(payload);
      final choices = map['choices'] as List? ?? const [];
      if (choices.isEmpty) continue;

      final choice = Map<String, dynamic>.from(choices.first as Map);
      final delta = Map<String, dynamic>.from(
        choice['delta'] as Map? ?? const {},
      );
      final chunk = delta['content']?.toString();
      if (chunk != null && chunk.isNotEmpty) {
        content.write(chunk);
        yield ChatEvent(
          type: 'chat.delta',
          payload: {
            'providerId': id,
            'model': request.model,
            'content': chunk,
            'raw': map,
          },
        );
      }

      final finishReason = choice['finish_reason']?.toString();
      if (finishReason != null && finishReason.isNotEmpty) {
        yield ChatEvent(
          type: 'chat.completed',
          payload: {
            'providerId': id,
            'model': request.model,
            'content': content.toString(),
            'finishReason': finishReason,
            'raw': map,
          },
        );
        return;
      }
    }
  } catch (error) {
    throw AiProviderException(
      providerId: id,
      message: 'OpenAI stream transport failed',
      cause: error,
    );
  }

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