openAiChunksFromEvents function

Stream<LLMStreamChunk> openAiChunksFromEvents(
  1. Stream<Map<String, dynamic>> events, {
  2. String? provider,
})

OpenAI 互換のストリーミングイベント(chat.completion.chunk)列を LLMStreamChunk 列へ変換する。

ストリーム中の error イベントは LLMException に正規化して投げる。

Implementation

Stream<LLMStreamChunk> openAiChunksFromEvents(
  Stream<Map<String, dynamic>> events, {
  String? provider,
}) async* {
  final text = StringBuffer();
  final reasoning = StringBuffer();
  final toolCalls = <int, _StreamToolCall>{};
  final images = <LLMContentPart>[];
  LLMFinishReason? finishReason;
  LLMUsage? usage;

  await for (final event in events) {
    final error = event['error'];
    if (error is Map<String, dynamic>) {
      throw _streamErrorToException(error, provider: provider, raw: event);
    }
    final rawUsage = event['usage'];
    if (rawUsage is Map<String, dynamic>) {
      usage = _parseOpenAiUsage(rawUsage);
    }
    final choices = event['choices'];
    if (choices is! List || choices.isEmpty) continue;
    final choice = choices.first;
    if (choice is! Map<String, dynamic>) continue;

    final rawFinish = choice['finish_reason'];
    if (rawFinish is String && rawFinish.isNotEmpty) {
      finishReason = LLMFinishReason.parse(rawFinish);
    }
    final delta = choice['delta'];
    if (delta is! Map<String, dynamic>) continue;

    final reasoningDelta = delta['reasoning_content'] ?? delta['reasoning'];
    if (reasoningDelta is String && reasoningDelta.isNotEmpty) {
      reasoning.write(reasoningDelta);
      yield LLMStreamChunk(reasoningDelta: reasoningDelta);
    }
    final contentDelta = delta['content'];
    if (contentDelta is String && contentDelta.isNotEmpty) {
      text.write(contentDelta);
      yield LLMStreamChunk(delta: contentDelta);
    }
    final imageDeltas = delta['images'];
    if (imageDeltas is List) _appendOpenAiImages(imageDeltas, images);
    final toolCallDeltas = delta['tool_calls'];
    if (toolCallDeltas is List) {
      _mergeToolCallDeltas(toolCallDeltas, toolCalls);
    }
  }

  yield LLMStreamChunk(
    parts: [
      if (reasoning.isNotEmpty) LLMReasoningPart(reasoning.toString()),
      if (text.isNotEmpty) LLMTextPart(text.toString()),
      ...images,
      for (final index in toolCalls.keys.toList()..sort())
        toolCalls[index]!.toPart(),
    ],
    finishReason: finishReason ?? LLMFinishReason.other,
    usage: usage,
  );
}