Chromium Code Reviews| Index: pkg/scheduled_test/lib/src/schedule.dart |
| diff --git a/pkg/scheduled_test/lib/src/schedule.dart b/pkg/scheduled_test/lib/src/schedule.dart |
| index e0ce6fc2cf23684df619d1ed5ece4657d3446dd5..1f600bcd9eb3ed3e6ab3958aed1fcffb2ff08461 100644 |
| --- a/pkg/scheduled_test/lib/src/schedule.dart |
| +++ b/pkg/scheduled_test/lib/src/schedule.dart |
| @@ -101,15 +101,6 @@ class Schedule { |
| heartbeat(); |
| } |
| - /// The number of out-of-band callbacks that have been registered with |
| - /// [wrapAsync] but have yet to be called. |
| - int _pendingCallbacks = 0; |
| - |
| - /// A completer that will be completed once [_pendingCallbacks] reaches zero. |
| - /// This will only be non-`null` if [_awaitPendingCallbacks] has been called |
| - /// while [_pendingCallbacks] is non-zero. |
| - Completer _noPendingCallbacks; |
| - |
| /// The timer for keeping track of task timeouts. This may be null. |
| Timer _timeoutTimer; |
| @@ -214,39 +205,19 @@ class Schedule { |
| /// the returned function has been called. It's used to ensure that |
| /// out-of-band callbacks are properly handled by the scheduled test. |
| /// |
| + /// [description] provides an optional description of the callback, which is |
| + /// used when generating error messages. |
| + /// |
| /// The top-level `wrapAsync` function should usually be used in preference to |
| /// this in test code. |
| - Function wrapAsync(fn(arg)) { |
| + Function wrapAsync(fn(arg), [String description]) { |
| if (_state == ScheduleState.DONE) { |
| throw new StateError("wrapAsync called after the schedule has finished " |
| "running."); |
| } |
| heartbeat(); |
| - var queue = currentQueue; |
| - // It's possible that the queue timed out before this. |
| - bool _timedOut() => queue != currentQueue || _pendingCallbacks == 0; |
| - |
| - _pendingCallbacks++; |
| - return (arg) { |
| - try { |
| - return fn(arg); |
| - } catch (e, stackTrace) { |
| - if (_timedOut()) { |
| - _signalPostTimeoutError(e, stackTrace); |
| - } else { |
| - signalError(e, stackTrace); |
| - } |
| - } finally { |
| - if (_timedOut()) return; |
| - |
| - _pendingCallbacks--; |
| - if (_pendingCallbacks == 0 && _noPendingCallbacks != null) { |
| - _noPendingCallbacks.complete(); |
| - _noPendingCallbacks = null; |
| - } |
| - } |
| - }; |
| + return currentQueue._wrapAsync(fn, description); |
| } |
| /// Like [wrapAsync], this ensures that the current task queue waits for |
| @@ -256,18 +227,20 @@ class Schedule { |
| /// |
| /// The returned [Future] completes to the same value or error as [future]. |
| /// |
| + /// [description] provides an optional description of the future, which is |
| + /// used when generating error messages. |
| + /// |
| /// The top-level `wrapFuture` function should usually be used in preference |
| /// to this in test code. |
| - Future wrapFuture(Future future) { |
| - var doneCallback = wrapAsync((_) => null); |
| - done() => new Future.immediate(null).then(doneCallback); |
| + Future wrapFuture(Future future, [String description]) { |
| + var done = wrapAsync((fn) => fn(), description); |
| - future = future.then((result) { |
| - done(); |
| - return result; |
| - }).catchError((e) { |
| - signalError(e); |
| - done(); |
| + future = future.then((result) => done(() => result)).catchError((e) { |
| + done(() { |
| + throw e; |
| + }); |
| + // wrapAsync will catch the first throw, so we throw [e] again so it |
| + // propagates through the Future chain. |
| throw e; |
| }); |
| @@ -303,38 +276,14 @@ class Schedule { |
| if (_timeout == null) { |
| _timeoutTimer = null; |
| } else { |
| - _timeoutTimer = mock_clock.newTimer(_timeout, _signalTimeout); |
| - } |
| - } |
| - |
| - /// The callback to run when the timeout timer fires. Notifies the current |
| - /// queue that a timeout has occurred. |
| - void _signalTimeout() { |
| - // Reset the timer so that we can detect timeouts in the onException and |
| - // onComplete queues. |
| - _timeoutTimer = null; |
| - |
| - var error = new ScheduleError.from(this, "The schedule timed out after " |
| - "$_timeout of inactivity."); |
| - |
| - _pendingCallbacks = 0; |
| - if (_noPendingCallbacks != null) { |
| - var noPendingCallbacks = _noPendingCallbacks; |
| - _noPendingCallbacks = null; |
| - noPendingCallbacks.completeError(error); |
| - } else { |
| - currentQueue._signalTimeout(error); |
| + _timeoutTimer = mock_clock.newTimer(_timeout, () { |
| + _timeoutTimer = null; |
| + currentQueue._signalTimeout(new ScheduleError.from(this, "The schedule " |
| + "timed out after $_timeout of inactivity.")); |
| + }); |
| } |
| } |
| - /// Returns a [Future] that will complete once there are no pending |
| - /// out-of-band callbacks. |
| - Future _awaitNoPendingCallbacks() { |
| - if (_pendingCallbacks == 0) return new Future.immediate(null); |
| - if (_noPendingCallbacks == null) _noPendingCallbacks = new Completer(); |
| - return _noPendingCallbacks.future; |
| - } |
| - |
| /// Register an error in the schedule's error list. This ensures that there |
| /// are no duplicate errors, and that all errors are wrapped in |
| /// [ScheduleError]. |
| @@ -388,6 +337,21 @@ class TaskQueue { |
| /// null if no task is currently running. |
| SubstituteFuture _taskFuture; |
| + /// The toal number of out-of-band callbacks that have been registered on |
| + /// [this]. |
| + int _totalCallbacks = 0; |
| + |
| + // TODO(nweiz): make this a read-only view when issue 8321 is fixed. |
| + /// The descriptions of all callbacks that are blocking the completion of |
| + /// [this]. |
| + Collection<String> get pendingCallbacks => _pendingCallbacks; |
| + final _pendingCallbacks = new Queue<String>(); |
| + |
| + /// A completer that will be completed once [_pendingCallbacks] becomes empty |
| + /// after the queue finished running its tasks. |
|
Bob Nystrom
2013/03/12 20:08:51
finished -> finishes.
nweiz
2013/03/12 20:56:51
Done.
|
| + Future get _noPendingCallbacks => _noPendingCallbacksCompleter.future; |
| + final Completer _noPendingCallbacksCompleter = new Completer(); |
| + |
| /// A [Future] that completes when the tasks in [this] are all complete. If an |
| /// error occurs while running this queue, the returned [Future] will complete |
| /// with that error. |
| @@ -407,6 +371,10 @@ class TaskQueue { |
| bool get isRunning => _schedule.state == ScheduleState.RUNNING && |
| _schedule.currentQueue == this; |
| + /// Whether this queue is running its tasks (as opposed to waiting for |
| + /// out-of-band callbacks or not running at all). |
| + bool get isRunningTasks => isRunning && _schedule.currentTask != null; |
| + |
| /// Schedules a task, [fn], to run asynchronously as part of this queue. Tasks |
| /// will be run in the order they're scheduled. In [fn] returns a [Future], |
| /// tasks after it won't be run until that [Future] completes. |
| @@ -459,14 +427,16 @@ class TaskQueue { |
| _signalError(error); |
| throw _error; |
| }); |
| + }).whenComplete(() { |
| + _schedule._currentTask = null; |
| }).then((_) { |
| _onTasksCompleteCompleter.complete(); |
| }).catchError((e) { |
| _onTasksCompleteCompleter.completeError(e); |
| throw e; |
| }).whenComplete(() { |
| - _schedule._currentTask = null; |
| - return _schedule._awaitNoPendingCallbacks().catchError((e) { |
| + if (pendingCallbacks.isEmpty) return; |
| + return _noPendingCallbacks.catchError((e) { |
| // Signal the error rather than passing it through directly so that if a |
| // timeout happens after an in-task error, both are reported. |
| _signalError(new ScheduleError.from(_schedule, e)); |
| @@ -480,6 +450,45 @@ class TaskQueue { |
| }); |
| } |
| + /// Returns a function wrapping [fn] that pipes any errors into the schedule |
| + /// chain. This will also block [this] from completing until the returned |
| + /// function has been called. It's used to ensure that out-of-band callbacks |
| + /// are properly handled by the scheduled test. |
| + Function _wrapAsync(fn(arg), String description) { |
| + assert(isRunning); |
| + |
| + // It's possible that the queue timed out before [fn] finished. |
| + bool _timedOut() => |
| + _schedule.currentQueue != this || pendingCallbacks.isEmpty; |
| + |
| + if (description == null) { |
| + description = "Out-of-band operation #${_totalCallbacks}"; |
| + } |
| + _totalCallbacks++; |
| + |
| + _pendingCallbacks.add(description);; |
|
Bob Nystrom
2013/03/12 20:08:51
;;
nweiz
2013/03/12 20:56:51
Done.
|
| + return (arg) { |
| + try { |
| + return fn(arg); |
| + } catch (e, stackTrace) { |
| + var error = new ScheduleError.from( |
| + _schedule, e, stackTrace: stackTrace); |
| + if (_timedOut()) { |
| + _schedule._signalPostTimeoutError(error); |
| + } else { |
| + _schedule.signalError(error); |
| + } |
| + } finally { |
| + if (_timedOut()) return; |
| + |
| + _pendingCallbacks.remove(description); |
| + if (_pendingCallbacks.isEmpty && !isRunningTasks) { |
| + _noPendingCallbacksCompleter.complete(); |
| + } |
| + } |
| + }; |
| + } |
| + |
| /// Signals that an out-of-band error has been detected and the queue should |
| /// stop running as soon as possible. |
| void _signalError(ScheduleError error) { |
| @@ -492,7 +501,10 @@ class TaskQueue { |
| /// Notifies the queue that it has timed out and it needs to terminate |
| /// immediately with a timeout error. |
| void _signalTimeout(ScheduleError error) { |
| - if (_taskFuture != null) { |
| + _pendingCallbacks.clear(); |
| + if (!isRunningTasks) { |
| + _noPendingCallbacksCompleter.completeError(error); |
| + } else if (_taskFuture != null) { |
| // Catch errors coming off the old task future, in case it completes after |
| // timing out. |
| _taskFuture.substitute(new Future.immediateError(error)).catchError((e) { |
| @@ -501,7 +513,7 @@ class TaskQueue { |
| } else { |
| // This branch probably won't be reached, but it's conceivable that the |
| // event loop might get pumped when _taskFuture is null but we haven't yet |
| - // called _awaitNoPendingCallbacks. |
| + // finished running all the tasks. |
| _signalError(error); |
| } |
| } |