Re-arranged the call order of connecting to avoid interrupted connections
This commit is contained in:
parent
dc8e867741
commit
b059cdf02c
1 changed files with 27 additions and 19 deletions
|
|
@ -24,6 +24,7 @@ _SocketIOSession = namedtuple('_SocketIOSession', [
|
||||||
])
|
])
|
||||||
_log = logging.getLogger(__name__)
|
_log = logging.getLogger(__name__)
|
||||||
RETRY_INTERVAL_IN_SECONDS = 1
|
RETRY_INTERVAL_IN_SECONDS = 1
|
||||||
|
MAX_UPGRADE_RETRIES = 3;
|
||||||
|
|
||||||
class BaseNamespace(object):
|
class BaseNamespace(object):
|
||||||
'Define client behavior'
|
'Define client behavior'
|
||||||
|
|
@ -75,15 +76,15 @@ class BaseNamespace(object):
|
||||||
|
|
||||||
def on_noop(self):
|
def on_noop(self):
|
||||||
'Called after server sends a noop; you can override this method'
|
'Called after server sends a noop; you can override this method'
|
||||||
_log.info('%s [noop]', self.path)
|
_log.debug('%s [noop]', self.path)
|
||||||
|
|
||||||
def on_ping(self):
|
def on_ping(self):
|
||||||
'Called after server sends a ping; you can override this method'
|
'Called after server sends a ping; you can override this method'
|
||||||
_log.info('%s [ping]', self.path)
|
_log.debug('%s [ping]', self.path)
|
||||||
|
|
||||||
def on_pong(self):
|
def on_pong(self):
|
||||||
'Called after server sends a pong; you can override this method'
|
'Called after server sends a pong; you can override this method'
|
||||||
_log.info('%s [pong]', self.path)
|
_log.debug('%s [pong]', self.path)
|
||||||
|
|
||||||
def on_open(self, *args):
|
def on_open(self, *args):
|
||||||
_log.info('%s [open] %s', self.path, args)
|
_log.info('%s [open] %s', self.path, args)
|
||||||
|
|
@ -434,16 +435,15 @@ class SocketIO(object):
|
||||||
if packet.type == PacketType.PONG:
|
if packet.type == PacketType.PONG:
|
||||||
_log.debug("[PONG] %s" % repr(packet));
|
_log.debug("[PONG] %s" % repr(packet));
|
||||||
|
|
||||||
self._terminate_heartbeat();
|
|
||||||
|
|
||||||
# Technically we would need to pause the current
|
# Technically we would need to pause the current
|
||||||
# transport (which should be polling in this
|
# transport (which should be polling in this
|
||||||
# implementation), but since we haven't actually
|
# implementation), but since we haven't actually
|
||||||
# started a polling yet, we can upgrade without that.
|
# started a polling yet, we can upgrade without that.
|
||||||
_log.debug("[upgrading] Sending upgrade request");
|
_log.debug("[upgrading] Sending upgrade request");
|
||||||
websocket.send_packet(PacketType.UPGRADE);
|
websocket.send_packet(PacketType.UPGRADE);
|
||||||
self._start_heartbeat(websocket);
|
|
||||||
return websocket;
|
return websocket;
|
||||||
|
else:
|
||||||
|
self._process_packet(packet)
|
||||||
|
|
||||||
def _terminate_heartbeat(self):
|
def _terminate_heartbeat(self):
|
||||||
"""Terminates the heartbeat thread.
|
"""Terminates the heartbeat thread.
|
||||||
|
|
@ -484,26 +484,33 @@ class SocketIO(object):
|
||||||
|
|
||||||
# Negotiate initial transport
|
# Negotiate initial transport
|
||||||
transport = transports.XHR_PollingTransport(self.session, self.is_secure, self.base_url, **self.kw);
|
transport = transports.XHR_PollingTransport(self.session, self.is_secure, self.base_url, **self.kw);
|
||||||
# Update namespaces
|
|
||||||
for path, namespace in self._namespace_by_path.items():
|
|
||||||
namespace._transport = transport
|
|
||||||
transport.connect(path)
|
|
||||||
|
|
||||||
transport.set_timeout(self.session.connection_timeout);
|
transport.set_timeout(self.session.connection_timeout);
|
||||||
|
|
||||||
# Start the heartbeat pacemaker (PING).
|
|
||||||
self._start_heartbeat(transport);
|
|
||||||
|
|
||||||
# If websocket is available, upgrade to it immediately.
|
# If websocket is available, upgrade to it immediately.
|
||||||
# TODO(sean): We could run this on a separate thread for
|
# TODO(sean): We could run this on a separate thread for
|
||||||
# maximum efficiency although that would require some
|
# maximum efficiency although that would require some
|
||||||
# synchronization to ensure buffers are flushed, etc.
|
# synchronization to ensure buffers are flushed, etc.
|
||||||
|
num_retries = 0;
|
||||||
if "websocket" in self.session.server_supported_transports:
|
if "websocket" in self.session.server_supported_transports:
|
||||||
try:
|
while num_retries < MAX_UPGRADE_RETRIES:
|
||||||
return self._upgrade();
|
try:
|
||||||
except:
|
transport = self._upgrade();
|
||||||
_log.warn("[websocket] Failed to upgrade to websocket")
|
break;
|
||||||
pass;
|
except:
|
||||||
|
pass;
|
||||||
|
time.sleep(1);
|
||||||
|
num_retries += 1;
|
||||||
|
|
||||||
|
if num_retries == MAX_UPGRADE_RETRIES:
|
||||||
|
_log.warn("[websocket] Failed to upgrade to websocket");
|
||||||
|
|
||||||
|
# Update namespaces
|
||||||
|
for path, namespace in self._namespace_by_path.items():
|
||||||
|
namespace._transport = transport
|
||||||
|
transport.connect(path)
|
||||||
|
|
||||||
|
# Start the heartbeat pacemaker (PING).
|
||||||
|
self._start_heartbeat(transport);
|
||||||
|
|
||||||
return transport
|
return transport
|
||||||
|
|
||||||
|
|
@ -669,6 +676,7 @@ def _get_socketIO_session(is_secure, base_url, **kw):
|
||||||
return None;
|
return None;
|
||||||
|
|
||||||
handshake = json.loads(packet.payload);
|
handshake = json.loads(packet.payload);
|
||||||
|
_log.info("socket.io client connected to: %s" % base_url);
|
||||||
|
|
||||||
return _SocketIOSession(
|
return _SocketIOSession(
|
||||||
id = handshake["sid"],
|
id = handshake["sid"],
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue