| Index: sdk/lib/io/websocket_impl.dart
|
| diff --git a/sdk/lib/io/websocket_impl.dart b/sdk/lib/io/websocket_impl.dart
|
| index 45b1f24f771a38cafd396726d42fadba4e181689..c054b1ec9e0bc301a1004a05f04052951e6be5a3 100644
|
| --- a/sdk/lib/io/websocket_impl.dart
|
| +++ b/sdk/lib/io/websocket_impl.dart
|
| @@ -572,7 +572,6 @@ class _WebSocketConsumer implements StreamConsumer {
|
| StreamController _controller;
|
| StreamSubscription _subscription;
|
| bool _issuedPause = false;
|
| - bool _reportError = false;
|
| Completer _closeCompleter = new Completer();
|
| Completer _completer;
|
|
|
| @@ -614,13 +613,8 @@ class _WebSocketConsumer implements StreamConsumer {
|
| _closeCompleter.complete(webSocket);
|
| },
|
| onError: (error) {
|
| - if (_reportError) {
|
| - if (!_done(error)) {
|
| - _closeCompleter.completeError(error);
|
| - }
|
| - } else {
|
| - _done();
|
| - _closeCompleter.complete(webSocket);
|
| + if (!_done(error)) {
|
| + _closeCompleter.completeError(error);
|
| }
|
| });
|
| }
|
| @@ -641,7 +635,6 @@ class _WebSocketConsumer implements StreamConsumer {
|
| _completer = new Completer();
|
| _subscription = stream.listen(
|
| (data) {
|
| - _reportError = true;
|
| _controller.add(data);
|
| },
|
| onDone: () {
|
| @@ -669,14 +662,13 @@ class _WebSocketConsumer implements StreamConsumer {
|
|
|
| void add(data) {
|
| _ensureController();
|
| - _reportError = false;
|
| _controller.add(data);
|
| }
|
| }
|
|
|
|
|
| class _WebSocketImpl extends Stream implements WebSocket {
|
| - StreamController _controller;
|
| + final StreamController _controller = new StreamController(sync: true);
|
| StreamSink _sink;
|
|
|
| final Socket _socket;
|
| @@ -688,7 +680,6 @@ class _WebSocketImpl extends Stream implements WebSocket {
|
| Duration _pingInterval;
|
| Timer _pingTimer;
|
| _WebSocketConsumer _consumer;
|
| - var _subscription;
|
|
|
| int _outCloseCode;
|
| String _outCloseReason;
|
| @@ -772,14 +763,13 @@ class _WebSocketImpl extends Stream implements WebSocket {
|
| _readyState = WebSocket.OPEN;
|
|
|
| var transformer = new _WebSocketProtocolTransformer(_serverSide);
|
| - _subscription = _socket.transform(transformer).listen(
|
| + _socket.transform(transformer).listen(
|
| (data) {
|
| - // Simply set pingInterval, as it'll cancel any timers.
|
| - pingInterval = _pingInterval;
|
| if (data is _WebSocketPing) {
|
| if (!_writeClosed) _consumer.add(new _WebSocketPong(data.payload));
|
| } else if (data is _WebSocketPong) {
|
| - // Already handled.
|
| + // Simply set pingInterval, as it'll cancel any timers.
|
| + pingInterval = _pingInterval;
|
| } else {
|
| _controller.add(data);
|
| }
|
| @@ -808,11 +798,6 @@ class _WebSocketImpl extends Stream implements WebSocket {
|
| _controller.close();
|
| },
|
| cancelOnError: true);
|
| - _subscription.pause();
|
| - _controller = new StreamController(sync: true,
|
| - onListen: _subscription.resume,
|
| - onPause: _subscription.pause,
|
| - onResume: _subscription.resume);
|
| }
|
|
|
| StreamSubscription listen(void onData(message),
|
| @@ -839,8 +824,6 @@ class _WebSocketImpl extends Stream implements WebSocket {
|
| _consumer.add(new _WebSocketPing());
|
| _pingTimer = new Timer(_pingInterval, () {
|
| // No pong received.
|
| - _closeCode = WebSocketStatus.GOING_AWAY;
|
| - _subscription.cancel();
|
| close(WebSocketStatus.GOING_AWAY);
|
| });
|
| });
|
|
|