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

Unified Diff: sdk/lib/io/http_parser.dart

Issue 15256002: Rewrite parts of http-parser, to better handle connection errors and pausing. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Created 7 years, 7 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
Index: sdk/lib/io/http_parser.dart
diff --git a/sdk/lib/io/http_parser.dart b/sdk/lib/io/http_parser.dart
index 18dc2b315cd00106ff9772a0d367c46fe1278453..8a0b1feac6ae5a36fcd50950c9fd5e190d62ae0e 100644
--- a/sdk/lib/io/http_parser.dart
+++ b/sdk/lib/io/http_parser.dart
@@ -197,10 +197,23 @@ class _HttpParser
_HttpParser._(this._requestParser) {
_controller = new StreamController<_HttpIncoming>(
- onListen: _updateParsePauseState,
- onPause: _updateParsePauseState,
- onResume: _updateParsePauseState,
- onCancel: _updateParsePauseState);
+ onListen: () {
+ _socketSubscription.resume();
+ },
+ onPause: () {
+ _paused = true;
+ _pauseStateChanged();
+ },
+ onResume: () {
+ _paused = false;
+ _pauseStateChanged();
+ },
+ onCancel: () {
+ try {
+ _socketSubscription.cancel();
+ } catch (e) {
+ }
+ });
_reset();
}
@@ -228,6 +241,7 @@ class _HttpParser
onDone: () {
completer.complete(this);
});
+ _socketSubscription.pause();
return completer.future;
}
@@ -267,7 +281,7 @@ class _HttpParser
_index < _buffer.length &&
_state != _State.FAILURE &&
_state != _State.UPGRADED) {
- if (_paused) {
+ if (_bodyPaused) {
Søren Gjesse 2013/05/17 07:09:28 Don't we need to check for both _paused and _bodyP
Anders Johnsen 2013/05/17 08:27:59 Done.
_parserCalled = false;
return;
}
@@ -562,8 +576,9 @@ class _HttpParser
_uri_or_reason_phrase.clear();
if (_connectionUpgrade) {
_incoming.upgraded = true;
+ _parserCalled = false;
_controller.add(_incoming);
- break;
+ return;
}
if (_transferLength == 0 ||
(_messageType == _MessageType.RESPONSE &&
@@ -584,8 +599,9 @@ class _HttpParser
// indicated by close.
_state = _State.BODY;
}
+ _parserCalled = false;
_controller.add(_incoming);
- break;
+ return;
case _State.CHUNK_SIZE_STARTING_CR:
_expect(byte, _CharCode.CR);
@@ -864,58 +880,71 @@ class _HttpParser
void _createIncoming(int transferLength) {
assert(_incoming == null);
assert(_bodyController == null);
+ assert(!_bodyPaused);
+ var incoming;
_bodyController = new StreamController<List<int>>(
- onListen: _bodySubscriptionStateChange,
- onPause: _updateParsePauseState,
- onResume: _updateParsePauseState,
- onCancel: _bodySubscriptionStateChange);
- _incoming = new _HttpIncoming(
+ onListen: () {
+ if (incoming != _incoming) return;
+ assert(_bodyPaused);
+ _bodyPaused = false;
+ _pauseStateChanged();
+ },
+ onPause: () {
+ if (incoming != _incoming) return;
+ assert(!_bodyPaused);
+ _bodyPaused = true;
+ _pauseStateChanged();
+ },
+ onResume: () {
+ if (incoming != _incoming) return;
+ assert(_bodyPaused);
+ _bodyPaused = false;
+ _pauseStateChanged();
+ },
+ onCancel: () {
+ if (incoming != _incoming) return;
+ if (_socketSubscription != null) {
+ _socketSubscription.cancel();
+ }
+ _closeIncoming();
+ _controller.close();
+ });
+ incoming = _incoming = new _HttpIncoming(
_headers, transferLength, _bodyController.stream);
- _pauseParsing(); // Needed to handle detaching - don't start on the body!
+ _bodyPaused = true;
+ _pauseStateChanged();
}
void _closeIncoming() {
// Ignore multiple close (can happend in re-entrance).
if (_incoming == null) return;
var tmp = _incoming;
- _incoming = null;
tmp.close();
+ _incoming = null;
if (_bodyController != null) {
_bodyController.close();
_bodyController = null;
}
- _updateParsePauseState();
- }
-
- void _continueParsing() {
- _paused = false;
- if (!_parserCalled && _buffer != null) _parse();
- }
-
- void _pauseParsing() {
- _paused = true;
+ _bodyPaused = false;
+ _pauseStateChanged();
}
- void _bodySubscriptionStateChange() {
- if (_incoming != null && !_bodyController.hasListener) {
- _closeIncoming();
- } else {
- _updateParsePauseState();
+ void _pauseStateChanged() {
+ void update(bool pause) {
+ if (pause && !_socketSubscription.isPaused) {
+ _socketSubscription.pause();
+ } else if (!pause && _socketSubscription.isPaused) {
+ _socketSubscription.resume();
+ }
}
- }
- void _updateParsePauseState() {
- if (_bodyController != null) {
- if (_bodyController.hasListener && !_bodyController.isPaused) {
- _continueParsing();
- } else {
- _pauseParsing();
+ if (_incoming != null) {
+ if (!_bodyPaused && !_parserCalled) {
+ _parse();
}
} else {
- if (_controller.hasListener && !_controller.isPaused) {
- _continueParsing();
- } else {
- _pauseParsing();
+ if (!_paused && !_parserCalled) {
+ _parse();
}
}
}
@@ -959,7 +988,8 @@ class _HttpParser
// The current incoming connection.
_HttpIncoming _incoming;
StreamSubscription _socketSubscription;
- bool _paused = false;
+ bool _paused = true;
+ bool _bodyPaused = false;
Completer _pauseCompleter;
StreamController<_HttpIncoming> _controller;
StreamController<List<int>> _bodyController;

Powered by Google App Engine
This is Rietveld 408576698