subscribeEvents method

Stream<SseEvent> subscribeEvents()

Opens /api/events and emits decoded SSE data payloads.

Implementation

Stream<SseEvent> subscribeEvents() async* {
  final request = await _httpClient.getUrl(_resolve('/api/events'));
  request.headers.set(HttpHeaders.acceptHeader, 'text/event-stream');
  final response = await request.close();

  if (response.statusCode < 200 || response.statusCode >= 300) {
    final body = await utf8.decoder.bind(response).join();
    throw ApiClientException(
      'SSE connection failed',
      statusCode: response.statusCode,
      body: _tryParseObject(body),
    );
  }

  final dataBuffer = StringBuffer();
  await for (final line in utf8.decoder.bind(response).transform(const LineSplitter())) {
    if (line.isEmpty) {
      if (dataBuffer.isNotEmpty) {
        final payload = dataBuffer.toString();
        dataBuffer.clear();
        final decoded = _tryParseObject(payload);
        if (decoded != null) {
          yield SseEvent.fromJson(decoded);
        }
      }
      continue;
    }

    if (line.startsWith('data:')) {
      if (dataBuffer.isNotEmpty) {
        dataBuffer.write('\n');
      }
      dataBuffer.write(line.substring(5).trimLeft());
    }
  }
}