stream method

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

Streams Gemini Server-Sent Events.

Emits chat.delta for each text part received and a final chat.completed event with the accumulated content after the stream ends.

Implementation

@override

/// Streams Gemini Server-Sent Events.
///
/// Emits `chat.delta` for each text part received and a final
/// `chat.completed` event with the accumulated content after the stream ends.
Stream<ChatEvent> stream(ChatRequest request) async* {
  final contents = request.messages
      .where((message) => message.role != 'system')
      .map((message) => {
            'role': message.role == 'assistant' ? 'model' : 'user',
            'parts': [
              {'text': message.content}
            ],
          })
      .toList(growable: false);

  final system = request.messages
      .where((message) => message.role == 'system')
      .map((message) => message.content)
      .join('\n');

  final content = StringBuffer();

  try {
    final chunks = httpClient.stream(
      AiHttpRequest(
        method: 'POST',
        uri: _streamEndpoint(endpoint),
        headers: {
          'Content-Type': 'application/json',
          'Accept': 'text/event-stream',
          if (bearerToken != null && bearerToken!.isNotEmpty)
            'Authorization': 'Bearer $bearerToken'
          else
            'x-goog-api-key': apiKey!,
        },
        body: jsonEncode({
          'contents': contents,
          if (system.isNotEmpty)
            'systemInstruction': {
              'parts': [
                {'text': system}
              ],
            },
          if (request.temperature != null || request.maxTokens != null)
            'generationConfig': {
              if (request.temperature != null)
                'temperature': request.temperature,
              if (request.maxTokens != null)
                'maxOutputTokens': request.maxTokens,
            },
          ...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 text = _geminiTextFromPayload(map);
      if (text != null && text.isNotEmpty) {
        content.write(text);
        yield ChatEvent(
          type: 'chat.delta',
          payload: {
            'providerId': id,
            'model': request.model,
            'content': text,
            'raw': map,
          },
        );
      }
    }
  } catch (error) {
    throw AiProviderException(
      providerId: id,
      message: 'Gemini stream transport failed',
      cause: error,
    );
  }

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