Chromium Code Reviews
chromiumcodereview-hr@appspot.gserviceaccount.com (chromiumcodereview-hr) | Please choose your nickname with Settings | Help | Chromium Project | Gerrit Changes | Sign out
(708)

Unified Diff: sdk/lib/async/future_impl.dart

Issue 23926011: Rewrite Futures. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Fix bad asserts. Created 7 years, 3 months ago
Use n/p to move between diff chunks; N/P to move between comments. Draft comments are only viewable by you.
Jump to:
View side-by-side diff with in-line comments
Download patch
« no previous file with comments | « sdk/lib/async/future.dart ('k') | sdk/lib/async/stream.dart » ('j') | no next file with comments »
Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
Index: sdk/lib/async/future_impl.dart
diff --git a/sdk/lib/async/future_impl.dart b/sdk/lib/async/future_impl.dart
index 4c846c67c6b5e7753ce1022d5dcc7a8c5eb1f206..ea8f9027bd8494c05685abfb90eb9fcdc69c88b0 100644
--- a/sdk/lib/async/future_impl.dart
+++ b/sdk/lib/async/future_impl.dart
@@ -4,155 +4,53 @@
part of dart.async;
-abstract class _Completer<T> implements Completer<T> {
- final Future<T> future;
- bool _isComplete = false;
+/** The onValue and onError handlers return either a value or a future */
+typedef dynamic _FutureOnValue<T>(T value);
+typedef dynamic _FutureOnError(error);
+/** Test used by [Future.catchError] to handle skip some errors. */
+typedef bool _FutureErrorTest(var error);
+/** Used by [WhenFuture]. */
+typedef _FutureAction();
- _Completer() : future = new _FutureImpl<T>() {
- _FutureImpl futureImpl = future;
- futureImpl._zone.expectCallback();
- }
+abstract class _Completer<T> implements Completer<T> {
+ final _Future<T> future = new _Future<T>();
- void _setFutureValue(T value);
- void _setFutureError(error);
+ void complete([T value]);
- void complete([T value]) {
- if (_isComplete) throw new StateError("Future already completed");
- _isComplete = true;
- _FutureImpl futureImpl = future;
- _setFutureValue(value);
- }
+ void completeError(Object error, [Object stackTrace = null]);
- void completeError(Object error, [Object stackTrace = null]) {
- if (_isComplete) throw new StateError("Future already completed");
- _isComplete = true;
- if (stackTrace != null) {
- // Force the stack trace onto the error, even if it already had one.
- _attachStackTrace(error, stackTrace);
- }
- _FutureImpl futureImpl = future;
- _setFutureError(error);
- }
-
- bool get isCompleted => _isComplete;
+ // The future's _isComplete doesn't take into account pending completions.
+ // We therefore use _mayComplete.
+ bool get isCompleted => !future._mayComplete;
}
class _AsyncCompleter<T> extends _Completer<T> {
- void _setFutureValue(T value) {
- _FutureImpl future = this.future;
- future._asyncSetValue(value);
- // The async-error will schedule another callback, so we can cancel
- // the expectation without shutting down the zone.
- future._zone.cancelCallbackExpectation();
- }
-
- void _setFutureError(error) {
- _FutureImpl future = this.future;
- future._asyncSetError(error);
- // The async-error will schedule another callback, so we can cancel
- // the expectation without shutting down the zone.
- future._zone.cancelCallbackExpectation();
- }
-}
-class _SyncCompleter<T> extends _Completer<T> {
- void _setFutureValue(T value) {
- _FutureImpl future = this.future;
- future._setValue(value);
- future._zone.cancelCallbackExpectation();
- }
-
- void _setFutureError(error) {
- _FutureImpl future = this.future;
- future._setError(error);
- future._zone.cancelCallbackExpectation();
+ void complete([T value]) {
+ if (!future._mayComplete) throw new StateError("Future already completed");
+ future._asyncComplete(value);
}
-}
-/**
- * A listener on a future.
- *
- * When the future completes, the [_sendValue] or [_sendError] method
- * is invoked with the result.
- *
- * Listeners are kept in a linked list.
- */
-abstract class _FutureListener<T> {
- _FutureListener _nextListener;
- factory _FutureListener.wrap(_FutureImpl future) {
- return new _FutureListenerWrapper(future);
+ void completeError(Object error, [Object stackTrace = null]) {
+ if (!future._mayComplete) throw new StateError("Future already completed");
+ future._asyncCompleteError(error, stackTrace);
}
- void _sendValue(T value);
- void _sendError(error);
-
- bool _inSameErrorZone(_Zone otherZone);
-}
-
-/** Adapter for a [_FutureImpl] to be a future result listener. */
-class _FutureListenerWrapper<T> implements _FutureListener<T> {
- _FutureImpl future;
- _FutureListener _nextListener;
- _FutureListenerWrapper(this.future);
- _sendValue(T value) { future._setValueUnchecked(value); }
- _sendError(error) { future._setErrorUnchecked(error); }
- bool _inSameErrorZone(_Zone otherZone) => future._inSameErrorZone(otherZone);
}
-/**
- * This listener is installed at error-zone boundaries. It signals an
- * uncaught error in the zone of origin when an error is sent from one error
- * zone to another.
- *
- * When a Future is listening to another Future and they have not been
- * instantiated in the same error-zone then Futures put an instance of this
- * class between them (see [_FutureImpl._addListener]).
- *
- * For example:
- *
- * var completer = new Completer();
- * var future = completer.future.then((x) => x);
- * catchErrors(() {
- * var future2 = future.catchError(print);
- * });
- * completer.completeError(499);
- *
- * In this example `future` and `future2` are in different error-zones. The
- * error (499) that originates outside `catchErrors` must not reach the
- * `catchError` future (`future2`) inside `catchErrors`.
- *
- * When invoking `catchError` on `future` the Future installs an
- * [_ErrorZoneBoundaryListener] between itself and the result, `future2`.
- *
- * Conceptually _ErrorZoneBoundaryListeners could be implemented as
- * `catchError`s on the origin future as well.
- */
-class _ErrorZoneBoundaryListener implements _FutureListener {
- _FutureListener _nextListener;
- final _FutureListener _listener;
-
- _ErrorZoneBoundaryListener(this._listener);
-
- bool _inSameErrorZone(_Zone otherZone) {
- // Should never be called. We use [_inSameErrorZone] to know if we have
- // to insert an instance of [_ErrorZoneBoundaryListener] (and in the
- // controller). Once we have inserted one we should never need to use it
- // anymore.
- throw new UnsupportedError(
- "A Zone boundary doesn't support the inSameErrorZone test.");
- }
+class _SyncCompleter<T> extends _Completer<T> {
- void _sendValue(value) {
- _listener._sendValue(value);
+ void complete([T value]) {
+ if (!future._mayComplete) throw new StateError("Future already completed");
+ future._complete(value);
}
- void _sendError(error) {
- // We are not allowed to send an error from one error-zone to another.
- // This is the whole purpose of this class.
- _Zone.current.handleUncaughtError(error);
+ void completeError(Object error, [Object stackTrace = null]) {
+ if (!future._mayComplete) throw new StateError("Future already completed");
+ future._completeError(error, stackTrace);
}
}
-class _FutureImpl<T> implements Future<T> {
+class _Future<T> implements Future<T> {
// State of the future. The state determines the interpretation of the
// [resultOrListeners] field.
// TODO(lrn): rename field since it can also contain a chained future.
@@ -161,35 +59,40 @@ class _FutureImpl<T> implements Future<T> {
/// [resultOrListeners] field holds a single-linked list of
/// [FutureListener] listeners.
static const int _INCOMPLETE = 0;
- /// Pending completion. Set when completed using [_asyncSetValue] or
- /// [_asyncSetError]. It is an error to try to complete it again.
+ /// Pending completion. Set when completed using [_asyncComplete] or
+ /// [_asyncCompleteError]. It is an error to try to complete it again.
static const int _PENDING_COMPLETE = 1;
/// The future has been chained to another future. The result of that
/// other future becomes the result of this future as well.
- /// In this state, the [resultOrListeners] field holds the future that
- /// will give the result to this future. Both existing and new listeners are
- /// forwarded directly to the other future.
+ /// In this state, no callback should be executed anymore.
+ // TODO(floitsch): we don't really need a special "_CHAINED" state. We could
+ // just use the PENDING_COMPLETE state instead.
static const int _CHAINED = 2;
- /// The future has been chained to another future, but there hasn't been
- /// any listeners added to this future yet. If it is completed with an
- /// error, the error will be considered unhandled.
- static const int _CHAINED_UNLISTENED = 6;
/// The future has been completed with a value result.
- static const int _VALUE = 8;
+ static const int _VALUE = 4;
/// The future has been completed with an error result.
- static const int _ERROR = 12;
+ static const int _ERROR = 8;
/** Whether the future is complete, and as what. */
int _state = _INCOMPLETE;
final _Zone _zone = _Zone.current.fork();
- bool get _isChained => (_state & _CHAINED) != 0;
- bool get _hasChainedListener => _state == _CHAINED;
- bool get _isComplete => _state >= _VALUE;
bool get _mayComplete => _state == _INCOMPLETE;
+ bool get _isChained => _state == _CHAINED;
+ bool get _isComplete => _state >= _VALUE;
bool get _hasValue => _state == _VALUE;
- bool get _hasError => _state >= _ERROR;
+ bool get _hasError => _state == _ERROR;
+
+ set _isChained(bool value) {
+ if (value) {
+ assert(!_isComplete);
+ _state = _CHAINED;
+ } else {
+ assert(_isChained);
+ _state = _INCOMPLETE;
+ }
+ }
/**
* Either the result, a list of listeners or another future.
@@ -212,175 +115,132 @@ class _FutureImpl<T> implements Future<T> {
*/
var _resultOrListeners;
- _FutureImpl();
-
- _FutureImpl.immediate(T value) {
- _state = _VALUE;
- _resultOrListeners = value;
+ /**
+ * A [_Future] implements a linked list. If a future has more than one
+ * listener the [_nextListener] field of the first listener points to the
+ * remaining listeners.
+ */
+ // TODO(floitsch): since single listeners are the common case we should
+ // use a bit to indicate that the _resultOrListeners contains a container.
+ _Future _nextListener;
+
+ // TODO(floitsch): we only need two closure fields to store the callbacks.
+ // If we store the type of a closure in the state field (where there are
+ // still bits left), we can just store two closures instead of using 4
+ // fields of which 2 are always null.
+ final _FutureOnValue _onValueCallback;
+ final _FutureErrorTest _errorTestCallback;
+ final _FutureOnError _onErrorCallback;
+ final _FutureAction _whenCompleteActionCallback;
+
+ _FutureOnValue get _onValue => _isChained ? null : _onValueCallback;
+ _FutureErrorTest get _errorTest => _isChained ? null : _errorTestCallback;
+ _FutureOnError get _onError => _isChained ? null : _onErrorCallback;
+ _FutureAction get _whenCompleteAction
+ => _isChained ? null : _whenCompleteActionCallback;
+
+ _Future()
+ : _onValueCallback = null, _errorTestCallback = null,
+ _onErrorCallback = null, _whenCompleteActionCallback = null;
+
+ _Future.immediate(T value)
+ : _onValueCallback = null, _errorTestCallback = null,
+ _onErrorCallback = null, _whenCompleteActionCallback = null {
+ _asyncComplete(value);
+ }
+
+ _Future.immediateError(var error, [Object stackTrace])
+ : _onValueCallback = null, _errorTestCallback = null,
+ _onErrorCallback = null, _whenCompleteActionCallback = null {
+ _asyncCompleteError(error, stackTrace);
+ }
+
+ _Future._then(this._onValueCallback, this._onErrorCallback)
+ : _errorTestCallback = null, _whenCompleteActionCallback = null {
+ _zone.expectCallback();
}
- _FutureImpl.immediateError(var error, [Object stackTrace]) {
- if (stackTrace != null) {
- // Force stack trace onto error, even if it had already one.
- _attachStackTrace(error, stackTrace);
- }
- _asyncSetError(error);
+ _Future._catchError(this._onErrorCallback, this._errorTestCallback)
+ : _onValueCallback = null, _whenCompleteActionCallback = null {
+ _zone.expectCallback();
}
- factory _FutureImpl.wait(Iterable<Future> futures) {
- Completer completer;
- // List collecting values from the futures.
- // Set to null if an error occurs.
- List values;
- void handleError(error) {
- if (values != null) {
- values = null;
- completer.completeError(error);
- }
- }
- // As each future completes, put its value into the corresponding
- // position in the list of values.
- int remaining = 0;
- for (Future future in futures) {
- int pos = remaining++;
- future.catchError(handleError).then((Object value) {
- if (values == null) return null;
- values[pos] = value;
- remaining--;
- if (remaining == 0) {
- completer.complete(values);
- }
- });
- }
- if (remaining == 0) {
- return new Future.value(const []);
- }
- values = new List(remaining);
- completer = new Completer<List>();
- return completer.future;
+ _Future._whenComplete(this._whenCompleteActionCallback)
+ : _onValueCallback = null, _errorTestCallback = null,
+ _onErrorCallback = null {
+ _zone.expectCallback();
}
Future then(f(T value), { onError(error) }) {
- if (onError == null) {
- return new _ThenFuture(f).._subscribeTo(this);
- }
- return new _SubscribeFuture(f, onError).._subscribeTo(this);
+ _Future result;
+ result = new _Future._then(f, onError);
+ _addListener(result);
+ return result;
}
Future catchError(f(error), { bool test(error) }) {
- return new _CatchErrorFuture(f, test).._subscribeTo(this);
+ _Future result = new _Future._catchError(f, test);
+ _addListener(result);
+ return result;
}
Future<T> whenComplete(action()) {
- return new _WhenFuture<T>(action).._subscribeTo(this);
+ _Future result = new _Future<T>._whenComplete(action);
+ _addListener(result);
+ return result;
}
Stream<T> asStream() => new Stream.fromFuture(this);
- bool _inSameErrorZone(_Zone otherZone) {
- return _zone.inSameErrorZone(otherZone);
+ void _markPendingCompletion() {
+ if (!_mayComplete) throw new StateError("Future already completed");
+ _state = _PENDING_COMPLETE;
}
- void _setValue(T value) {
- if (!_mayComplete) throw new StateError("Future already completed");
- _setValueUnchecked(value);
+ T get _value {
+ assert(_isComplete && _hasValue);
+ return _resultOrListeners;
+ }
+
+ Object get _error {
+ assert(_isComplete && _hasError);
+ return _resultOrListeners;
}
- void _setValueUnchecked(T value) {
- _FutureListener listeners = _isChained ? null : _removeListeners();
+ void _setValue(T value) {
+ assert(!_isComplete); // But may have a completion pending.
_state = _VALUE;
_resultOrListeners = value;
- while (listeners != null) {
- _FutureListener listener = listeners;
- listeners = listener._nextListener;
- listener._nextListener = null;
- listener._sendValue(value);
- }
}
void _setError(Object error) {
- if (!_mayComplete) throw new StateError("Future already completed");
- _setErrorUnchecked(error);
- }
-
- void _setErrorUnchecked(Object error) {
- _FutureListener listeners;
- bool hasListeners;
- if (_isChained) {
- listeners = null;
- hasListeners = (_state == _CHAINED); // and not _CHAINED_UNLISTENED.
- } else {
- listeners = _removeListeners();
- hasListeners = (listeners != null);
- }
-
+ assert(!_isComplete); // But may have a completion pending.
_state = _ERROR;
_resultOrListeners = error;
-
- if (!hasListeners) {
- // TODO(floitsch): Hook this into unhandled error handling.
- var error = _resultOrListeners;
- _zone.handleUncaughtError(error);
- return;
- }
- while (listeners != null) {
- _FutureListener listener = listeners;
- listeners = listener._nextListener;
- listener._nextListener = null;
- listener._sendError(error);
- }
- }
-
- void _asyncSetValue(T value) {
- if (!_mayComplete) throw new StateError("Future already completed");
- _state = _PENDING_COMPLETE;
- runAsync(() { _setValueUnchecked(value); });
- }
-
- void _asyncSetError(Object error) {
- if (!_mayComplete) throw new StateError("Future already completed");
- _state = _PENDING_COMPLETE;
- runAsync(() { _setErrorUnchecked(error); });
}
- void _addListener(_FutureListener listener) {
+ void _addListener(_Future listener) {
assert(listener._nextListener == null);
- if (!listener._inSameErrorZone(_zone)) {
- listener = new _ErrorZoneBoundaryListener(listener);
- }
- if (_isChained) {
- _state = _CHAINED; // In case it was _CHAINED_UNLISTENED.
- _FutureImpl resultSource = _chainSource;
- resultSource._addListener(listener);
- return;
- }
if (_isComplete) {
// Handle late listeners asynchronously.
runAsync(() {
- if (_hasValue) {
- T value = _resultOrListeners;
- listener._sendValue(value);
- } else {
- assert(_hasError);
- listener._sendError(_resultOrListeners);
- }
+ _propagateToListeners(this, listener);
});
} else {
- assert(!_isComplete);
listener._nextListener = _resultOrListeners;
_resultOrListeners = listener;
}
}
- _FutureListener _removeListeners() {
+ _Future _removeListeners() {
// Reverse listeners before returning them, so the resulting list is in
// subscription order.
assert(!_isComplete);
- _FutureListener current = _resultOrListeners;
+ _Future current = _resultOrListeners;
_resultOrListeners = null;
- _FutureListener prev = null;
+ _Future prev = null;
while (current != null) {
- _FutureListener next = current._nextListener;
+ _Future next = current._nextListener;
current._nextListener = prev;
prev = current;
current = next;
@@ -388,270 +248,268 @@ class _FutureImpl<T> implements Future<T> {
return prev;
}
- /**
- * Make another [_FutureImpl] receive the result of this one.
- *
- * If this future is already complete, the [future] is notified
- * immediately. This function is only called during event resolution
- * where it's acceptable to send an event.
- */
- void _chain(_FutureImpl future) {
- if (!_isComplete) {
- future._chainFromFuture(this);
- } else if (_hasValue) {
- future._setValue(_resultOrListeners);
- } else {
- assert(_hasError);
- future._setError(_resultOrListeners);
- }
- }
-
- /**
- * Returns the future that this future is chained to.
- *
- * If that future is itself chained to something else,
- * get the [_chainSource] of that future instead, and make this
- * future chain directly to the earliest source.
- */
- _FutureImpl get _chainSource {
- assert(_isChained);
- _FutureImpl future = _resultOrListeners;
- if (future._isChained) {
- future = _resultOrListeners = future._chainSource;
- }
- return future;
- }
+ static void _chainFutures(Future source, _Future target) {
+ assert(!target._isComplete);
- /**
- * Make this incomplete future end up with the same result as [resultSource].
- *
- * This is done by moving all listeners to [resultSource] and forwarding all
- * future [_addListener] calls to [resultSource] directly.
- */
- void _chainFromFuture(_FutureImpl resultSource) {
- assert(!_isComplete);
- assert(!_isChained);
- if (resultSource._isChained) {
- resultSource = resultSource._chainSource;
- }
- assert(!resultSource._isChained);
- if (identical(this, resultSource)) {
- // The only unchained future in a future dependency tree (as defined
- // by the chain-relations) is the "root" that every other future depends
- // on. The future we are adding is unchained, so if it is already in the
- // tree, it must be the root, so that's the only one we need to check
- // against to detect a cycle.
- _setError(new StateError("Cyclic future dependency."));
- return;
- }
- _FutureListener cursor = _removeListeners();
- bool hadListeners = cursor != null;
- while (cursor != null) {
- _FutureListener listener = cursor;
- cursor = cursor._nextListener;
- listener._nextListener = null;
- resultSource._addListener(listener);
- }
- // Listen with this future as well, so that when the other future completes,
- // this future will be completed as well.
- resultSource._addListener(this._asListener());
- _resultOrListeners = resultSource;
- _state = hadListeners ? _CHAINED : _CHAINED_UNLISTENED;
- }
-
- /**
- * Helper function to handle the result of transforming an incoming event.
- *
- * If the result is itself a [Future], this future is linked to that
- * future's output. If not, this future is completed with the result.
- */
- void _setOrChainValue(var result) {
- assert(!_isChained);
- assert(!_isComplete);
- if (result is Future) {
- // Result should be a Future<T>.
- if (result is _FutureImpl) {
- _FutureImpl chainFuture = result;
- chainFuture._chain(this);
- return;
+ // Mark the target as chained (and as such half-completed).
+ target._isChained = true;
+ if (source is _Future) {
+ _Future internalFuture = source;
+ if (internalFuture._isComplete) {
+ _propagateToListeners(internalFuture, target);
} else {
- Future future = result;
- future.then(_setValue,
- onError: _setError);
- return;
+ internalFuture._addListener(target);
}
} else {
- // Result must be of type T.
- _setValue(result);
+ source.then((value) {
+ assert(target._isChained);
+ target._complete(value);
+ },
+ onError: (error) {
+ assert(target._isChained);
+ target._completeError(error);
+ });
}
}
- _FutureListener _asListener() => new _FutureListener.wrap(this);
-}
+ void _complete(value) {
+ assert(!_isComplete);
+ assert(_onValue == null);
+ assert(_onError == null);
+ assert(_whenCompleteAction == null);
+ assert(_errorTest == null);
-/**
- * Transforming future base class.
- *
- * A transforming future is itself a future and a future listener.
- * Subclasses override [_sendValue]/[_sendError] to intercept
- * the results of a previous future.
- */
-abstract class _TransformFuture<S, T> extends _FutureImpl<T>
- implements _FutureListener<S> {
- // _FutureListener implementation.
- _FutureListener _nextListener;
-
- _TransformFuture() {
- _zone.expectCallback();
+ if (value is Future) {
+ _chainFutures(value, this);
+ return;
+ }
+ _Future listeners = _removeListeners();
+ _setValue(value);
+ _propagateToListeners(this, listeners);
}
- void _sendValue(S value) {
- _zone.executeCallback(() => _zonedSendValue(value));
- }
+ void _completeError(error, [StackTrace stackTrace]) {
+ assert(!_isComplete);
+ assert(_onValue == null);
+ assert(_onError == null);
+ assert(_whenCompleteAction == null);
+ assert(_errorTest == null);
- void _sendError(error) {
- _zone.executeCallback(() => _zonedSendError(error));
- }
+ if (stackTrace != null) {
+ // Force the stack trace onto the error, even if it already had one.
+ _attachStackTrace(error, stackTrace);
+ }
- void _subscribeTo(_FutureImpl future) {
- future._addListener(this);
+ _Future listeners = _isChained ? null : _removeListeners();
+ _setError(error);
+ _propagateToListeners(this, listeners);
}
- void _zonedSendValue(S value);
- void _zonedSendError(error);
-}
-
-/** The onValue and onError handlers return either a value or a future */
-typedef dynamic _FutureOnValue<T>(T value);
-typedef dynamic _FutureOnError(error);
-/** Test used by [Future.catchError] to handle skip some errors. */
-typedef bool _FutureErrorTest(var error);
-/** Used by [WhenFuture]. */
-typedef _FutureAction();
-
-/** Future returned by [Future.then] with no [:onError:] parameter. */
-class _ThenFuture<S, T> extends _TransformFuture<S, T> {
- // TODO(ahe): Restore type when feature is implemented in dart2js
- // checked mode.
- final /* _FutureOnValue<S> */ _onValue;
-
- _ThenFuture(this._onValue);
-
- _zonedSendValue(S value) {
- assert(_onValue != null);
- var result;
- try {
- result = _onValue(value);
- } catch (e, s) {
- _setError(_asyncError(e, s));
+ void _asyncComplete(value) {
+ assert(!_isComplete);
+ assert(_onValue == null);
+ assert(_onError == null);
+ assert(_whenCompleteAction == null);
+ assert(_errorTest == null);
+ // Two corner cases if the value is a future:
+ // 1. the future is already completed and an error.
+ // 2. the future is not yet completed but might become an error.
+ // The first case means that we must not immediately complete the Future,
+ // as our code would immediately start propagating the error without
+ // giving the time to install error-handlers.
+ // However the second case requires us to deal with the value immediately.
+ // Otherwise the value could complete with an error and report an
+ // unhandled error, even though we know we are already going to listen to
+ // it.
+ if (value is Future &&
+ (value is! _Future || !(value as _Future)._isComplete)) {
+ // Case 2 from above. We need to register.
+ // Note that we are still completing asynchronously: either we register
+ // through .then (in which case the completing is asynchronous), or we
+ // have a _Future which isn't complete yet.
+ _complete(value);
return;
}
- _setOrChainValue(result);
- }
- void _zonedSendError(error) {
- _setError(error);
+ _markPendingCompletion();
+ runAsync(() {
+ _complete(value);
+ });
}
-}
-/** Future returned by [Future.catchError]. */
-class _CatchErrorFuture<T> extends _TransformFuture<T,T> {
- final _FutureErrorTest _test;
- final _FutureOnError _onError;
+ void _asyncCompleteError(error, [StackTrace stackTrace]) {
+ assert(!_isComplete);
+ assert(_onValue == null);
+ assert(_onError == null);
+ assert(_whenCompleteAction == null);
+ assert(_errorTest == null);
- _CatchErrorFuture(this._onError, this._test);
+ _markPendingCompletion();
+ runAsync(() {
+ _completeError(error, stackTrace);
+ });
+ }
- _zonedSendValue(T value) {
- _setValue(value);
+ /**
+ * Propagates the value/error of [source] to its [listeners].
+ *
+ * Unlinks all listeners and propagates the source to each listener
+ * separately.
+ */
+ static void _propagateMultipleListeners(_Future source, _Future listeners) {
+ assert(listeners != null);
+ assert(listeners._nextListener != null);
+ do {
+ _Future listener = listeners;
+ listeners = listener._nextListener;
+ listener._nextListener = null;
+ _propagateToListeners(source, listener);
+ } while (listeners != null);
}
- _zonedSendError(error) {
- assert(_onError != null);
- // if _test is supplied, check if it returns true, otherwise just
- // forward the error unmodified.
- if (_test != null) {
- bool matchesTest;
- try {
- matchesTest = _test(error);
- } catch (e, s) {
- _setError(_asyncError(e, s));
+ /**
+ * Propagates the value/error of [source] to its [listeners], executing the
+ * listeners' callbacks.
+ *
+ * If [runCallback] is true (which should be the default) it executes
+ * the registered action of listeners. If it is `false` then the callback is
+ * skipped. This is used to complete futures with chained futures.
+ */
+ static void _propagateToListeners(_Future source, _Future listeners) {
+ while (true) {
+ if (!source._isComplete) return; // Chained future.
+ bool hasError = source._hasError;
+ if (hasError && listeners == null) {
+ source._zone.handleUncaughtError(source._error);
return;
}
- if (!matchesTest) {
- _setError(error);
+ if (listeners == null) return;
+ _Future listener = listeners;
+ if (listener._nextListener != null) {
+ // Usually futures only have one listener. If they have several, we
+ // handle them specially.
+ _propagateMultipleListeners(source, listeners);
return;
}
- }
- // Act on the error, and use the result as this future's result.
- var result;
- try {
- result = _onError(error);
- } catch (e, s) {
- _setError(_asyncError(e, s));
- return;
- }
- _setOrChainValue(result);
- }
-}
-
-/** Future returned by [Future.then] with an [:onError:] parameter. */
-class _SubscribeFuture<S, T> extends _ThenFuture<S, T> {
- final _FutureOnError _onError;
-
- _SubscribeFuture(onValue(S value), this._onError) : super(onValue);
-
- // The _sendValue method is inherited from ThenFuture.
-
- void _zonedSendError(error) {
- assert(_onError != null);
- var result;
- try {
- result = _onError(error);
- } catch (e, s) {
- _setError(_asyncError(e, s));
- return;
- }
- _setOrChainValue(result);
- }
-}
-
-/** Future returned by [Future.whenComplete]. */
-class _WhenFuture<T> extends _TransformFuture<T, T> {
- final _FutureAction _action;
-
- _WhenFuture(this._action);
-
- void _zonedSendValue(T value) {
- try {
- var result = _action();
- if (result is Future) {
- Future resultFuture = result;
- resultFuture.then((_) {
- _setValue(value);
- }, onError: _setError);
+ if (hasError && !source._zone.inSameErrorZone(listener._zone)) {
+ // Don't cross zone boundaries with errors.
+ source._zone.handleUncaughtError(source._error);
+ return;
+ }
+ if (!identical(_Zone.current, listener._zone)) {
+ // Run the propagation in the listener's zone to avoid
+ // zone transitions. The idea is that many chained futures will
+ // be in the same zone.
+ listener._zone.executePeriodicCallback(() {
+ _propagateToListeners(source, listener);
+ });
return;
}
- } catch (e, s) {
- _setError(_asyncError(e, s));
- return;
- }
- _setValue(value);
- }
- void _zonedSendError(error) {
- try {
- var result = _action();
- if (result is Future) {
- Future resultFuture = result;
- // TODO(lrn): Find a way to combine [error] into [e].
- resultFuture.then((_) {
- _setError(error);
- }, onError: _setError);
+ // Do the actual propagation.
+ // TODO(floitsch): Do we need to go through the zone even if we
+ // don't have a callback to execute?
+ bool listenerHasValue;
+ var listenerValueOrError;
+ // Set to true if a whenComplete needs to wait for a future.
+ // The whenComplete action will resume the propagation by itself.
+ bool isPropagationAborted = false;
+ // Even though we are already in the right zone (due to the optimization
+ // above), we still need to go through the zone. The overhead of
+ // executeCallback is however smaller when it is already in the correct
+ // zone.
+ // TODO(floitsch): only run callbacks in the zone, not the whole
+ // handling code.
+ listener._zone.executeCallback(() {
+ // TODO(floitsch): mark the listener as pending completion. Currently
+ // we can't do this, since the markPendingCompletion verifies that
+ // the future is not already marked (or chained).
+ try {
+ if (!hasError) {
+ var value = source._value;
+ if (listener._onValue != null) {
+ listenerValueOrError = listener._onValue(value);
+ listenerHasValue = true;
+ } else {
+ // Copy over the value from the source.
+ listenerValueOrError = value;
+ listenerHasValue = true;
+ }
+ } else {
+ Object error = source._error;
+ _FutureErrorTest test = listener._errorTest;
+ bool matchesTest = true;
+ if (test != null) {
+ matchesTest = test(error);
+ }
+ if (matchesTest && listener._onError != null) {
+ listenerValueOrError = listener._onError(error);
+ listenerHasValue = true;
+ } else {
+ // Copy over the error from the source.
+ listenerValueOrError = error;
+ listenerHasValue = false;
+ }
+ }
+
+ if (listener._whenCompleteAction != null) {
+ var completeResult = listener._whenCompleteAction();
+ if (completeResult is Future) {
+ listener._isChained = true;
+ completeResult.then((ignored) {
+ // Try again, but this time don't run the whenComplete callback.
+ _propagateToListeners(source, listener);
+ }, onError: (error) {
+ // When there is an error, we have to make the error the new
+ // result of the current listener.
+ if (completeResult is! _Future) {
+ // This should be a rare case.
+ completeResult = new _Future();
+ completeResult._setError(error);
+ }
+ _propagateToListeners(completeResult, listener);
+ });
+ isPropagationAborted = true;
+ // We will reenter the listener's zone.
+ listener._zone.expectCallback();
+ }
+ }
+ } catch (e, s) {
+ // Set the exception as error.
+ listenerValueOrError = _asyncError(e, s);
+ listenerHasValue = false;
+ }
+ if (listenerHasValue && listenerValueOrError is Future) {
+ // We are going to reenter the zone to finish what we started.
+ listener._zone.expectCallback();
+ }
+ });
+ if (isPropagationAborted) return;
+ // If the listener's value is a future we need to chain it.
+ if (listenerHasValue && listenerValueOrError is Future) {
+ Future chainSource = listenerValueOrError;
+ // Shortcut if the chain-source is already completed. Just continue the
+ // loop.
+ if (chainSource is _Future && (chainSource as _Future)._isComplete) {
+ // propagate the value (simulating a tail call).
+ listener._isChained = true;
+ source = chainSource;
+ listeners = listener;
+ continue;
+ }
+ _chainFutures(chainSource, listener);
return;
}
- } catch (e, s) {
- error = _asyncError(e, s);
+
+ if (listenerHasValue) {
+ listeners = listener._removeListeners();
+ listener._setValue(listenerValueOrError);
+ } else {
+ listeners = listener._removeListeners();
+ listener._setError(listenerValueOrError);
+ }
+ // Prepare for next round.
+ source = listener;
}
- _setError(error);
}
}
« no previous file with comments | « sdk/lib/async/future.dart ('k') | sdk/lib/async/stream.dart » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698