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

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

Issue 137793004: Fix scheduled stream usage of asBroadcastStream (Closed) Base URL: http://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/pkg.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
===================================================================
--- pkg/scheduled_test/lib/scheduled_stream.dart (revision 31773)
+++ pkg/scheduled_test/lib/scheduled_stream.dart (working copy)
@@ -40,6 +40,7 @@
/// 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 @@
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 @@
_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 @@
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 @@
/// 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/pkg.status ('k') | no next file » | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698