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