stream method
Stream<RpcStreamEvent<RpcReply> >
stream(
- String method, {
- Map<
String, Object?> params = const {}, - Duration timeout = const Duration(minutes: 5),
- 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;
}