Chromium Code Reviews| 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..3d76aeabbfeab6dcafa903117651009b0e38a51d 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); |
|
Søren Gjesse
2013/07/24 11:12:45
Oops!
Anders Johnsen
2013/07/24 12:14:23
Yeah :) Copy paste at its worsts.
|
| 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)); |
| + _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; |
|
Søren Gjesse
2013/07/24 11:12:45
This means that the next ping is sent "pingInterva
Anders Johnsen
2013/07/24 12:14:23
Done.
|
| } 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; |