codekingpro/portable-devtools
114k
1import random2 3import engineio4 5from . import base_client6from . import exceptions7from . import packet8 9 10class Client(base_client.BaseClient):11 """A Socket.IO client.12 13 This class implements a fully compliant Socket.IO web client with support14 for websocket and long-polling transports.15 16 :param reconnection: ``True`` if the client should automatically attempt to17 reconnect to the server after an interruption, or18 ``False`` to not reconnect. The default is ``True``.19 :param reconnection_attempts: How many reconnection attempts to issue20 before giving up, or 0 for infinite attempts.21 The default is 0.22 :param reconnection_delay: How long to wait in seconds before the first23 reconnection attempt. Each successive attempt24 doubles this delay.25 :param reconnection_delay_max: The maximum delay between reconnection26 attempts.27 :param randomization_factor: Randomization amount for each delay between28 reconnection attempts. The default is 0.5,29 which means that each delay is randomly30 adjusted by +/- 50%.31 :param logger: To enable logging set to ``True`` or pass a logger object to32 use. To disable logging set to ``False``. The default is33 ``False``. Note that fatal errors are logged even when34 ``logger`` is ``False``.35 :param serializer: The serialization method to use when transmitting36 packets. Valid values are ``'default'``, ``'pickle'``,37 ``'msgpack'`` and ``'cbor'``. Alternatively, a subclass38 of the :class:`Packet` class with custom implementations39 of the ``encode()`` and ``decode()`` methods can be40 provided. Client and server must use compatible41 serializers.42 :param json: An alternative json module to use for encoding and decoding43 packets. Custom json modules must have ``dumps`` and ``loads``44 functions that are compatible with the standard library45 versions.46 :param handle_sigint: Set to ``True`` to automatically handle disconnection47 when the process is interrupted, or to ``False`` to48 leave interrupt handling to the calling application.49 Interrupt handling can only be enabled when the50 client instance is created in the main thread.51 52 The Engine.IO configuration supports the following settings:53 54 :param request_timeout: A timeout in seconds for requests. The default is55 5 seconds.56 :param http_session: an initialized ``requests.Session`` object to be used57 when sending requests to the server. Use it if you58 need to add special client options such as proxy59 servers, SSL certificates, etc.60 :param ssl_verify: ``True`` to verify SSL certificates, or ``False`` to61 skip SSL certificate verification, allowing62 connections to servers with self signed certificates.63 The default is ``True``.64 :param engineio_logger: To enable Engine.IO logging set to ``True`` or pass65 a logger object to use. To disable logging set to66 ``False``. The default is ``False``. Note that67 fatal errors are logged even when68 ``engineio_logger`` is ``False``.69 """70 def connect(self, url, headers={}, auth=None, transports=None,71 namespaces=None, socketio_path='socket.io', wait=True,72 wait_timeout=1, retry=False):73 """Connect to a Socket.IO server.74 75 :param url: The URL of the Socket.IO server. It can include custom76 query string parameters if required by the server. If a77 function is provided, the client will invoke it to obtain78 the URL each time a connection or reconnection is79 attempted.80 :param headers: A dictionary with custom headers to send with the81 connection request. If a function is provided, the82 client will invoke it to obtain the headers dictionary83 each time a connection or reconnection is attempted.84 :param auth: Authentication data passed to the server with the85 connection request, normally a dictionary with one or86 more string key/value pairs. If a function is provided,87 the client will invoke it to obtain the authentication88 data each time a connection or reconnection is attempted.89 :param transports: The list of allowed transports. Valid transports90 are ``'polling'`` and ``'websocket'``. If not91 given, the polling transport is connected first,92 then an upgrade to websocket is attempted.93 :param namespaces: The namespaces to connect as a string or list of94 strings. If not given, the namespaces that have95 registered event handlers are connected.96 :param socketio_path: The endpoint where the Socket.IO server is97 installed. The default value is appropriate for98 most cases.99 :param wait: if set to ``True`` (the default) the call only returns100 when all the namespaces are connected. If set to101 ``False``, the call returns as soon as the Engine.IO102 transport is connected, and the namespaces will connect103 in the background.104 :param wait_timeout: How long the client should wait for the105 connection. The default is 1 second. This106 argument is only considered when ``wait`` is set107 to ``True``.108 :param retry: Apply the reconnection logic if the initial connection109 attempt fails. The default is ``False``.110 111 Example usage::112 113 sio = socketio.Client()114 sio.connect('http://localhost:5000')115 """116 if self.connected:117 raise exceptions.ConnectionError('Already connected')118 119 self.connection_url = url120 self.connection_headers = headers121 self.connection_auth = auth122 self.connection_transports = transports123 self.connection_namespaces = namespaces124 self.socketio_path = socketio_path125 126 if namespaces is None:127 namespaces = list(set(self.handlers.keys()).union(128 set(self.namespace_handlers.keys())))129 if len(namespaces) == 0:130 namespaces = ['/']131 elif isinstance(namespaces, str):132 namespaces = [namespaces]133 self.connection_namespaces = namespaces134 self.namespaces = {}135 if self._connect_event is None:136 self._connect_event = self.eio.create_event()137 else:138 self._connect_event.clear()139 real_url = self._get_real_value(self.connection_url)140 real_headers = self._get_real_value(self.connection_headers)141 try:142 self.eio.connect(real_url, headers=real_headers,143 transports=transports,144 engineio_path=socketio_path)145 except engineio.exceptions.ConnectionError as exc:146 for n in self.connection_namespaces:147 self._trigger_event(148 'connect_error', n,149 exc.args[1] if len(exc.args) > 1 else exc.args[0])150 if retry: # pragma: no cover151 self._handle_reconnect()152 if self.eio.state == 'connected':153 return154 raise exceptions.ConnectionError(exc.args[0]) from None155 156 if wait:157 while self._connect_event.wait(timeout=wait_timeout):158 self._connect_event.clear()159 if set(self.namespaces) == set(self.connection_namespaces):160 break161 if set(self.namespaces) != set(self.connection_namespaces):162 self.disconnect()163 raise exceptions.ConnectionError(164 'One or more namespaces failed to connect')165 166 self.connected = True167 168 def wait(self):169 """Wait until the connection with the server ends.170 171 Client applications can use this function to block the main thread172 during the life of the connection.173 """174 while True:175 self.eio.wait()176 self.sleep(1) # give the reconnect task time to start up177 if not self._reconnect_task:178 break179 self._reconnect_task.join()180 if self.eio.state != 'connected':181 break182 183 def emit(self, event, data=None, namespace=None, callback=None):184 """Emit a custom event to the server.185 186 :param event: The event name. It can be any string. The event names187 ``'connect'``, ``'message'`` and ``'disconnect'`` are188 reserved and should not be used.189 :param data: The data to send to the server. Data can be of190 type ``str``, ``bytes``, ``list`` or ``dict``. To send191 multiple arguments, use a tuple where each element is of192 one of the types indicated above.193 :param namespace: The Socket.IO namespace for the event. If this194 argument is omitted the event is emitted to the195 default namespace.196 :param callback: If given, this function will be called to acknowledge197 the server has received the message. The arguments198 that will be passed to the function are those provided199 by the server.200 201 Note: this method is not thread safe. If multiple threads are emitting202 at the same time on the same client connection, messages composed of203 multiple packets may end up being sent in an incorrect sequence. Use204 standard concurrency solutions (such as a Lock object) to prevent this205 situation.206 """207 namespace = namespace or '/'208 if namespace not in self.namespaces:209 raise exceptions.BadNamespaceError(210 namespace + ' is not a connected namespace.')211 self.logger.info('Emitting event "%s" [%s]', event, namespace)212 if callback is not None:213 id = self._generate_ack_id(namespace, callback)214 else:215 id = None216 # tuples are expanded to multiple arguments, everything else is sent217 # as a single argument218 if isinstance(data, tuple):219 data = list(data)220 elif data is not None:221 data = [data]222 else:223 data = []224 self._send_packet(self.packet_class(packet.EVENT, namespace=namespace,225 data=[event] + data, id=id))226 227 def send(self, data, namespace=None, callback=None):228 """Send a message to the server.229 230 This function emits an event with the name ``'message'``. Use231 :func:`emit` to issue custom event names.232 233 :param data: The data to send to the server. Data can be of234 type ``str``, ``bytes``, ``list`` or ``dict``. To send235 multiple arguments, use a tuple where each element is of236 one of the types indicated above.237 :param namespace: The Socket.IO namespace for the event. If this238 argument is omitted the event is emitted to the239 default namespace.240 :param callback: If given, this function will be called to acknowledge241 the server has received the message. The arguments242 that will be passed to the function are those provided243 by the server.244 """245 self.emit('message', data=data, namespace=namespace,246 callback=callback)247 248 def call(self, event, data=None, namespace=None, timeout=60):249 """Emit a custom event to the server and wait for the response.250 251 This method issues an emit with a callback and waits for the callback252 to be invoked before returning. If the callback isn't invoked before253 the timeout, then a ``TimeoutError`` exception is raised. If the254 Socket.IO connection drops during the wait, this method still waits255 until the specified timeout.256 257 :param event: The event name. It can be any string. The event names258 ``'connect'``, ``'message'`` and ``'disconnect'`` are259 reserved and should not be used.260 :param data: The data to send to the server. Data can be of261 type ``str``, ``bytes``, ``list`` or ``dict``. To send262 multiple arguments, use a tuple where each element is of263 one of the types indicated above.264 :param namespace: The Socket.IO namespace for the event. If this265 argument is omitted the event is emitted to the266 default namespace.267 :param timeout: The waiting timeout. If the timeout is reached before268 the server acknowledges the event, then a269 ``TimeoutError`` exception is raised.270 271 Note: this method is not thread safe. If multiple threads are emitting272 at the same time on the same client connection, messages composed of273 multiple packets may end up being sent in an incorrect sequence. Use274 standard concurrency solutions (such as a Lock object) to prevent this275 situation.276 """277 callback_event = self.eio.create_event()278 callback_args = []279 280 def event_callback(*args):281 callback_args.append(args)282 callback_event.set()283 284 self.emit(event, data=data, namespace=namespace,285 callback=event_callback)286 if not callback_event.wait(timeout=timeout):287 raise exceptions.TimeoutError()288 return callback_args[0] if len(callback_args[0]) > 1 \289 else callback_args[0][0] if len(callback_args[0]) == 1 \290 else None291 292 def disconnect(self):293 """Disconnect from the server."""294 # here we just request the disconnection295 # later in _handle_eio_disconnect we invoke the disconnect handler296 for n in self.namespaces:297 self._send_packet(self.packet_class(298 packet.DISCONNECT, namespace=n))299 self.eio.disconnect(abort=True)300 301 def start_background_task(self, target, *args, **kwargs):302 """Start a background task using the appropriate async model.303 304 This is a utility function that applications can use to start a305 background task using the method that is compatible with the306 selected async mode.307 308 :param target: the target function to execute.309 :param args: arguments to pass to the function.310 :param kwargs: keyword arguments to pass to the function.311 312 This function returns an object that represents the background task,313 on which the ``join()`` methond can be invoked to wait for the task to314 complete.315 """316 return self.eio.start_background_task(target, *args, **kwargs)317 318 def sleep(self, seconds=0):319 """Sleep for the requested amount of time using the appropriate async320 model.321 322 This is a utility function that applications can use to put a task to323 sleep without having to worry about using the correct call for the324 selected async mode.325 """326 return self.eio.sleep(seconds)327 328 def _get_real_value(self, value):329 """Return the actual value, for parameters that can also be given as330 callables."""331 if not callable(value):332 return value333 return value()334 335 def _send_packet(self, pkt):336 """Send a Socket.IO packet to the server."""337 encoded_packet = pkt.encode()338 if isinstance(encoded_packet, list):339 for ep in encoded_packet:340 self.eio.send(ep)341 else:342 self.eio.send(encoded_packet)343 344 def _handle_connect(self, namespace, data):345 namespace = namespace or '/'346 if namespace not in self.namespaces:347 self.logger.info('Namespace {} is connected'.format(namespace))348 self.namespaces[namespace] = (data or {}).get('sid', self.sid)349 self._trigger_event('connect', namespace=namespace)350 self._connect_event.set()351 352 def _handle_disconnect(self, namespace):353 if not self.connected:354 return355 namespace = namespace or '/'356 self._trigger_event('disconnect', namespace=namespace)357 self._trigger_event('__disconnect_final', namespace=namespace)358 if namespace in self.namespaces:359 del self.namespaces[namespace]360 if not self.namespaces:361 self.connected = False362 self.eio.disconnect(abort=True)363 364 def _handle_event(self, namespace, id, data):365 namespace = namespace or '/'366 self.logger.info('Received event "%s" [%s]', data[0], namespace)367 r = self._trigger_event(data[0], namespace, *data[1:])368 if id is not None:369 # send ACK packet with the response returned by the handler370 # tuples are expanded as multiple arguments371 if r is None:372 data = []373 elif isinstance(r, tuple):374 data = list(r)375 else:376 data = [r]377 self._send_packet(self.packet_class(378 packet.ACK, namespace=namespace, id=id, data=data))379 380 def _handle_ack(self, namespace, id, data):381 namespace = namespace or '/'382 self.logger.info('Received ack [%s]', namespace)383 callback = None384 try:385 callback = self.callbacks[namespace][id]386 except KeyError:387 # if we get an unknown callback we just ignore it388 self.logger.warning('Unknown callback received, ignoring.')389 else:390 del self.callbacks[namespace][id]391 if callback is not None:392 callback(*data)393 394 def _handle_error(self, namespace, data):395 namespace = namespace or '/'396 self.logger.info('Connection to namespace {} was rejected'.format(397 namespace))398 if data is None:399 data = tuple()400 elif not isinstance(data, (tuple, list)):401 data = (data,)402 self._trigger_event('connect_error', namespace, *data)403 self._connect_event.set()404 if namespace in self.namespaces:405 del self.namespaces[namespace]406 if namespace == '/':407 self.namespaces = {}408 self.connected = False409 410 def _trigger_event(self, event, namespace, *args):411 """Invoke an application event handler."""412 # first see if we have an explicit handler for the event413 handler, args = self._get_event_handler(event, namespace, args)414 if handler:415 return handler(*args)416 417 # or else, forward the event to a namespace handler if one exists418 handler, args = self._get_namespace_handler(namespace, args)419 if handler:420 return handler.trigger_event(event, *args)421 422 def _handle_reconnect(self):423 if self._reconnect_abort is None: # pragma: no cover424 self._reconnect_abort = self.eio.create_event()425 self._reconnect_abort.clear()426 base_client.reconnecting_clients.append(self)427 attempt_count = 0428 current_delay = self.reconnection_delay429 while True:430 delay = current_delay431 current_delay *= 2432 if delay > self.reconnection_delay_max:433 delay = self.reconnection_delay_max434 delay += self.randomization_factor * (2 * random.random() - 1)435 self.logger.info(436 'Connection failed, new attempt in {:.02f} seconds'.format(437 delay))438 if self._reconnect_abort.wait(delay):439 self.logger.info('Reconnect task aborted')440 for n in self.connection_namespaces:441 self._trigger_event('__disconnect_final', namespace=n)442 break443 attempt_count += 1444 try:445 self.connect(self.connection_url,446 headers=self.connection_headers,447 auth=self.connection_auth,448 transports=self.connection_transports,449 namespaces=self.connection_namespaces,450 socketio_path=self.socketio_path,451 retry=False)452 except (exceptions.ConnectionError, ValueError):453 pass454 else:455 self.logger.info('Reconnection successful')456 self._reconnect_task = None457 break458 if self.reconnection_attempts and \459 attempt_count >= self.reconnection_attempts:460 self.logger.info(461 'Maximum reconnection attempts reached, giving up')462 for n in self.connection_namespaces:463 self._trigger_event('__disconnect_final', namespace=n)464 break465 base_client.reconnecting_clients.remove(self)466 467 def _handle_eio_connect(self):468 """Handle the Engine.IO connection event."""469 self.logger.info('Engine.IO connection established')470 self.sid = self.eio.sid471 real_auth = self._get_real_value(self.connection_auth) or {}472 for n in self.connection_namespaces:473 self._send_packet(self.packet_class(474 packet.CONNECT, data=real_auth, namespace=n))475 476 def _handle_eio_message(self, data):477 """Dispatch Engine.IO messages."""478 if self._binary_packet:479 pkt = self._binary_packet480 if pkt.add_attachment(data):481 self._binary_packet = None482 if pkt.packet_type == packet.BINARY_EVENT:483 self._handle_event(pkt.namespace, pkt.id, pkt.data)484 else:485 self._handle_ack(pkt.namespace, pkt.id, pkt.data)486 else:487 pkt = self.packet_class(encoded_packet=data)488 if pkt.packet_type == packet.CONNECT:489 self._handle_connect(pkt.namespace, pkt.data)490 elif pkt.packet_type == packet.DISCONNECT:491 self._handle_disconnect(pkt.namespace)492 elif pkt.packet_type == packet.EVENT:493 self._handle_event(pkt.namespace, pkt.id, pkt.data)494 elif pkt.packet_type == packet.ACK:495 self._handle_ack(pkt.namespace, pkt.id, pkt.data)496 elif pkt.packet_type == packet.BINARY_EVENT or \497 pkt.packet_type == packet.BINARY_ACK:498 self._binary_packet = pkt499 elif pkt.packet_type == packet.CONNECT_ERROR:500 self._handle_error(pkt.namespace, pkt.data)501 else:502 raise ValueError('Unknown packet type.')503 504 def _handle_eio_disconnect(self):505 """Handle the Engine.IO disconnection event."""506 self.logger.info('Engine.IO connection dropped')507 will_reconnect = self.reconnection and self.eio.state == 'connected'508 if self.connected:509 for n in self.namespaces:510 self._trigger_event('disconnect', namespace=n)511 if not will_reconnect:512 self._trigger_event('__disconnect_final', namespace=n)513 self.namespaces = {}514 self.connected = False515 self.callbacks = {}516 self._binary_packet = None517 self.sid = None518 if will_reconnect:519 self._reconnect_task = self.start_background_task(520 self._handle_reconnect)521 522 def _engineio_client_class(self):523 return engineio.Client524 