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

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

Issue 19970004: Add ability to send a ping interval on a WebSocket. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Add pingInterval test for server-side and add close-check when sending Created 7 years, 5 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 | « sdk/lib/io/websocket.dart ('k') | tests/standalone/io/web_socket_ping_test.dart » ('j') | no next file with comments »
Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
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;
« no previous file with comments | « sdk/lib/io/websocket.dart ('k') | tests/standalone/io/web_socket_ping_test.dart » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698