parseSseStream function

Stream<SseEvent> parseSseStream(
  1. Stream<List<int>> byteStream, {
  2. String? doneSignal = '[DONE]',
})

Parse a byte stream of Server-Sent Events into SseEvents.

Events where data equals doneSignal (default '[DONE]') are filtered out and the stream completes. Set doneSignal to null to disable this behavior (the stream ends when the connection closes).

See the SSE specification.

Implementation

Stream<SseEvent> parseSseStream(
  Stream<List<int>> byteStream, {
  String? doneSignal = '[DONE]',
}) async* {
  // Decode UTF-8, then split on lines. SSE uses \n, \r, or \r\n.
  final lines = byteStream
      .transform(const Utf8Decoder(allowMalformed: true))
      .transform(const LineSplitter());

  String? eventType;
  final dataLines = <String>[];
  String? lastId;

  await for (final line in lines) {
    if (line.isEmpty) {
      // Empty line = event boundary. Dispatch if we have data.
      if (dataLines.isNotEmpty) {
        final data = dataLines.join('\n');
        if (doneSignal != null && data == doneSignal) return;
        yield SseEvent(event: eventType, data: data, id: lastId);
      }
      eventType = null;
      dataLines.clear();
      continue;
    }

    // Lines starting with ':' are comments - ignore.
    if (line.startsWith(':')) continue;

    final colonIndex = line.indexOf(':');
    final String field;
    final String value;
    if (colonIndex == -1) {
      field = line;
      value = '';
    } else {
      field = line.substring(0, colonIndex);
      // Skip optional single space after colon.
      final valueStart =
          (colonIndex + 1 < line.length && line[colonIndex + 1] == ' ')
          ? colonIndex + 2
          : colonIndex + 1;
      value = line.substring(valueStart);
    }

    switch (field) {
      case 'data':
        dataLines.add(value);
      case 'event':
        eventType = value;
      case 'id':
        lastId = value;
      // 'retry' and unknown fields are ignored.
    }
  }

  // If the stream ends without a trailing blank line, flush any pending event.
  if (dataLines.isNotEmpty) {
    final data = dataLines.join('\n');
    if (doneSignal != null && data == doneSignal) return;
    yield SseEvent(event: eventType, data: data, id: lastId);
  }
}