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