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]);
+          });
+        }
+      });
+    }
+  }
 }