stream method
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(),
},
);
}