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

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

Issue 12321010: Increase size for slow-consumer test (overflows memory on X64). (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Done Created 7 years, 10 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/slow_consumer_test.dart » ('j') | no next file with comments »
Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
Index: sdk/lib/async/stream_impl.dart
diff --git a/sdk/lib/async/stream_impl.dart b/sdk/lib/async/stream_impl.dart
index be661e7bd5eded48d8df6c4fe43c9b88ee44fdf4..bbfc98ebb90661fd96348873c1ae0f6a34ced1ae 100644
--- a/sdk/lib/async/stream_impl.dart
+++ b/sdk/lib/async/stream_impl.dart
@@ -25,7 +25,7 @@ const int _STREAM_FIRING = 8;
const int _STREAM_CALLBACK = 16;
/// The count of times a stream has paused is stored in the
/// state, shifted by this amount.
-const int _STREAM_PAUSE_COUNT_SHIFT = 8;
+const int _STREAM_PAUSE_COUNT_SHIFT = 5;
// States for listeners.
@@ -148,6 +148,9 @@ abstract class _StreamImpl<T> extends Stream<T> {
/** Whether one or more active subscribers have requested a pause. */
bool get _isPaused => _state >= (1 << _STREAM_PAUSE_COUNT_SHIFT);
+ /** Whether we are currently executing a state-chance callback. */
+ bool get _isInCallback => (_state & _STREAM_CALLBACK) != 0;
+
/** Check whether the pending event queue is non-empty */
bool get _hasPendingEvent =>
_pendingEvents != null && !_pendingEvents.isEmpty;
@@ -231,7 +234,6 @@ abstract class _StreamImpl<T> extends Stream<T> {
void _endFiring() {
assert(_isFiring);
_state ^= _STREAM_FIRING;
-
if (!_hasSubscribers) {
_callOnSubscriptionStateChange();
} else if (_isPaused) {
@@ -332,16 +334,24 @@ abstract class _StreamImpl<T> extends Stream<T> {
void _callOnPauseStateChange() {
// After calling [_close], all pauses are handled internally by the Stream.
if (_isClosed) return;
- _state |= _STREAM_CALLBACK;
- _onPauseStateChange();
- _state ^= _STREAM_CALLBACK;
+ if (!_isInCallback) {
+ _state |= _STREAM_CALLBACK;
+ _onPauseStateChange();
+ _state ^= _STREAM_CALLBACK;
+ } else {
+ _onPauseStateChange();
+ }
}
/** Calls [_onSubscriptionStateChange] while setting callback bit. */
void _callOnSubscriptionStateChange() {
- _state |= _STREAM_CALLBACK;
- _onSubscriptionStateChange();
- _state ^= _STREAM_CALLBACK;
+ if (!_isInCallback) {
+ _state |= _STREAM_CALLBACK;
+ _onSubscriptionStateChange();
+ _state ^= _STREAM_CALLBACK;
+ } else {
+ _onSubscriptionStateChange();
+ }
}
/**
@@ -471,9 +481,12 @@ class _SingleStreamImpl<T> extends _StreamImpl<T> {
bool get _isPaused => (!_hasSubscribers && !_isComplete) || super._isPaused;
- // A single-stream is considered paused when it has no subscriber.
- bool get _canFireEvent =>
- _mayFireState && !_hasPendingEvent && _hasSubscribers;
+ bool get _canFireEvent {
+ // A single-stream is considered paused when it has no subscriber and
+ // isn't complete, so it can't fire if there is no subscriber.
+ // It won't try to fire when completed, so no need to check that.
+ return _mayFireState && !_hasPendingEvent && _hasSubscribers;
+ }
/** Whether there is currently a subscriber on this [Stream]. */
« no previous file with comments | « no previous file | tests/lib/async/slow_consumer_test.dart » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698