openSseStream method
Opens a long-lived Server-Sent Events stream.
Implementation
@override
CancellableSubscription openSseStream(
NetworkRequest request,
SseStreamHandler handler,
) {
final cancelToken = CancelToken();
final subscription = _DioCancellableSubscription(cancelToken);
Future<void> run() async {
try {
final response = await _dio.request<ResponseBody>(
request.url,
data: request.body,
cancelToken: cancelToken,
options: Options(
method: request.method.name.toUpperCase(),
responseType: ResponseType.stream,
headers: {
'Accept': 'text/event-stream',
'Cache-Control': 'no-cache',
...request.headers,
},
// Network plan §4.2: connect within 10 s; receiveTimeout is the
// longest silent gap between chunks before the stream is dead.
connectTimeout: request.connectTimeout ?? _restTimeout,
receiveTimeout: request.receiveTimeout ?? _sseSilenceTimeout,
validateStatus: (_) => true,
),
);
if (subscription.isCancelled) return;
// validateStatus accepts everything so the status reaches us here:
// an error page is not an open stream (network plan §4.3).
final status = response.statusCode ?? 0;
if (status < 200 || status >= 300) {
// Dio already reads the body; only the token closes the response
// and its connection.
cancelToken.cancel('SSE stream rejected (status=$status)');
handler.onError(SseHttpException(status));
return;
}
final stream = response.data?.stream;
if (stream == null) {
handler.onError(StateError('No response stream available'));
return;
}
handler.onOpen();
final parser = SseParser(handler.onEvent);
// A streaming decoder: a multi-byte character split across two
// chunks is decoded whole, not as two malformed halves.
subscription.streamSubscription = stream
.cast<List<int>>()
.transform(const Utf8Decoder(allowMalformed: true))
.listen(
(text) {
if (subscription.isCancelled) return;
parser.add(text);
},
onError: (Object error, StackTrace _) {
if (!subscription.isCancelled) {
handler.onError(error);
}
},
onDone: () {
if (!subscription.isCancelled) {
handler.onClosed();
}
},
cancelOnError: true,
);
} catch (error) {
if (!subscription.isCancelled) {
handler.onError(error);
}
}
}
unawaited(run());
return subscription;
}