Allow null values in asyncMapSample and buffer (dart-lang/stream_transform#124)
Closes dart-lang/stream_transform#121
- Add a boolean to track whether we have a value since we can't use
`null` as a sentinel value for a potentially nullable generic type.
- Replace `!` with `as S` to allow for nullable `S`.
diff --git a/pkgs/stream_transform/lib/src/aggregate_sample.dart b/pkgs/stream_transform/lib/src/aggregate_sample.dart
index cd0a7d0..0eea768 100644
--- a/pkgs/stream_transform/lib/src/aggregate_sample.dart
+++ b/pkgs/stream_transform/lib/src/aggregate_sample.dart
@@ -21,6 +21,7 @@
: StreamController<S>(sync: true);
S? currentResults;
+ var hasCurrentResults = false;
var waitingForTrigger = true;
var isTriggerDone = false;
var isValueDone = false;
@@ -28,13 +29,15 @@
StreamSubscription<void>? triggerSub;
void emit() {
- controller.add(currentResults!);
+ controller.add(currentResults as S);
currentResults = null;
+ hasCurrentResults = false;
waitingForTrigger = true;
}
void onValue(T value) {
currentResults = aggregate(value, currentResults);
+ hasCurrentResults = true;
if (!waitingForTrigger) emit();
@@ -46,7 +49,7 @@
void onValuesDone() {
isValueDone = true;
- if (currentResults == null) {
+ if (!hasCurrentResults) {
triggerSub?.cancel();
controller.close();
}
@@ -55,7 +58,7 @@
void onTrigger(_) {
waitingForTrigger = false;
- if (currentResults != null) emit();
+ if (hasCurrentResults) emit();
if (isValueDone) {
triggerSub!.cancel();
diff --git a/pkgs/stream_transform/lib/src/rate_limit.dart b/pkgs/stream_transform/lib/src/rate_limit.dart
index 6a02703..ebb55a2 100644
--- a/pkgs/stream_transform/lib/src/rate_limit.dart
+++ b/pkgs/stream_transform/lib/src/rate_limit.dart
@@ -240,27 +240,33 @@
{required bool leading, required bool trailing}) {
Timer? timer;
S? soFar;
+ var hasPending = false;
var shouldClose = false;
var emittedLatestAsLeading = false;
+
return transformByHandlers(onData: (value, sink) {
+ void emit() {
+ sink.add(soFar as S);
+ soFar = null;
+ hasPending = false;
+ }
+
timer?.cancel();
soFar = collect(value, soFar);
+ hasPending = true;
if (timer == null && leading) {
emittedLatestAsLeading = true;
- sink.add(soFar as S);
+ emit();
} else {
emittedLatestAsLeading = false;
}
timer = Timer(duration, () {
- if (trailing && !emittedLatestAsLeading) sink.add(soFar as S);
- if (shouldClose) {
- sink.close();
- }
- soFar = null;
+ if (trailing && !emittedLatestAsLeading) emit();
+ if (shouldClose) sink.close();
timer = null;
});
}, onDone: (EventSink<S> sink) {
- if (soFar != null && trailing) {
+ if (hasPending && trailing) {
shouldClose = true;
} else {
timer?.cancel();
diff --git a/pkgs/stream_transform/test/async_map_sample_test.dart b/pkgs/stream_transform/test/async_map_sample_test.dart
index 06457d8..0d2d846 100644
--- a/pkgs/stream_transform/test/async_map_sample_test.dart
+++ b/pkgs/stream_transform/test/async_map_sample_test.dart
@@ -200,4 +200,9 @@
}
});
}
+
+ test('allows nulls', () async {
+ var stream = Stream<int?>.value(null);
+ await stream.asyncMapSample(expectAsync1((_) async {})).drain();
+ });
}
diff --git a/pkgs/stream_transform/test/debounce_test.dart b/pkgs/stream_transform/test/debounce_test.dart
index f24bfb9..9031db5 100644
--- a/pkgs/stream_transform/test/debounce_test.dart
+++ b/pkgs/stream_transform/test/debounce_test.dart
@@ -231,4 +231,9 @@
});
});
}
+ test('allows nulls', () async {
+ final values = Stream<int?>.fromIterable([null]);
+ final transformed = values.debounce(const Duration(milliseconds: 1));
+ expect(await transformed.toList(), [null]);
+ });
}