onResponse method

  1. @override
void onResponse(
  1. Response response,
  2. ResponseInterceptorHandler handler
)

Called when the response is about to be resolved.

Implementation

@override
void onResponse(Response response, ResponseInterceptorHandler handler) async {
  Future appendStreamDone() async {
    List<ServerSentEvent> filteredList =
        await _filterService.resolve(SSEAutoRemoveInterceptor.genAutoRemoveSSE());
    _resolveSseEvent(filteredList, response.requestOptions.path, toPeek: true);
  }

  Future appendStreamOpen() async {
    ServerSentEvent sse = ServerSentEvent(
        elementType: eventStreamOpen, sessionLogId: eventStreamOpenLogId, result: "client_mock");
    List<ServerSentEvent> filteredList = await _filterService.resolve(sse);
    _resolveSseEvent(filteredList, response.requestOptions.path, toPeek: true);
  }

  slog.d("start resolve response", tag: tag);
  if (response.requestOptions.isStream() && response.data is ResponseBody) {
    slog.d("open stream url : ${response.requestOptions.uri}", tag: tag);
    Stream stream = (response.data as ResponseBody)
        .stream
        .map((event) => event.toList())
        .transform(utf8.decoder);
    if (response.requestOptions.isOfflineStream()) {
      _sseStream = response.requestOptions.getOfflineStream()!.onOfflineStream();
    } else {
      _sseStream = SSECore(stream, adapterService: _adapterService);
    }
    _sseCacheDeliverer.reset();
    _sseStream?.open(
        onOpened: () async {
          await appendStreamOpen();
        },
        onReceive: (event) async {
          List<ServerSentEvent> filteredList = await _filterService.resolve(event);
          _resolveSseEvent(filteredList, response.requestOptions.path, toPeek: true);
        },
        onDone: (list) async {
          await appendStreamDone();
          slog.df("stream transform done", tag: tag);
          _streamTransforming.value = false;
          _changeConnectState(ConnectState.connectSuspend);
          _sseCacheDeliverer.flushPeek(pop: (sseC) {
            _interceptorManager.deliver(sseC, isPeek: true);
          });
          _filterService.reset();
          _sseBridge.offWork();
        },
        onError: (e) async {
          await appendStreamDone();
          slog.df("stream transform error $e", tag: tag);
          _handleConnectError();
          _streamTransforming.value = false;
          _filterService.reset();
          _sseBridge.offWork();
        },
        requestOptions: response.requestOptions);
  } else {
    /// A specific scenario where the client requests stream data and the server does not return stream data.
    if (response.requestOptions.isStream()) {
      _streamTransforming.value = false;
    }
    slog.d("response is not stream format", tag: tag);
  }
  handler.next(response);
}