codekingpro/portable-devtools
114k
1import asyncio2import logging3import random4 5import engineio6 7from . import base_client8from . import exceptions9from . import packet10 11default_logger = logging.getLogger('socketio.client')12 13 14class AsyncClient(base_client.BaseClient):15 """A Socket.IO client for asyncio.16 17 This class implements a fully compliant Socket.IO web client with support18 for websocket and long-polling transports.19 20 :param reconnection: ``True`` if the client should automatically attempt to21 reconnect to the server after an interruption, or22 ``False`` to not reconnect. The default is ``True``.23 :param reconnection_attempts: How many reconnection attempts to issue24 before giving up, or 0 for infinite attempts.25 The default is 0.26 :param reconnection_delay: How long to wait in seconds before the first27 reconnection attempt. Each successive attempt28 doubles this delay.29 :param reconnection_delay_max: The maximum delay between reconnection30 attempts.31 :param randomization_factor: Randomization amount for each delay between32 reconnection attempts. The default is 0.5,33 which means that each delay is randomly34 adjusted by +/- 50%.35 :param logger: To enable logging set to ``True`` or pass a logger object to36 use. To disable logging set to ``False``. The default is37 ``False``. Note that fatal errors are logged even when38 ``logger`` is ``False``.39 :param json: An alternative json module to use for encoding and decoding40 packets. Custom json modules must have ``dumps`` and ``loads``41 functions that are compatible with the standard library42 versions.43 :param handle_sigint: Set to ``True`` to automatically handle disconnection44 when the process is interrupted, or to ``False`` to45 leave interrupt handling to the calling application.46 Interrupt handling can only be enabled when the47 client instance is created in the main thread.48 49 The Engine.IO configuration supports the following settings:50 51 :param request_timeout: A timeout in seconds for requests. The default is52 5 seconds.53 :param http_session: an initialized ``aiohttp.ClientSession`` object to be54 used when sending requests to the server. Use it if55 you need to add special client options such as proxy56 servers, SSL certificates, etc.57 :param ssl_verify: ``True`` to verify SSL certificates, or ``False`` to58 skip SSL certificate verification, allowing59 connections to servers with self signed certificates.60 The default is ``True``.61 :param engineio_logger: To enable Engine.IO logging set to ``True`` or pass62 a logger object to use. To disable logging set to63 ``False``. The default is ``False``. Note that64 fatal errors are logged even when65 ``engineio_logger`` is ``False``.66 """67 def is_asyncio_based(self):68 return True69 70 async 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 Note: this method is a coroutine.112 113 Example usage::114 115 sio = socketio.AsyncClient()116 sio.connect('http://localhost:5000')117 """118 if self.connected:119 raise exceptions.ConnectionError('Already connected')120 121 self.connection_url = url122 self.connection_headers = headers123 self.connection_auth = auth124 self.connection_transports = transports125 self.connection_namespaces = namespaces126 self.socketio_path = socketio_path127 128 if namespaces is None:129 namespaces = list(set(self.handlers.keys()).union(130 set(self.namespace_handlers.keys())))131 if len(namespaces) == 0:132 namespaces = ['/']133 elif isinstance(namespaces, str):134 namespaces = [namespaces]135 self.connection_namespaces = namespaces136 self.namespaces = {}137 if self._connect_event is None:138 self._connect_event = self.eio.create_event()139 else:140 self._connect_event.clear()141 real_url = await self._get_real_value(self.connection_url)142 real_headers = await self._get_real_value(self.connection_headers)143 try:144 await self.eio.connect(real_url, headers=real_headers,145 transports=transports,146 engineio_path=socketio_path)147 except engineio.exceptions.ConnectionError as exc:148 for n in self.connection_namespaces:149 await self._trigger_event(150 'connect_error', n,151 exc.args[1] if len(exc.args) > 1 else exc.args[0])152 if retry: # pragma: no cover153 await self._handle_reconnect()154 if self.eio.state == 'connected':155 return156 raise exceptions.ConnectionError(exc.args[0]) from None157 158 if wait:159 try:160 while True:161 await asyncio.wait_for(self._connect_event.wait(),162 wait_timeout)163 self._connect_event.clear()164 if set(self.namespaces) == set(self.connection_namespaces):165 break166 except asyncio.TimeoutError:167 pass168 if set(self.namespaces) != set(self.connection_namespaces):169 await self.disconnect()170 raise exceptions.ConnectionError(171 'One or more namespaces failed to connect')172 173 self.connected = True174 175 async def wait(self):176 """Wait until the connection with the server ends.177 178 Client applications can use this function to block the main thread179 during the life of the connection.180 181 Note: this method is a coroutine.182 """183 while True:184 await self.eio.wait()185 await self.sleep(1) # give the reconnect task time to start up186 if not self._reconnect_task:187 break188 await self._reconnect_task189 if self.eio.state != 'connected':190 break191 192 async def emit(self, event, data=None, namespace=None, callback=None):193 """Emit a custom event to the server.194 195 :param event: The event name. It can be any string. The event names196 ``'connect'``, ``'message'`` and ``'disconnect'`` are197 reserved and should not be used.198 :param data: The data to send to the server. Data can be of199 type ``str``, ``bytes``, ``list`` or ``dict``. To send200 multiple arguments, use a tuple where each element is of201 one of the types indicated above.202 :param namespace: The Socket.IO namespace for the event. If this203 argument is omitted the event is emitted to the204 default namespace.205 :param callback: If given, this function will be called to acknowledge206 the server has received the message. The arguments207 that will be passed to the function are those provided208 by the server.209 210 Note: this method is not designed to be used concurrently. If multiple211 tasks are emitting at the same time on the same client connection, then212 messages composed of multiple packets may end up being sent in an213 incorrect sequence. Use standard concurrency solutions (such as a Lock214 object) to prevent this situation.215 216 Note 2: this method is a coroutine.217 """218 namespace = namespace or '/'219 if namespace not in self.namespaces:220 raise exceptions.BadNamespaceError(221 namespace + ' is not a connected namespace.')222 self.logger.info('Emitting event "%s" [%s]', event, namespace)223 if callback is not None:224 id = self._generate_ack_id(namespace, callback)225 else:226 id = None227 # tuples are expanded to multiple arguments, everything else is sent228 # as a single argument229 if isinstance(data, tuple):230 data = list(data)231 elif data is not None:232 data = [data]233 else:234 data = []235 await self._send_packet(self.packet_class(236 packet.EVENT, namespace=namespace, data=[event] + data, id=id))237 238 async def send(self, data, namespace=None, callback=None):239 """Send a message to the server.240 241 This function emits an event with the name ``'message'``. Use242 :func:`emit` to issue custom event names.243 244 :param data: The data to send to the server. Data can be of245 type ``str``, ``bytes``, ``list`` or ``dict``. To send246 multiple arguments, use a tuple where each element is of247 one of the types indicated above.248 :param namespace: The Socket.IO namespace for the event. If this249 argument is omitted the event is emitted to the250 default namespace.251 :param callback: If given, this function will be called to acknowledge252 the server has received the message. The arguments253 that will be passed to the function are those provided254 by the server.255 256 Note: this method is a coroutine.257 """258 await self.emit('message', data=data, namespace=namespace,259 callback=callback)260 261 async def call(self, event, data=None, namespace=None, timeout=60):262 """Emit a custom event to the server and wait for the response.263 264 This method issues an emit with a callback and waits for the callback265 to be invoked before returning. If the callback isn't invoked before266 the timeout, then a ``TimeoutError`` exception is raised. If the267 Socket.IO connection drops during the wait, this method still waits268 until the specified timeout.269 270 :param event: The event name. It can be any string. The event names271 ``'connect'``, ``'message'`` and ``'disconnect'`` are272 reserved and should not be used.273 :param data: The data to send to the server. Data can be of274 type ``str``, ``bytes``, ``list`` or ``dict``. To send275 multiple arguments, use a tuple where each element is of276 one of the types indicated above.277 :param namespace: The Socket.IO namespace for the event. If this278 argument is omitted the event is emitted to the279 default namespace.280 :param timeout: The waiting timeout. If the timeout is reached before281 the server acknowledges the event, then a282 ``TimeoutError`` exception is raised.283 284 Note: this method is not designed to be used concurrently. If multiple285 tasks are emitting at the same time on the same client connection, then286 messages composed of multiple packets may end up being sent in an287 incorrect sequence. Use standard concurrency solutions (such as a Lock288 object) to prevent this situation.289 290 Note 2: this method is a coroutine.291 """292 callback_event = self.eio.create_event()293 callback_args = []294 295 def event_callback(*args):296 callback_args.append(args)297 callback_event.set()298 299 await self.emit(event, data=data, namespace=namespace,300 callback=event_callback)301 try:302 await asyncio.wait_for(callback_event.wait(), timeout)303 except asyncio.TimeoutError:304 raise exceptions.TimeoutError() from None305 return callback_args[0] if len(callback_args[0]) > 1 \306 else callback_args[0][0] if len(callback_args[0]) == 1 \307 else None308 309 async def disconnect(self):310 """Disconnect from the server.311 312 Note: this method is a coroutine.313 """314 # here we just request the disconnection315 # later in _handle_eio_disconnect we invoke the disconnect handler316 for n in self.namespaces:317 await self._send_packet(self.packet_class(packet.DISCONNECT,318 namespace=n))319 await self.eio.disconnect(abort=True)320 321 def start_background_task(self, target, *args, **kwargs):322 """Start a background task using the appropriate async model.323 324 This is a utility function that applications can use to start a325 background task using the method that is compatible with the326 selected async mode.327 328 :param target: the target function to execute.329 :param args: arguments to pass to the function.330 :param kwargs: keyword arguments to pass to the function.331 332 The return value is a ``asyncio.Task`` object.333 """334 return self.eio.start_background_task(target, *args, **kwargs)335 336 async def sleep(self, seconds=0):337 """Sleep for the requested amount of time using the appropriate async338 model.339 340 This is a utility function that applications can use to put a task to341 sleep without having to worry about using the correct call for the342 selected async mode.343 344 Note: this method is a coroutine.345 """346 return await self.eio.sleep(seconds)347 348 async def _get_real_value(self, value):349 """Return the actual value, for parameters that can also be given as350 callables."""351 if not callable(value):352 return value353 if asyncio.iscoroutinefunction(value):354 return await value()355 return value()356 357 async def _send_packet(self, pkt):358 """Send a Socket.IO packet to the server."""359 encoded_packet = pkt.encode()360 if isinstance(encoded_packet, list):361 for ep in encoded_packet:362 await self.eio.send(ep)363 else:364 await self.eio.send(encoded_packet)365 366 async def _handle_connect(self, namespace, data):367 namespace = namespace or '/'368 if namespace not in self.namespaces:369 self.logger.info('Namespace {} is connected'.format(namespace))370 self.namespaces[namespace] = (data or {}).get('sid', self.sid)371 await self._trigger_event('connect', namespace=namespace)372 self._connect_event.set()373 374 async def _handle_disconnect(self, namespace):375 if not self.connected:376 return377 namespace = namespace or '/'378 await self._trigger_event('disconnect', namespace=namespace)379 await self._trigger_event('__disconnect_final', namespace=namespace)380 if namespace in self.namespaces:381 del self.namespaces[namespace]382 if not self.namespaces:383 self.connected = False384 await self.eio.disconnect(abort=True)385 386 async def _handle_event(self, namespace, id, data):387 namespace = namespace or '/'388 self.logger.info('Received event "%s" [%s]', data[0], namespace)389 r = await self._trigger_event(data[0], namespace, *data[1:])390 if id is not None:391 # send ACK packet with the response returned by the handler392 # tuples are expanded as multiple arguments393 if r is None:394 data = []395 elif isinstance(r, tuple):396 data = list(r)397 else:398 data = [r]399 await self._send_packet(self.packet_class(400 packet.ACK, namespace=namespace, id=id, data=data))401 402 async def _handle_ack(self, namespace, id, data):403 namespace = namespace or '/'404 self.logger.info('Received ack [%s]', namespace)405 callback = None406 try:407 callback = self.callbacks[namespace][id]408 except KeyError:409 # if we get an unknown callback we just ignore it410 self.logger.warning('Unknown callback received, ignoring.')411 else:412 del self.callbacks[namespace][id]413 if callback is not None:414 if asyncio.iscoroutinefunction(callback):415 await callback(*data)416 else:417 callback(*data)418 419 async def _handle_error(self, namespace, data):420 namespace = namespace or '/'421 self.logger.info('Connection to namespace {} was rejected'.format(422 namespace))423 if data is None:424 data = tuple()425 elif not isinstance(data, (tuple, list)):426 data = (data,)427 await self._trigger_event('connect_error', namespace, *data)428 self._connect_event.set()429 if namespace in self.namespaces:430 del self.namespaces[namespace]431 if namespace == '/':432 self.namespaces = {}433 self.connected = False434 435 async def _trigger_event(self, event, namespace, *args):436 """Invoke an application event handler."""437 # first see if we have an explicit handler for the event438 handler, args = self._get_event_handler(event, namespace, args)439 if handler:440 if asyncio.iscoroutinefunction(handler):441 try:442 ret = await handler(*args)443 except asyncio.CancelledError: # pragma: no cover444 ret = None445 else:446 ret = handler(*args)447 return ret448 449 # or else, forward the event to a namepsace handler if one exists450 handler, args = self._get_namespace_handler(namespace, args)451 if handler:452 return await handler.trigger_event(event, *args)453 454 async def _handle_reconnect(self):455 if self._reconnect_abort is None: # pragma: no cover456 self._reconnect_abort = self.eio.create_event()457 self._reconnect_abort.clear()458 base_client.reconnecting_clients.append(self)459 attempt_count = 0460 current_delay = self.reconnection_delay461 while True:462 delay = current_delay463 current_delay *= 2464 if delay > self.reconnection_delay_max:465 delay = self.reconnection_delay_max466 delay += self.randomization_factor * (2 * random.random() - 1)467 self.logger.info(468 'Connection failed, new attempt in {:.02f} seconds'.format(469 delay))470 try:471 await asyncio.wait_for(self._reconnect_abort.wait(), delay)472 self.logger.info('Reconnect task aborted')473 for n in self.connection_namespaces:474 await self._trigger_event('__disconnect_final',475 namespace=n)476 break477 except (asyncio.TimeoutError, asyncio.CancelledError):478 pass479 attempt_count += 1480 try:481 await self.connect(self.connection_url,482 headers=self.connection_headers,483 auth=self.connection_auth,484 transports=self.connection_transports,485 namespaces=self.connection_namespaces,486 socketio_path=self.socketio_path,487 retry=False)488 except (exceptions.ConnectionError, ValueError):489 pass490 else:491 self.logger.info('Reconnection successful')492 self._reconnect_task = None493 break494 if self.reconnection_attempts and \495 attempt_count >= self.reconnection_attempts:496 self.logger.info(497 'Maximum reconnection attempts reached, giving up')498 for n in self.connection_namespaces:499 await self._trigger_event('__disconnect_final',500 namespace=n)501 break502 base_client.reconnecting_clients.remove(self)503 504 async def _handle_eio_connect(self):505 """Handle the Engine.IO connection event."""506 self.logger.info('Engine.IO connection established')507 self.sid = self.eio.sid508 real_auth = await self._get_real_value(self.connection_auth) or {}509 for n in self.connection_namespaces:510 await self._send_packet(self.packet_class(511 packet.CONNECT, data=real_auth, namespace=n))512 513 async def _handle_eio_message(self, data):514 """Dispatch Engine.IO messages."""515 if self._binary_packet:516 pkt = self._binary_packet517 if pkt.add_attachment(data):518 self._binary_packet = None519 if pkt.packet_type == packet.BINARY_EVENT:520 await self._handle_event(pkt.namespace, pkt.id, pkt.data)521 else:522 await self._handle_ack(pkt.namespace, pkt.id, pkt.data)523 else:524 pkt = self.packet_class(encoded_packet=data)525 if pkt.packet_type == packet.CONNECT:526 await self._handle_connect(pkt.namespace, pkt.data)527 elif pkt.packet_type == packet.DISCONNECT:528 await self._handle_disconnect(pkt.namespace)529 elif pkt.packet_type == packet.EVENT:530 await self._handle_event(pkt.namespace, pkt.id, pkt.data)531 elif pkt.packet_type == packet.ACK:532 await self._handle_ack(pkt.namespace, pkt.id, pkt.data)533 elif pkt.packet_type == packet.BINARY_EVENT or \534 pkt.packet_type == packet.BINARY_ACK:535 self._binary_packet = pkt536 elif pkt.packet_type == packet.CONNECT_ERROR:537 await self._handle_error(pkt.namespace, pkt.data)538 else:539 raise ValueError('Unknown packet type.')540 541 async def _handle_eio_disconnect(self):542 """Handle the Engine.IO disconnection event."""543 self.logger.info('Engine.IO connection dropped')544 will_reconnect = self.reconnection and self.eio.state == 'connected'545 if self.connected:546 for n in self.namespaces:547 await self._trigger_event('disconnect', namespace=n)548 if not will_reconnect:549 await self._trigger_event('__disconnect_final',550 namespace=n)551 self.namespaces = {}552 self.connected = False553 self.callbacks = {}554 self._binary_packet = None555 self.sid = None556 if will_reconnect:557 self._reconnect_task = self.start_background_task(558 self._handle_reconnect)559 560 def _engineio_client_class(self):561 return engineio.AsyncClient562 