| Index: sdk/lib/io/websocket_impl.dart
|
| diff --git a/sdk/lib/io/websocket_impl.dart b/sdk/lib/io/websocket_impl.dart
|
| index 33aec3f51f2fa1e8bd2c56a34e62d6f3b35c6016..c054b1ec9e0bc301a1004a05f04052951e6be5a3 100644
|
| --- a/sdk/lib/io/websocket_impl.dart
|
| +++ b/sdk/lib/io/websocket_impl.dart
|
| @@ -455,7 +455,7 @@ class _WebSocketOutgoingTransformer extends StreamEventTransformer {
|
| return;
|
| }
|
| if (message is _WebSocketPing) {
|
| - addFrame(_WebSocketOpcode.PONG, message.payload, sink);
|
| + addFrame(_WebSocketOpcode.PING, message.payload, sink);
|
| return;
|
| }
|
| List<int> data;
|
| @@ -571,6 +571,7 @@ class _WebSocketConsumer implements StreamConsumer {
|
| final Socket socket;
|
| StreamController _controller;
|
| StreamSubscription _subscription;
|
| + bool _issuedPause = false;
|
| Completer _closeCompleter = new Completer();
|
| Completer _completer;
|
|
|
| @@ -582,11 +583,27 @@ class _WebSocketConsumer implements StreamConsumer {
|
| }
|
| }
|
|
|
| + void _onPause() {
|
| + if (_subscription != null) {
|
| + _subscription.pause();
|
| + } else {
|
| + _issuedPause = true;
|
| + }
|
| + }
|
| +
|
| + void _onResume() {
|
| + if (_subscription != null) {
|
| + _subscription.resume();
|
| + } else {
|
| + _issuedPause = false;
|
| + }
|
| + }
|
| +
|
| _ensureController() {
|
| if (_controller != null) return;
|
| _controller = new StreamController(sync: true,
|
| - onPause: () => _subscription.pause(),
|
| - onResume: () => _subscription.resume(),
|
| + onPause: _onPause,
|
| + onResume: _onResume,
|
| onCancel: _onListen);
|
| var stream = _controller.stream.transform(
|
| new _WebSocketOutgoingTransformer(webSocket));
|
| @@ -627,6 +644,10 @@ class _WebSocketConsumer implements StreamConsumer {
|
| _done(error);
|
| },
|
| cancelOnError: true);
|
| + if (_issuedPause) {
|
| + _subscription.pause();
|
| + _issuedPause = false;
|
| + }
|
| return _completer.future;
|
| }
|
|
|
| @@ -656,6 +677,9 @@ class _WebSocketImpl extends Stream implements WebSocket {
|
| bool _writeClosed = false;
|
| int _closeCode;
|
| String _closeReason;
|
| + Duration _pingInterval;
|
| + Timer _pingTimer;
|
| + _WebSocketConsumer _consumer;
|
|
|
| int _outCloseCode;
|
| String _outCloseReason;
|
| @@ -734,17 +758,18 @@ class _WebSocketImpl extends Stream implements WebSocket {
|
|
|
| _WebSocketImpl._fromSocket(Socket this._socket,
|
| [bool this._serverSide = false]) {
|
| - var consumer = new _WebSocketConsumer(this, _socket);
|
| - _sink = new _StreamSinkImpl(consumer);
|
| + _consumer = new _WebSocketConsumer(this, _socket);
|
| + _sink = new _StreamSinkImpl(_consumer);
|
| _readyState = WebSocket.OPEN;
|
|
|
| var transformer = new _WebSocketProtocolTransformer(_serverSide);
|
| _socket.transform(transformer).listen(
|
| (data) {
|
| if (data is _WebSocketPing) {
|
| - consumer.add(new _WebSocketPong(data.payload));
|
| + if (!_writeClosed) _consumer.add(new _WebSocketPong(data.payload));
|
| } else if (data is _WebSocketPong) {
|
| - // TODO(ajohnsen): Notify pong?
|
| + // Simply set pingInterval, as it'll cancel any timers.
|
| + pingInterval = _pingInterval;
|
| } else {
|
| _controller.add(data);
|
| }
|
| @@ -785,6 +810,25 @@ class _WebSocketImpl extends Stream implements WebSocket {
|
| cancelOnError: cancelOnError);
|
| }
|
|
|
| + Duration get pingnterval => _pingInterval;
|
| +
|
| + void set pingInterval(Duration interval) {
|
| + if (_writeClosed) return;
|
| + if (_pingTimer != null) _pingTimer.cancel();
|
| + _pingInterval = interval;
|
| +
|
| + if (_pingInterval == null) return;
|
| +
|
| + _pingTimer = new Timer(_pingInterval, () {
|
| + if (_writeClosed) return;
|
| + _consumer.add(new _WebSocketPing());
|
| + _pingTimer = new Timer(_pingInterval, () {
|
| + // No pong received.
|
| + close(WebSocketStatus.GOING_AWAY);
|
| + });
|
| + });
|
| + }
|
| +
|
| int get readyState => _readyState;
|
|
|
| String get extensions => null;
|
|
|