Fix asyncWhere for broadcast streams (dart-lang/stream_transform#33)
Previous implementation may have failed to send onDone for some
listeners since each listener would increment `valuesWaiting` and only
the last listener would decrement it back to zero. There were also
potential bugs if the predicate was non-deterministic.
diff --git a/pkgs/stream_transform/CHANGELOG.md b/pkgs/stream_transform/CHANGELOG.md
index 25d39a5..1abe605 100644
--- a/pkgs/stream_transform/CHANGELOG.md
+++ b/pkgs/stream_transform/CHANGELOG.md
@@ -1,7 +1,7 @@
## 0.0.6
- Bug Fix: Some transformers did not correctly add data to all listeners on
- broadcast streams. Fixed for `throttle`, `debounce`, and `audit`.
+ broadcast streams. Fixed for `throttle`, `debounce`, `asyncWhere` and `audit`.
- Bug Fix: Only call the `tap` data callback once per event rather than once per
listener.
- Bug Fix: Allow canceling and re-listening to broadcast streams after a
diff --git a/pkgs/stream_transform/lib/src/async_where.dart b/pkgs/stream_transform/lib/src/async_where.dart
index 1bb2d64..6a4fc32 100644
--- a/pkgs/stream_transform/lib/src/async_where.dart
+++ b/pkgs/stream_transform/lib/src/async_where.dart
@@ -3,11 +3,13 @@
// BSD-style license that can be found in the LICENSE file.
import 'dart:async';
+import 'from_handlers.dart';
+
/// Like [Stream.where] but allows the [test] to return a [Future].
StreamTransformer<T, T> asyncWhere<T>(FutureOr<bool> test(T element)) {
var valuesWaiting = 0;
var sourceDone = false;
- return new StreamTransformer<T, T>.fromHandlers(handleData: (element, sink) {
+ return fromHandlers(handleData: (element, sink) {
valuesWaiting++;
() async {
if (await test(element)) sink.add(element);
diff --git a/pkgs/stream_transform/test/async_where_test.dart b/pkgs/stream_transform/test/async_where_test.dart
index 919f931..18c34bc 100644
--- a/pkgs/stream_transform/test/async_where_test.dart
+++ b/pkgs/stream_transform/test/async_where_test.dart
@@ -34,4 +34,36 @@
var filtered = values.transform(asyncWhere((e) => e > 4));
expect(await filtered.isEmpty, true);
});
+
+ test('forwards values to multiple listeners', () async {
+ var values = new StreamController.broadcast();
+ var filtered = values.stream.transform(asyncWhere((e) async => e > 2));
+ var firstValues = [];
+ var secondValues = [];
+ filtered..listen(firstValues.add)..listen(secondValues.add);
+ values..add(1)..add(2)..add(3)..add(4);
+ await new Future(() {});
+ expect(firstValues, [3, 4]);
+ expect(secondValues, [3, 4]);
+ });
+
+ test('closes streams with multiple listeners', () async {
+ var values = new StreamController.broadcast();
+ var predicate = new Completer<bool>();
+ var filtered = values.stream.transform(asyncWhere((_) => predicate.future));
+ var firstDone = false;
+ var secondDone = false;
+ filtered
+ ..listen(null, onDone: () => firstDone = true)
+ ..listen(null, onDone: () => secondDone = true);
+ values.add(1);
+ await values.close();
+ expect(firstDone, false);
+ expect(secondDone, false);
+
+ predicate.complete(true);
+ await new Future(() {});
+ expect(firstDone, true);
+ expect(secondDone, true);
+ });
}