| Index: sdk/lib/async/future_impl.dart
|
| diff --git a/sdk/lib/async/future_impl.dart b/sdk/lib/async/future_impl.dart
|
| index f3cd5e41c6b9b8ee4e6e69f0d14e1753b6652519..a0cfeb16efe9b05e4778085a0988824e46168236 100644
|
| --- a/sdk/lib/async/future_impl.dart
|
| +++ b/sdk/lib/async/future_impl.dart
|
| @@ -57,21 +57,82 @@ class _SyncCompleter<T> extends _Completer<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.
|
| +class _FutureListener {
|
| + static const int MASK_VALUE = 1;
|
| + static const int MASK_ERROR = 2;
|
| + static const int MASK_TEST_ERROR = 4;
|
| + static const int MASK_WHENCOMPLETE = 8;
|
| + static const int STATE_CHAIN = 0;
|
| + static const int STATE_THEN = MASK_VALUE;
|
| + static const int STATE_THEN_ONERROR = MASK_VALUE | MASK_ERROR;
|
| + static const int STATE_CATCHERROR = MASK_ERROR;
|
| + static const int STATE_CATCHERROR_TEST = MASK_ERROR | MASK_TEST_ERROR;
|
| + static const int STATE_WHENCOMPLETE = MASK_WHENCOMPLETE;
|
| + // Listeners on the same future are linked through this link.
|
| + _FutureListener _nextListener = null;
|
| + // The future to complete when this listener is activated.
|
| + final _Future result;
|
| + // Which fields means what.
|
| + final int state;
|
| + // Used for then/whenDone callback and error test
|
| + final Function callback;
|
| + // Used for error callbacks.
|
| + final Function errorCallback;
|
| +
|
| + _FutureListener.then(this.result,
|
| + _FutureOnValue onValue, Function errorCallback)
|
| + : callback = onValue,
|
| + errorCallback = errorCallback,
|
| + state = (errorCallback == null) ? STATE_THEN : STATE_THEN_ONERROR;
|
| +
|
| + _FutureListener.catchError(this.result,
|
| + this.errorCallback, _FutureErrorTest test)
|
| + : callback = test,
|
| + state = (test == null) ? STATE_CATCHERROR : STATE_CATCHERROR_TEST;
|
| +
|
| + _FutureListener.whenComplete(this.result, _FutureAction onComplete)
|
| + : callback = onComplete,
|
| + errorCallback = null,
|
| + state = STATE_WHENCOMPLETE;
|
| +
|
| + _FutureListener.chain(this.result)
|
| + : callback = null,
|
| + errorCallback = null,
|
| + state = STATE_CHAIN;
|
| +
|
| + Zone get _zone => result._zone;
|
| +
|
| + bool get handlesValue => (state & MASK_VALUE != 0);
|
| + bool get handlesError => (state & MASK_ERROR != 0);
|
| + bool get hasErrorTest => (state == STATE_CATCHERROR_TEST);
|
| + bool get handlesComplete => (state == STATE_WHENCOMPLETE);
|
| +
|
| + _FutureOnValue get _onValue {
|
| + assert(handlesValue);
|
| + return callback;
|
| + }
|
| + Function get _onError => errorCallback;
|
| + _FutureErrorTest get _errorTest {
|
| + assert(hasErrorTest);
|
| + return callback;
|
| + }
|
| + _FutureAction get _whenCompleteAction {
|
| + assert(handlesComplete);
|
| + return callback;
|
| + }
|
| +}
|
|
|
| +class _Future<T> implements Future<T> {
|
| /// Initial state, waiting for a result. In this state, the
|
| /// [resultOrListeners] field holds a single-linked list of
|
| - /// [FutureListener] listeners.
|
| + /// [_FutureListener] listeners.
|
| static const int _INCOMPLETE = 0;
|
| /// Pending completion. Set when completed using [_asyncComplete] or
|
| /// [_asyncCompleteError]. It is an error to try to complete it again.
|
| + /// [resultOrListeners] holds listeners.
|
| 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, 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;
|
| @@ -83,23 +144,14 @@ class _Future<T> implements Future<T> {
|
| /** Whether the future is complete, and as what. */
|
| int _state = _INCOMPLETE;
|
|
|
| - final Zone _zone;
|
| -
|
| - 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;
|
| -
|
| - set _isChained(bool value) {
|
| - if (value) {
|
| - assert(!_isComplete);
|
| - _state = _CHAINED;
|
| - } else {
|
| - assert(_isChained);
|
| - _state = _INCOMPLETE;
|
| - }
|
| - }
|
| + /**
|
| + * Zone that the future was completed from.
|
| + * This is the zone that an error result belongs to.
|
| + *
|
| + * Until the future is completed, the field may hold the zone that
|
| + * listener callbacks used to create this future should be run in.
|
| + */
|
| + final Zone _zone = Zone.current;
|
|
|
| /**
|
| * Either the result, a list of listeners or another future.
|
| @@ -116,95 +168,68 @@ class _Future<T> implements Future<T> {
|
| * will complete with the same result.
|
| * All listeners are forwarded to the other future.
|
| *
|
| - * The cases are disjoint (incomplete and unchained, incomplete and
|
| - * chained, or completed with value or error), so the field only needs to hold
|
| + * The cases are disjoint - incomplete and unchained ([_INCOMPLETE]),
|
| + * incomplete and chained ([_CHAINED]), or completed with value or error
|
| + * ([_VALUE] or [_ERROR]) - so the field only needs to hold
|
| * one value at a time.
|
| */
|
| var _resultOrListeners;
|
|
|
| - /**
|
| - * 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.
|
| - _FutureOnValue _onValueCallback;
|
| - _FutureErrorTest _errorTestCallback;
|
| - Function _onErrorCallback;
|
| - _FutureAction _whenCompleteActionCallback;
|
| -
|
| - _FutureOnValue get _onValue => _isChained ? null : _onValueCallback;
|
| - _FutureErrorTest get _errorTest => _isChained ? null : _errorTestCallback;
|
| - Function get _onError => _isChained ? null : _onErrorCallback;
|
| - _FutureAction get _whenCompleteAction
|
| - => _isChained ? null : _whenCompleteActionCallback;
|
| -
|
| - _Future()
|
| - : _zone = Zone.current,
|
| - _onValueCallback = null, _errorTestCallback = null,
|
| - _onErrorCallback = null, _whenCompleteActionCallback = null;
|
| + _Future();
|
|
|
| /// Valid types for value: `T` or `Future<T>`.
|
| - _Future.immediate(value)
|
| - : _zone = Zone.current,
|
| - _onValueCallback = null, _errorTestCallback = null,
|
| - _onErrorCallback = null, _whenCompleteActionCallback = null {
|
| + _Future.immediate(value) {
|
| _asyncComplete(value);
|
| }
|
|
|
| - _Future.immediateError(var error, [StackTrace stackTrace])
|
| - : _zone = Zone.current,
|
| - _onValueCallback = null, _errorTestCallback = null,
|
| - _onErrorCallback = null, _whenCompleteActionCallback = null {
|
| + _Future.immediateError(var error, [StackTrace stackTrace]) {
|
| _asyncCompleteError(error, stackTrace);
|
| }
|
|
|
| - _Future._then(onValueCallback(value), Function onErrorCallback)
|
| - : _zone = Zone.current,
|
| - _onValueCallback = Zone.current.registerUnaryCallback(onValueCallback),
|
| - _onErrorCallback = _registerErrorHandler(onErrorCallback, Zone.current),
|
| - _errorTestCallback = null,
|
| - _whenCompleteActionCallback = null;
|
| -
|
| - _Future._catchError(Function onErrorCallback, bool errorTestCallback(e))
|
| - : _zone = Zone.current,
|
| - _onErrorCallback = _registerErrorHandler(onErrorCallback, Zone.current),
|
| - _errorTestCallback =
|
| - Zone.current.registerUnaryCallback(errorTestCallback),
|
| - _onValueCallback = null,
|
| - _whenCompleteActionCallback = null;
|
| -
|
| - _Future._whenComplete(whenCompleteActionCallback())
|
| - : _zone = Zone.current,
|
| - _whenCompleteActionCallback =
|
| - Zone.current.registerCallback(whenCompleteActionCallback),
|
| - _onValueCallback = null,
|
| - _errorTestCallback = null,
|
| - _onErrorCallback = null;
|
| + 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;
|
| +
|
| + set _isChained(bool value) {
|
| + if (value) {
|
| + assert(!_isComplete);
|
| + _state = _CHAINED;
|
| + } else {
|
| + assert(_isChained);
|
| + _state = _INCOMPLETE;
|
| + }
|
| + }
|
|
|
| Future then(f(T value), { Function onError }) {
|
| - _Future result;
|
| - result = new _Future._then(f, onError);
|
| - _addListener(result);
|
| + _Future result = new _Future();
|
| + if (!identical(result._zone, _ROOT_ZONE)) {
|
| + f = result._zone.registerUnaryCallback(f);
|
| + if (onError != null) {
|
| + onError = _registerErrorHandler(onError, result._zone);
|
| + }
|
| + }
|
| + _addListener(new _FutureListener.then(result, f, onError));
|
| return result;
|
| }
|
|
|
| Future catchError(Function onError, { bool test(error) }) {
|
| - _Future result = new _Future._catchError(onError, test);
|
| - _addListener(result);
|
| + _Future result = new _Future();
|
| + if (!identical(result._zone, _ROOT_ZONE)) {
|
| + onError = _registerErrorHandler(onError, result._zone);
|
| + if (test != null) test = result._zone.registerUnaryCallback(test);
|
| + }
|
| + _addListener(new _FutureListener.catchError(result, onError, test));
|
| return result;
|
| }
|
|
|
| Future<T> whenComplete(action()) {
|
| - _Future result = new _Future<T>._whenComplete(action);
|
| - _addListener(result);
|
| + _Future result = new _Future<T>();
|
| + if (!identical(result._zone, _ROOT_ZONE)) {
|
| + action = result._zone.registerCallback(action);
|
| + }
|
| + _addListener(new _FutureListener.whenComplete(result, action));
|
| return result;
|
| }
|
|
|
| @@ -231,13 +256,17 @@ class _Future<T> implements Future<T> {
|
| _resultOrListeners = value;
|
| }
|
|
|
| - void _setError(Object error, StackTrace stackTrace) {
|
| + void _setErrorObject(AsyncError error) {
|
| assert(!_isComplete); // But may have a completion pending.
|
| _state = _ERROR;
|
| - _resultOrListeners = new AsyncError(error, stackTrace);
|
| + _resultOrListeners = error;
|
| }
|
|
|
| - void _addListener(_Future listener) {
|
| + void _setError(Object error, StackTrace stackTrace) {
|
| + _setErrorObject(new AsyncError(error, stackTrace));
|
| + }
|
| +
|
| + void _addListener(_FutureListener listener) {
|
| assert(listener._nextListener == null);
|
| if (_isComplete) {
|
| // Handle late listeners asynchronously.
|
| @@ -250,15 +279,15 @@ class _Future<T> implements Future<T> {
|
| }
|
| }
|
|
|
| - _Future _removeListeners() {
|
| + _FutureListener _removeListeners() {
|
| // Reverse listeners before returning them, so the resulting list is in
|
| // subscription order.
|
| assert(!_isComplete);
|
| - _Future current = _resultOrListeners;
|
| + _FutureListener current = _resultOrListeners;
|
| _resultOrListeners = null;
|
| - _Future prev = null;
|
| + _FutureListener prev = null;
|
| while (current != null) {
|
| - _Future next = current._nextListener;
|
| + _FutureListener next = current._nextListener;
|
| current._nextListener = prev;
|
| prev = current;
|
| current = next;
|
| @@ -297,21 +326,16 @@ class _Future<T> implements Future<T> {
|
|
|
| // Mark the target as chained (and as such half-completed).
|
| target._isChained = true;
|
| - _Future internalFuture = source;
|
| - if (internalFuture._isComplete) {
|
| - _propagateToListeners(internalFuture, target);
|
| + _FutureListener listener = new _FutureListener.chain(target);
|
| + if (source._isComplete) {
|
| + _propagateToListeners(source, listener);
|
| } else {
|
| - internalFuture._addListener(target);
|
| + source._addListener(listener);
|
| }
|
| }
|
|
|
| void _complete(value) {
|
| assert(!_isComplete);
|
| - assert(_onValue == null);
|
| - assert(_onError == null);
|
| - assert(_whenCompleteAction == null);
|
| - assert(_errorTest == null);
|
| -
|
| if (value is Future) {
|
| if (value is _Future) {
|
| _chainCoreFuture(value, this);
|
| @@ -319,7 +343,7 @@ class _Future<T> implements Future<T> {
|
| _chainForeignFuture(value, this);
|
| }
|
| } else {
|
| - _Future listeners = _removeListeners();
|
| + _FutureListener listeners = _removeListeners();
|
| _setValue(value);
|
| _propagateToListeners(this, listeners);
|
| }
|
| @@ -327,35 +351,23 @@ class _Future<T> implements Future<T> {
|
|
|
| void _completeWithValue(value) {
|
| assert(!_isComplete);
|
| - assert(_onValue == null);
|
| - assert(_onError == null);
|
| - assert(_whenCompleteAction == null);
|
| - assert(_errorTest == null);
|
| assert(value is! Future);
|
|
|
| - _Future listeners = _removeListeners();
|
| + _FutureListener listeners = _removeListeners();
|
| _setValue(value);
|
| _propagateToListeners(this, listeners);
|
| }
|
|
|
| void _completeError(error, [StackTrace stackTrace]) {
|
| assert(!_isComplete);
|
| - assert(_onValue == null);
|
| - assert(_onError == null);
|
| - assert(_whenCompleteAction == null);
|
| - assert(_errorTest == null);
|
|
|
| - _Future listeners = _removeListeners();
|
| + _FutureListener listeners = _removeListeners();
|
| _setError(error, stackTrace);
|
| _propagateToListeners(this, listeners);
|
| }
|
|
|
| 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.
|
| @@ -403,10 +415,6 @@ class _Future<T> implements Future<T> {
|
|
|
| void _asyncCompleteError(error, StackTrace stackTrace) {
|
| assert(!_isComplete);
|
| - assert(_onValue == null);
|
| - assert(_onError == null);
|
| - assert(_whenCompleteAction == null);
|
| - assert(_errorTest == null);
|
|
|
| _markPendingCompletion();
|
| _zone.scheduleMicrotask(() {
|
| @@ -415,53 +423,36 @@ class _Future<T> implements Future<T> {
|
| }
|
|
|
| /**
|
| - * 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);
|
| - }
|
| -
|
| - /**
|
| * 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) {
|
| + static void _propagateToListeners(_Future source, _FutureListener listeners) {
|
| while (true) {
|
| - if (!source._isComplete) return; // Chained future.
|
| + assert(source._isComplete);
|
| bool hasError = source._hasError;
|
| - if (hasError && listeners == null) {
|
| - AsyncError asyncError = source._error;
|
| - source._zone.handleUncaughtError(
|
| - asyncError.error, asyncError.stackTrace);
|
| + if (listeners == null) {
|
| + if (hasError) {
|
| + AsyncError asyncError = source._error;
|
| + source._zone.handleUncaughtError(
|
| + asyncError.error, asyncError.stackTrace);
|
| + }
|
| return;
|
| }
|
| - 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;
|
| + // Usually futures only have one listener. If they have several, we
|
| + // call handle them separately in recursive calls, continuing
|
| + // here only when there is only one listener left.
|
| + while (listeners._nextListener != null) {
|
| + _FutureListener listener = listeners;
|
| + listeners = listener._nextListener;
|
| + listener._nextListener = null;
|
| + _propagateToListeners(source, listener);
|
| }
|
| + _FutureListener listener = listeners;
|
| // Do the actual propagation.
|
| // Set initial state of listenerHasValue and listenerValueOrError. These
|
| // variables are updated, with the outcome of potential callbacks.
|
| bool listenerHasValue = true;
|
| - final sourceValue = source._hasValue ? source._value : null;
|
| + final sourceValue = hasError ? null : source._value;
|
| var listenerValueOrError = sourceValue;
|
| // Set to true if a whenComplete needs to wait for a future.
|
| // The whenComplete action will resume the propagation by itself.
|
| @@ -472,9 +463,7 @@ class _Future<T> implements Future<T> {
|
| // Only if we either have an error or callbacks, go into this, somewhat
|
| // expensive, branch. Here we'll enter/leave the zone. Many futures
|
| // doesn't have callbacks, so this is a significant optimization.
|
| - if (hasError ||
|
| - listener._onValue != null ||
|
| - listener._whenCompleteAction != null) {
|
| + if (hasError || (listener.handlesValue || listener.handlesComplete)) {
|
| Zone zone = listener._zone;
|
| if (hasError && !source._zone.inSameErrorZone(zone)) {
|
| // Don't cross zone boundaries with errors.
|
| @@ -503,13 +492,12 @@ class _Future<T> implements Future<T> {
|
|
|
| void handleError() {
|
| AsyncError asyncError = source._error;
|
| - _FutureErrorTest test = listener._errorTest;
|
| bool matchesTest = true;
|
| - if (test != null) {
|
| + if (listener.hasErrorTest) {
|
| + _FutureErrorTest test = listener._errorTest;
|
| try {
|
| matchesTest = zone.runUnary(test, asyncError.error);
|
| } catch (e, s) {
|
| - // TODO(ajohnsen): Should we suport rethrow for test throws?
|
| listenerValueOrError = identical(asyncError.error, e) ?
|
| asyncError : new AsyncError(e, s);
|
| listenerHasValue = false;
|
| @@ -552,14 +540,14 @@ class _Future<T> implements Future<T> {
|
| listenerValueOrError = new AsyncError(e, s);
|
| }
|
| listenerHasValue = false;
|
| + return;
|
| }
|
| if (completeResult is Future) {
|
| - listener._isChained = true;
|
| + _Future result = listener.result;
|
| + result._isChained = true;
|
| isPropagationAborted = true;
|
| completeResult.then((ignored) {
|
| - // Try again. Since the future is marked as chained it won't run
|
| - // the whenComplete again.
|
| - _propagateToListeners(source, listener);
|
| + _propagateToListeners(source, new _FutureListener.chain(result));
|
| }, onError: (error, [stackTrace]) {
|
| // When there is an error, we have to make the error the new
|
| // result of the current listener.
|
| @@ -568,27 +556,24 @@ class _Future<T> implements Future<T> {
|
| completeResult = new _Future();
|
| completeResult._setError(error, stackTrace);
|
| }
|
| - _propagateToListeners(completeResult, listener);
|
| + _propagateToListeners(completeResult,
|
| + new _FutureListener.chain(result));
|
| });
|
| }
|
| }
|
|
|
| if (!hasError) {
|
| - if (listener._onValue != null) {
|
| + if (listener.handlesValue) {
|
| listenerHasValue = handleValueCallback();
|
| }
|
| } else {
|
| handleError();
|
| }
|
| - if (listener._whenCompleteAction != null) {
|
| + if (listener.handlesComplete) {
|
| handleWhenCompleteCallback();
|
| }
|
| // If we changed zone, oldZone will not be null.
|
| if (oldZone != null) Zone._leave(oldZone);
|
| - listener._onValueCallback = null;
|
| - listener._errorTestCallback = null;
|
| - listener._onErrorCallback = null;
|
| - listener._whenCompleteActionCallback = null;
|
|
|
| if (isPropagationAborted) return;
|
| // If the listener's value is a future we need to chain it. Note that
|
| @@ -600,32 +585,33 @@ class _Future<T> implements Future<T> {
|
| Future chainSource = listenerValueOrError;
|
| // Shortcut if the chain-source is already completed. Just continue
|
| // the loop.
|
| + _Future result = listener.result;
|
| if (chainSource is _Future) {
|
| if (chainSource._isComplete) {
|
| // propagate the value (simulating a tail call).
|
| - listener._isChained = true;
|
| + result._isChained = true;
|
| source = chainSource;
|
| - listeners = listener;
|
| + listeners = new _FutureListener.chain(result);
|
| continue;
|
| } else {
|
| - _chainCoreFuture(chainSource, listener);
|
| + _chainCoreFuture(chainSource, result);
|
| }
|
| } else {
|
| - _chainForeignFuture(chainSource, listener);
|
| + _chainForeignFuture(chainSource, result);
|
| }
|
| return;
|
| }
|
| }
|
| + _Future result = listener.result;
|
| + listeners = result._removeListeners();
|
| if (listenerHasValue) {
|
| - listeners = listener._removeListeners();
|
| - listener._setValue(listenerValueOrError);
|
| + result._setValue(listenerValueOrError);
|
| } else {
|
| - listeners = listener._removeListeners();
|
| AsyncError asyncError = listenerValueOrError;
|
| - listener._setError(asyncError.error, asyncError.stackTrace);
|
| + result._setErrorObject(asyncError);
|
| }
|
| // Prepare for next round.
|
| - source = listener;
|
| + source = result;
|
| }
|
| }
|
|
|
|
|