openSseStream method

  1. @override
CancellableSubscription openSseStream(
  1. NetworkRequest request,
  2. SseStreamHandler handler
)
override

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;
}