Add public library and buffer utility (dart-lang/stream_transform#2)

- Add `buffer` utility and tests
- Export buffer from package library
- Add travis config
- Add README, CHANGELOG, and version
diff --git a/pkgs/stream_transform/.travis.yml b/pkgs/stream_transform/.travis.yml
new file mode 100644
index 0000000..ed6c5bb
--- /dev/null
+++ b/pkgs/stream_transform/.travis.yml
@@ -0,0 +1,13 @@
+language: dart
+sudo: false
+dart:
+  - dev
+  - stable
+  - 1.22.1
+cache:
+  directories:
+    - $HOME/.pub-cache
+dart_task:
+  - test
+  - dartfmt
+  - dartanalyzer
diff --git a/pkgs/stream_transform/CHANGELOG.md b/pkgs/stream_transform/CHANGELOG.md
new file mode 100644
index 0000000..af394b6
--- /dev/null
+++ b/pkgs/stream_transform/CHANGELOG.md
@@ -0,0 +1,4 @@
+## 0.0.1
+
+- Add `buffer` utility. Collects events in a `List` until a `trigger` stream
+  fires.
diff --git a/pkgs/stream_transform/README.md b/pkgs/stream_transform/README.md
new file mode 100644
index 0000000..ecab8a4
--- /dev/null
+++ b/pkgs/stream_transform/README.md
@@ -0,0 +1,7 @@
+Contains utility methods to create `StreamTransfomer` instances to manipulate
+Streams.
+
+# buffer
+
+Collects values from a source stream until a `trigger` stream fires and the
+collected values are emitted.
diff --git a/pkgs/stream_transform/lib/src/buffer.dart b/pkgs/stream_transform/lib/src/buffer.dart
new file mode 100644
index 0000000..4f2976d
--- /dev/null
+++ b/pkgs/stream_transform/lib/src/buffer.dart
@@ -0,0 +1,144 @@
+// Copyright (c) 2017, the Dart project authors.  Please see the AUTHORS file
+// for details. All rights reserved. Use of this source code is governed by a
+// BSD-style license that can be found in the LICENSE file.
+import 'dart:async';
+
+/// Creates a [StreamTransformer] which collects values and emits when it sees a
+/// value on [trigger].
+///
+/// If there are no pending values when [trigger] emits, the next value on the
+/// source Stream will immediately flow through. Otherwise, the pending values
+/// are released when [trigger] emits.
+///
+/// Errors from the source stream or the trigger are immediately forwarded to
+/// the output.
+StreamTransformer<T, List<T>> buffer<T>(Stream trigger) => new _Buffer(trigger);
+
+List<T> _collectToList<T>(T element, List<T> soFar) {
+  soFar ??= <T>[];
+  soFar.add(element);
+  return soFar;
+}
+
+/// A StreamTransformer which aggregates values and emits when it sees a value
+/// on [_trigger].
+///
+/// If there are no pending values when [_trigger] emits the first value on the
+/// source Stream will immediately flow through. Otherwise, the pending values
+/// and released when [_trigger] emits.
+///
+/// Errors from the source stream or the trigger are immediately forwarded to
+/// the output.
+class _Buffer<T> implements StreamTransformer<T, List<T>> {
+  final Stream _trigger;
+
+  _Buffer(this._trigger);
+
+  @override
+  Stream<List<T>> bind(Stream<T> values) {
+    StreamController<List<T>> controller;
+    if (values.isBroadcast) {
+      controller = new StreamController<List<T>>.broadcast();
+    } else {
+      controller = new StreamController<List<T>>();
+    }
+
+    List<T> currentResults;
+    bool waitingForTrigger = true;
+    StreamSubscription valuesSub;
+    StreamSubscription triggerSub;
+
+    cancelValues() {
+      var sub = valuesSub;
+      valuesSub = null;
+      return sub?.cancel() ?? new Future.value();
+    }
+
+    cancelTrigger() {
+      var sub = triggerSub;
+      triggerSub = null;
+      return sub?.cancel() ?? new Future.value();
+    }
+
+    closeController() {
+      var ctl = controller;
+      controller = null;
+      return ctl?.close() ?? new Future.value();
+    }
+
+    emit() {
+      controller.add(currentResults);
+      currentResults = null;
+      waitingForTrigger = true;
+    }
+
+    onValue(T value) {
+      currentResults = _collectToList(value, currentResults);
+      if (!waitingForTrigger) {
+        emit();
+      }
+    }
+
+    valuesDone() {
+      valuesSub = null;
+      if (currentResults == null) {
+        closeController();
+        cancelTrigger();
+      }
+    }
+
+    onTrigger(_) {
+      if (currentResults == null) {
+        waitingForTrigger = false;
+        return;
+      }
+      emit();
+      if (valuesSub == null) {
+        closeController();
+        cancelTrigger();
+      }
+    }
+
+    triggerDone() {
+      cancelValues();
+      closeController();
+    }
+
+    controller.onListen = () {
+      if (valuesSub != null) return;
+      valuesSub = values.listen(onValue,
+          onError: controller.addError, onDone: valuesDone);
+      if (triggerSub != null) {
+        if (triggerSub.isPaused) triggerSub.resume();
+      } else {
+        triggerSub = _trigger.listen(onTrigger,
+            onError: controller.addError, onDone: triggerDone);
+      }
+    };
+
+    // Forward methods from listener
+    if (!values.isBroadcast) {
+      controller.onPause = () {
+        valuesSub?.pause();
+        triggerSub?.pause();
+      };
+      controller.onResume = () {
+        valuesSub?.resume();
+        triggerSub?.resume();
+      };
+      controller.onCancel =
+          () => Future.wait([cancelValues(), cancelTrigger()]);
+    } else {
+      controller.onCancel = () {
+        if (controller?.hasListener ?? false) return;
+        if (_trigger.isBroadcast) {
+          cancelTrigger();
+        } else {
+          triggerSub.pause();
+        }
+        cancelValues();
+      };
+    }
+    return controller.stream;
+  }
+}
diff --git a/pkgs/stream_transform/lib/stream_transform.dart b/pkgs/stream_transform/lib/stream_transform.dart
new file mode 100644
index 0000000..02069f5
--- /dev/null
+++ b/pkgs/stream_transform/lib/stream_transform.dart
@@ -0,0 +1,5 @@
+// Copyright (c) 2017, the Dart project authors.  Please see the AUTHORS file
+// for details. All rights reserved. Use of this source code is governed by a
+// BSD-style license that can be found in the LICENSE file.
+
+export 'src/buffer.dart';
diff --git a/pkgs/stream_transform/pubspec.yaml b/pkgs/stream_transform/pubspec.yaml
index 129aa9c..b260f12 100644
--- a/pkgs/stream_transform/pubspec.yaml
+++ b/pkgs/stream_transform/pubspec.yaml
@@ -2,6 +2,7 @@
 description: A collection of utilities to transform and manipulate streams.
 author: Dart Team <misc@dartlang.org>
 homepage: https://www.github.com/dart-lang/stream_transform
+version: 0.0.1-dev
 
 environment:
   sdk: ">=1.22.0 <2.0.0"
diff --git a/pkgs/stream_transform/test/buffer_test.dart b/pkgs/stream_transform/test/buffer_test.dart
new file mode 100644
index 0000000..c371ac5
--- /dev/null
+++ b/pkgs/stream_transform/test/buffer_test.dart
@@ -0,0 +1,185 @@
+// Copyright (c) 2017, the Dart project authors.  Please see the AUTHORS file
+// for details. All rights reserved. Use of this source code is governed by a
+// BSD-style license that can be found in the LICENSE file.
+import 'dart:async';
+
+import 'package:test/test.dart';
+
+import 'package:stream_transform/stream_transform.dart';
+
+void main() {
+  var streamTypes = {
+    'single subscription': () => new StreamController(),
+    'broadcast': () => new StreamController.broadcast()
+  };
+  for (var triggerType in streamTypes.keys) {
+    for (var valuesType in streamTypes.keys) {
+      group('Trigger type: [$triggerType], Values type: [$valuesType]', () {
+        StreamController trigger;
+        StreamController values;
+        List emittedValues;
+        bool valuesCanceled;
+        bool triggerCanceled;
+        bool isDone;
+        List errors;
+        StreamSubscription subscription;
+
+        setUp(() async {
+          valuesCanceled = false;
+          triggerCanceled = false;
+          trigger = streamTypes[triggerType]()
+            ..onCancel = () {
+              triggerCanceled = true;
+            };
+          values = streamTypes[triggerType]()
+            ..onCancel = () {
+              valuesCanceled = true;
+            };
+          emittedValues = [];
+          errors = [];
+          isDone = false;
+          subscription = values.stream
+              .transform(buffer(trigger.stream))
+              .listen(emittedValues.add, onError: errors.add, onDone: () {
+            isDone = true;
+          });
+        });
+
+        test('does not emit before `trigger`', () async {
+          values.add(1);
+          await new Future(() {});
+          expect(emittedValues, isEmpty);
+          trigger.add(null);
+          await new Future(() {});
+          expect(emittedValues, [
+            [1]
+          ]);
+        });
+
+        test('emits immediately if trigger emits before a value', () async {
+          trigger.add(null);
+          await new Future(() {});
+          expect(emittedValues, isEmpty);
+          values.add(1);
+          await new Future(() {});
+          expect(emittedValues, [
+            [1]
+          ]);
+        });
+
+        test('two triggers in a row - emit then emit next value', () async {
+          values.add(1);
+          values.add(2);
+          await new Future(() {});
+          trigger.add(null);
+          trigger.add(null);
+          await new Future(() {});
+          values.add(3);
+          await new Future(() {});
+          expect(emittedValues, [
+            [1, 2],
+            [3]
+          ]);
+        });
+
+        test('pre-emptive trigger then trigger after values', () async {
+          trigger.add(null);
+          await new Future(() {});
+          values.add(1);
+          values.add(2);
+          await new Future(() {});
+          trigger.add(null);
+          await new Future(() {});
+          expect(emittedValues, [
+            [1],
+            [2]
+          ]);
+        });
+
+        test('multiple pre-emptive triggers, only emits first value', () async {
+          trigger.add(null);
+          trigger.add(null);
+          await new Future(() {});
+          values.add(1);
+          values.add(2);
+          await new Future(() {});
+          expect(emittedValues, [
+            [1]
+          ]);
+        });
+
+        test('groups values between trigger', () async {
+          values.add(1);
+          values.add(2);
+          await new Future(() {});
+          trigger.add(null);
+          values.add(3);
+          values.add(4);
+          await new Future(() {});
+          trigger.add(null);
+          await new Future(() {});
+          expect(emittedValues, [
+            [1, 2],
+            [3, 4]
+          ]);
+        });
+
+        test('cancels value subscription when output canceled', () async {
+          expect(valuesCanceled, false);
+          await subscription.cancel();
+          expect(valuesCanceled, true);
+        });
+
+        test('cancels trigger subscription when output canceled', () async {
+          expect(triggerCanceled, false);
+          await subscription.cancel();
+          expect(triggerCanceled, true);
+        });
+
+        test('closes when trigger ends', () async {
+          expect(isDone, false);
+          await trigger.close();
+          await new Future(() {});
+          expect(isDone, true);
+        });
+
+        test('closes after outputting final values when source closes',
+            () async {
+          expect(isDone, false);
+          values.add(1);
+          await values.close();
+          expect(isDone, false);
+          trigger.add(null);
+          await new Future(() {});
+          expect(emittedValues, [
+            [1]
+          ]);
+          expect(isDone, true);
+        });
+
+        test(
+            'closes immediately if there are no pending values when source closes',
+            () async {
+          expect(isDone, false);
+          values.add(1);
+          trigger.add(null);
+          await values.close();
+          await new Future(() {});
+          expect(isDone, true);
+        });
+
+        test('forwards errors from trigger', () async {
+          trigger.addError('error');
+          await new Future(() {});
+          expect(errors, ['error']);
+        });
+
+        test('forwards errors from values', () async {
+          values.addError('error');
+          await new Future(() {});
+          expect(errors, ['error']);
+        });
+      });
+    }
+  }
+}