Chromium Code Reviews
chromiumcodereview-hr@appspot.gserviceaccount.com (chromiumcodereview-hr) | Please choose your nickname with Settings | Help | Chromium Project | Gerrit Changes | Sign out
(649)

Unified Diff: sdk/lib/async/stream.dart

Issue 12381064: Add Stream.periodic. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Address comments. Add test. Created 7 years, 9 months ago
Use n/p to move between diff chunks; N/P to move between comments. Draft comments are only viewable by you.
Jump to:
View side-by-side diff with in-line comments
Download patch
« no previous file with comments | « no previous file | tests/lib/async/stream_periodic2_test.dart » ('j') | no next file with comments »
Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
Index: sdk/lib/async/stream.dart
diff --git a/sdk/lib/async/stream.dart b/sdk/lib/async/stream.dart
index 8ac0541908681a94f19156baf1051b542fafbdaa..ff133cc3e1164c15b538dd93583b725002657b1b 100644
--- a/sdk/lib/async/stream.dart
+++ b/sdk/lib/async/stream.dart
@@ -83,6 +83,67 @@ abstract class Stream<T> {
}
/**
+ * Creates a stream that repeatedly emits events at [period] intervals.
+ *
+ * The event values are computed by invoking [computation]. The argument to
+ * this callback is an integer that starts with 0 and is incremented for
+ * every event.
+ *
+ * If [computation] is omitted the event values will all be `null`.
+ */
+ factory Stream.periodic(Duration period,
+ [T computation(int computationCount)]) {
+ if (computation == null) computation = ((i) => null);
+
+ Timer timer;
+ int computationCount = 0;
+ StreamController<T> controller;
+ // Counts the time that the Stream was running (and not paused).
+ Stopwatch watch = new Stopwatch();
+
+ void sendEvent() {
+ watch.reset();
+ T data = computation(computationCount++);
+ controller.add(data);
+ }
+
+ void startPeriodicTimer() {
+ assert(timer == null);
+ timer = new Timer.periodic(period, (Timer timer) {
+ sendEvent();
+ });
+ }
+
+ controller = new StreamController<T>(
+ onPauseStateChange: () {
+ if (controller.isPaused) {
+ timer.cancel();
+ timer = null;
+ watch.stop();
+ } else {
+ assert(timer == null);
+ Duration elapsed = watch.elapsed;
+ watch.start();
+ timer = new Timer(period - elapsed, () {
+ timer = null;
+ startPeriodicTimer();
+ sendEvent();
+ });
+ }
+ },
+ onSubscriptionStateChange: () {
+ if (controller.hasSubscribers) {
+ watch.start();
+ startPeriodicTimer();
+ } else {
+ if (timer != null) timer.cancel();
+ timer = null;
+ }
+ });
+ return controller.stream;
+ }
+
+ /**
* Reports whether this stream is a broadcast stream.
*/
bool get isBroadcast => false;
« no previous file with comments | « no previous file | tests/lib/async/stream_periodic2_test.dart » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698