Метод timeout
Создает новый поток с теми же событиями, что и этот поток.
Когда кто-то подписывается на возвращаемый поток, и более чем timeLimit проходит без получения каких-либо событий от этого потока, вызывается функция onTimeout, которая может затем генерировать дополнительные события в возвращаемом потоке.
Отсчет времени начинается при подписке на возвращаемый поток и перезапускается при получении события от этого потока или при приостановке и возобновлении подписки на возвращаемый поток. Отсчет времени останавливается при приостановке или отмене подписки на возвращаемый поток. Новый отсчет не запускается, когда отсчет завершается, и вызывается функция onTimeout, даже если события генерируются. Если задержка между событиями этого потока — это многократное значение timeLimit, произойдет не более одного таймаута между событиями.
Функция onTimeout вызывается с одним аргументом: EventSink, который позволяет помещать события в возвращаемый поток. Этот EventSink действителен только во время вызова onTimeout. Вызов EventSink.close на передаваемом в onTimeout sink закрывает возвращаемый поток, и дальнейшая обработка событий не выполняется.
Если onTimeout опущено, таймаут отправит TimeoutException в канал ошибок возвращаемого потока. Если вызов onTimeout вызывает ошибку, она будет отправлена как ошибка в возвращаемый поток.
Возвращаемый поток является потоком широковещательной передачи, если таковым является исходный поток. Если на поток широковещательной передачи подписано более одного раза, каждая подписка будет иметь свой таймер, который начинает отсчет при подписке, и таймеры подписок можно приостановить индивидуально.
Пример:
Future<String> waitTask() async {
return await Future.delayed(
const Duration(seconds: 4), () => 'Complete');
}
final stream = Stream<String>.fromFuture(waitTask())
.timeout(const Duration(seconds: 2), onTimeout: (controller) {
print('TimeOut occurred');
controller.close();
});
stream.listen(print, onDone: () => print('Done'));
// Outputs:
// TimeOut occurred
// Done Реализация
Stream<T> timeout(Duration timeLimit, {void onTimeout(EventSink<T> sink)?}) {
_StreamControllerBase<T> controller;
if (isBroadcast) {
controller = new _SyncBroadcastStreamController<T>(null, null);
} else {
controller = new _SyncStreamController<T>(null, null, null, null);
}
Zone zone = Zone.current;
// Register callback immediately.
_TimerCallback timeoutCallback;
if (onTimeout == null) {
timeoutCallback = () {
controller.addError(
new TimeoutException("No stream event", timeLimit), null);
};
} else {
var registeredOnTimeout =
zone.registerUnaryCallback<void, EventSink<T>>(onTimeout);
var wrapper = new _ControllerEventSinkWrapper<T>(null);
timeoutCallback = () {
wrapper._sink = controller; // Only valid during call.
zone.runUnaryGuarded(registeredOnTimeout, wrapper);
wrapper._sink = null;
};
}
// All further setup happens inside `onListen`.
controller.onListen = () {
Timer timer = zone.createTimer(timeLimit, timeoutCallback);
var subscription = this.listen(null);
// Set up event forwarding. Each data or error event resets the timer
subscription
..onData((T event) {
timer.cancel();
timer = zone.createTimer(timeLimit, timeoutCallback);
// Controller is synchronous, and the call might close the stream
// and cancel the timer,
// so create the Timer before calling into add();
// issue: https://github.com/dart-lang/sdk/issues/37565
controller.add(event);
})
..onError((Object error, StackTrace stackTrace) {
timer.cancel();
timer = zone.createTimer(timeLimit, timeoutCallback);
controller._addError(
error, stackTrace); // Avoid Zone error replacement.
})
..onDone(() {
timer.cancel();
controller.close();
});
// Set up further controller callbacks.
controller.onCancel = () {
timer.cancel();
return subscription.cancel();
};
if (!isBroadcast) {
controller
..onPause = () {
timer.cancel();
subscription.pause();
}
..onResume = () {
subscription.resume();
timer = zone.createTimer(timeLimit, timeoutCallback);
};
}
};
return controller.stream;
}
© 2012 the Dart project authors
Licensed under the BSD 3-Clause "New" or "Revised" License.
https://api.dart.dev/stable/2.18.5/dart-async/Stream/timeout.html