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

Unified Diff: pkg/scheduled_test/lib/scheduled_stream.dart

Issue 139363002: Reapply r31820 with a 'slow' status file marker. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Created 6 years, 11 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 | « pkg/pkgbuild.status ('k') | no next file » | no next file with comments »
Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
Index: pkg/scheduled_test/lib/scheduled_stream.dart
diff --git a/pkg/scheduled_test/lib/scheduled_stream.dart b/pkg/scheduled_test/lib/scheduled_stream.dart
index b6c51847737470ca3018de9e40affba0fb650af2..2b8be676ee949fbc36419f8ec344d284d1e4caea 100644
--- a/pkg/scheduled_test/lib/scheduled_stream.dart
+++ b/pkg/scheduled_test/lib/scheduled_stream.dart
@@ -40,6 +40,7 @@ class ScheduledStream<T> {
/// The set of all streams forked from this one.
final _forks = new Set<ScheduledStream<T>>();
+ final _forkControllers = new Set<StreamController<T>>();
/// The queue of values emitted by [_stream] but not yet emitted through
/// [next].
@@ -72,9 +73,11 @@ class ScheduledStream<T> {
bool _isNextPending = false;
/// Creates a new scheduled stream wrapping [stream].
- ScheduledStream(Stream<T> stream)
- : _stream = stream.asBroadcastStream() {
+ ScheduledStream(this._stream) {
_subscription = _stream.listen((value) {
+ for (var c in _forkControllers) {
+ c.add(value);
+ }
if (_hasNextCompleter != null) {
_hasNextCompleter.complete(true);
_hasNextCompleter = null;
@@ -89,6 +92,9 @@ class ScheduledStream<T> {
_pendingValues.add(new Fallible.withValue(value));
}
}, onError: (error, stackTrace) {
+ for (var c in _forkControllers) {
+ c.addError(error, stackTrace);
+ }
if (_hasNextCompleter != null) {
_hasNextCompleter.completeError(error, stackTrace);
_hasNextCompleter = null;
@@ -205,11 +211,10 @@ class ScheduledStream<T> {
controller.addError(valueOrError.error, valueOrError.stackTrace);
}
}
-
if (_isDone) {
controller.close();
} else {
- _stream.pipe(controller);
+ _forkControllers.add(controller);
}
var fork = new ScheduledStream<T>(controller.stream);
@@ -237,6 +242,11 @@ class ScheduledStream<T> {
/// Handles a "done" event from the underlying stream, as well as [this] being
/// closed.
void _onDone() {
+ for (var c in _forkControllers) {
+ c.close();
+ }
+ _forkControllers.clear();
+
if (_hasNextCompleter != null) {
_hasNextCompleter.complete(false);
_hasNextCompleter = null;
« no previous file with comments | « pkg/pkgbuild.status ('k') | no next file » | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698