stream method

  1. @override
Stream<RpcStreamEvent<RpcReply>> stream(
  1. String method, {
  2. Map<String, Object?> params = const {},
  3. Duration timeout = const Duration(minutes: 5),
  4. RpcCancellationToken? cancellationToken,
})
override

Implementation

@override
Stream<RpcStreamEvent<RpcReply>> stream(
  String method, {
  Map<String, Object?> params = const {},
  Duration timeout = const Duration(minutes: 5),
  RpcCancellationToken? cancellationToken,
}) {
  if (_state == _SidecarState.closing ||
      _state == _SidecarState.closed ||
      _state == _SidecarState.failed) {
    return Stream.error(_terminalError ?? const BackendClosedException());
  }
  if (cancellationToken?.isCancelled ?? false) {
    return Stream.error(RpcCancelledException(method));
  }

  final id = '${++_nextID}';
  final pending = _PendingSidecarStream();
  StreamSubscription<RpcStreamEvent<RpcReply>>? relay;
  late final StreamController<RpcStreamEvent<RpcReply>> output;
  output = StreamController<RpcStreamEvent<RpcReply>>(
    sync: true,
    onListen: () {
      _streams[id] = pending;
      pending.timeout = Timer(
        timeout,
        () => _cancelPendingStream(
          id,
          TimeoutException('Sidecar stream $method timed out.', timeout),
        ),
      );
      pending.cancellationSubscription = cancellationToken?.onCancel.listen(
        (_) => _cancelPendingStream(id, RpcCancelledException(method)),
      );
      relay = pending.controller.stream.listen(
        (event) {
          output.add(event);
          if (_streams[id] == pending) {
            unawaited(_sendStreamAcknowledgement(id, event.sequence));
          }
        },
        onError: output.addError,
        onDone: output.close,
      );
      unawaited(
        _sendRequestWhenReady(
          id,
          encodeRpcRequest(
            id: id,
            method: method,
            params: params,
            token: _token,
            stream: true,
            streamWindow: _streamWindow,
          ),
        ),
      );
    },
    onPause: () => relay?.pause(),
    onResume: () => relay?.resume(),
    onCancel: () async {
      final abandoned = _streams.remove(id);
      if (identical(abandoned, pending)) {
        pending.dispose();
        if (_state == _SidecarState.running) {
          unawaited(_sendCancellation(id));
        }
      }
      await relay?.cancel();
    },
  );
  return output.stream;
}