Rewrite concat (dart-lang/stream_transform#32) Uses sync forwarders for dart-lang/stream_transform#24 The simple implementation leaves out a lot of edge cases around single-subscription vs broadcast streams and how they combine. Rewrite to handle pause, resume, and cancel manually. - Loop over stream types in tests - Add test for onDone and canceling and relistening - Force _next to broadcast when first is broadcast - Document behavior of broadcast stream as second stream - Add tests for pause/resume
diff --git a/pkgs/stream_transform/CHANGELOG.md b/pkgs/stream_transform/CHANGELOG.md index d6bd7a9..da8fd25 100644 --- a/pkgs/stream_transform/CHANGELOG.md +++ b/pkgs/stream_transform/CHANGELOG.md
@@ -6,6 +6,8 @@ listener. - Bug Fix: Allow canceling and re-listening to broadcast streams after a `merge` transform. +- Bug Fix: Single-subscription streams concatted after broadcast streams are + handled correctly. - Use sync `StreamControllers` for forwarding where possible. ## 0.0.5
diff --git a/pkgs/stream_transform/lib/src/concat.dart b/pkgs/stream_transform/lib/src/concat.dart index 0a2c27a..955918a 100644 --- a/pkgs/stream_transform/lib/src/concat.dart +++ b/pkgs/stream_transform/lib/src/concat.dart
@@ -8,6 +8,15 @@ /// /// If the initial stream never finishes, the [next] stream will never be /// listened to. +/// +/// If a single-subscription is concatted to the end of a broadcast stream it +/// may be listened to and never canceled. +/// +/// If a broadcast stream is concatted to any other stream it will miss any +/// events which occur before the first stream is done. If a broadcast stream is +/// concatted to a single-subscription stream, pausing the stream while it is +/// listening to the second stream will cause events to be dropped rather than +/// buffered. StreamTransformer<T, T> concat<T>(Stream<T> next) => new _Concat<T>(next); class _Concat<T> implements StreamTransformer<T, T> { @@ -17,11 +26,61 @@ @override Stream<T> bind(Stream<T> first) { - var controller = new StreamController<T>(); - controller - .addStream(first) - .then((_) => controller.addStream(_next)) - .then((_) => controller.close()); + var controller = first.isBroadcast + ? new StreamController<T>.broadcast(sync: true) + : new StreamController<T>(sync: true); + + var next = first.isBroadcast && !_next.isBroadcast + ? _next.asBroadcastStream() + : _next; + + StreamSubscription subscription; + var currentStream = first; + var firstDone = false; + var secondDone = false; + + Function currentDoneHandler; + + listen() { + subscription = currentStream.listen(controller.add, + onError: controller.addError, onDone: () => currentDoneHandler()); + } + + onSecondDone() { + secondDone = true; + controller.close(); + } + + onFirstDone() { + firstDone = true; + currentStream = next; + currentDoneHandler = onSecondDone; + listen(); + } + + currentDoneHandler = onFirstDone; + + controller.onListen = () { + if (subscription != null) return; + listen(); + if (!first.isBroadcast) { + controller.onPause = () { + if (!firstDone || !next.isBroadcast) return subscription.pause(); + subscription.cancel(); + subscription = null; + }; + controller.onResume = () { + if (!firstDone || !next.isBroadcast) return subscription.resume(); + listen(); + }; + } + controller.onCancel = () { + if (secondDone) return null; + var toCancel = subscription; + subscription = null; + return toCancel.cancel(); + }; + }; return controller.stream; } }
diff --git a/pkgs/stream_transform/test/concat_test.dart b/pkgs/stream_transform/test/concat_test.dart index 6cf041b..37a46f7 100644 --- a/pkgs/stream_transform/test/concat_test.dart +++ b/pkgs/stream_transform/test/concat_test.dart
@@ -8,40 +8,151 @@ import 'package:stream_transform/stream_transform.dart'; void main() { - group('concat', () { - test('adds all values from both streams', () async { - var first = new Stream.fromIterable([1, 2, 3]); - var second = new Stream.fromIterable([4, 5, 6]); - var all = await first.transform(concat(second)).toList(); - expect(all, [1, 2, 3, 4, 5, 6]); - }); + var streamTypes = { + 'single subscription': () => new StreamController(), + 'broadcast': () => new StreamController.broadcast() + }; + for (var firstType in streamTypes.keys) { + for (var secondType in streamTypes.keys) { + group('concat [$firstType] with [$secondType]', () { + StreamController first; + StreamController second; - test('closes first stream on cancel', () async { - var firstStreamClosed = false; - var first = new StreamController() - ..onCancel = () { - firstStreamClosed = true; - }; - var second = new StreamController(); - var subscription = - first.stream.transform(concat(second.stream)).listen((_) {}); - await subscription.cancel(); - expect(firstStreamClosed, true); - }); + List emittedValues; + bool firstCanceled; + bool secondCanceled; + bool secondListened; + bool isDone; + List errors; + Stream transformed; + StreamSubscription subscription; - test('closes second stream on cancel if first stream done', () async { - var first = new StreamController(); - var secondStreamClosed = false; - var second = new StreamController() - ..onCancel = () { - secondStreamClosed = true; - }; - var subscription = - first.stream.transform(concat(second.stream)).listen((_) {}); - await first.close(); - await new Future(() {}); - await subscription.cancel(); - expect(secondStreamClosed, true); - }); - }); + setUp(() async { + firstCanceled = false; + secondCanceled = false; + secondListened = false; + first = streamTypes[firstType]() + ..onCancel = () { + firstCanceled = true; + }; + second = streamTypes[secondType]() + ..onCancel = () { + secondCanceled = true; + } + ..onListen = () { + secondListened = true; + }; + emittedValues = []; + errors = []; + isDone = false; + transformed = first.stream.transform(concat(second.stream)); + subscription = transformed + .listen(emittedValues.add, onError: errors.add, onDone: () { + isDone = true; + }); + }); + + test('adds all values from both streams', () async { + first..add(1)..add(2); + await first.close(); + await new Future(() {}); + second..add(3)..add(4); + await new Future(() {}); + expect(emittedValues, [1, 2, 3, 4]); + }); + + test('Does not listen to second stream before first stream finishes', + () async { + expect(secondListened, false); + await first.close(); + expect(secondListened, true); + }); + + test('closes stream after both inputs close', () async { + await first.close(); + await second.close(); + expect(isDone, true); + }); + + test('cancels any type of first stream on cancel', () async { + await subscription.cancel(); + expect(firstCanceled, true); + }); + + if (firstType == 'single subscription') { + test( + 'cancels any type of second stream on cancel if first is ' + 'broadcast', () async { + await first.close(); + await subscription.cancel(); + expect(secondCanceled, true); + }); + + if (secondType == 'broadcast') { + test('can pause and resume during second stream - dropping values', + () async { + await first.close(); + subscription.pause(); + second.add(1); + await new Future(() {}); + subscription.resume(); + second.add(2); + await new Future(() {}); + expect(emittedValues, [2]); + }); + } else { + test('can pause and resume during second stream - buffering values', + () async { + await first.close(); + subscription.pause(); + second.add(1); + await new Future(() {}); + subscription.resume(); + second.add(2); + await new Future(() {}); + expect(emittedValues, [1, 2]); + }); + } + } + + if (firstType == 'broadcast') { + test('can cancel and relisten during first stream', () async { + await subscription.cancel(); + first.add(1); + subscription = transformed.listen(emittedValues.add); + first.add(2); + await new Future(() {}); + expect(emittedValues, [2]); + }); + + test('can cancel and relisten during second stream', () async { + await first.close(); + await subscription.cancel(); + second.add(2); + await new Future(() {}); + subscription = transformed.listen(emittedValues.add); + second.add(3); + await new Future((() {})); + expect(emittedValues, [3]); + }); + + test('forwards values to multiple listeners', () async { + var otherValues = []; + transformed.listen(otherValues.add); + first.add(1); + await first.close(); + second.add(2); + await new Future(() {}); + var thirdValues = []; + transformed.listen(thirdValues.add); + second.add(3); + await new Future(() {}); + expect(emittedValues, [1, 2, 3]); + expect(otherValues, [1, 2, 3]); + expect(thirdValues, [3]); + }); + } + }); + } + } }