Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
async_redis_manager.py108 linesDownload Raw Back to socketio
1import asyncio2import pickle3 4try:  # pragma: no cover5    from redis import asyncio as aioredis6    from redis.exceptions import RedisError7except ImportError:  # pragma: no cover8    try:9        import aioredis10        from aioredis.exceptions import RedisError11    except ImportError:12        aioredis = None13        RedisError = None14 15from .async_pubsub_manager import AsyncPubSubManager16 17 18class AsyncRedisManager(AsyncPubSubManager):  # pragma: no cover19    """Redis based client manager for asyncio servers.20 21    This class implements a Redis backend for event sharing across multiple22    processes.23 24    To use a Redis backend, initialize the :class:`AsyncServer` instance as25    follows::26 27        url = 'redis://hostname:port/0'28        server = socketio.AsyncServer(29            client_manager=socketio.AsyncRedisManager(url))30 31    :param url: The connection URL for the Redis server. For a default Redis32                store running on the same host, use ``redis://``.  To use an33                SSL connection, use ``rediss://``.34    :param channel: The channel name on which the server sends and receives35                    notifications. Must be the same in all the servers.36    :param write_only: If set to ``True``, only initialize to emit events. The37                       default of ``False`` initializes the class for emitting38                       and receiving.39    :param redis_options: additional keyword arguments to be passed to40                          ``aioredis.from_url()``.41    """42    name = 'aioredis'43 44    def __init__(self, url='redis://localhost:6379/0', channel='socketio',45                 write_only=False, logger=None, redis_options=None):46        if aioredis is None:47            raise RuntimeError('Redis package is not installed '48                               '(Run "pip install redis" in your virtualenv).')49        if not hasattr(aioredis.Redis, 'from_url'):50            raise RuntimeError('Version 2 of aioredis package is required.')51        self.redis_url = url52        self.redis_options = redis_options or {}53        self._redis_connect()54        super().__init__(channel=channel, write_only=write_only, logger=logger)55 56    def _redis_connect(self):57        self.redis = aioredis.Redis.from_url(self.redis_url,58                                             **self.redis_options)59        self.pubsub = self.redis.pubsub(ignore_subscribe_messages=True)60 61    async def _publish(self, data):62        retry = True63        while True:64            try:65                if not retry:66                    self._redis_connect()67                return await self.redis.publish(68                    self.channel, pickle.dumps(data))69            except RedisError:70                if retry:71                    self._get_logger().error('Cannot publish to redis... '72                                             'retrying')73                    retry = False74                else:75                    self._get_logger().error('Cannot publish to redis... '76                                             'giving up')77                    break78 79    async def _redis_listen_with_retries(self):80        retry_sleep = 181        connect = False82        while True:83            try:84                if connect:85                    self._redis_connect()86                    await self.pubsub.subscribe(self.channel)87                    retry_sleep = 188                async for message in self.pubsub.listen():89                    yield message90            except RedisError:91                self._get_logger().error('Cannot receive from redis... '92                                         'retrying in '93                                         '{} secs'.format(retry_sleep))94                connect = True95                await asyncio.sleep(retry_sleep)96                retry_sleep *= 297                if retry_sleep > 60:98                    retry_sleep = 6099 100    async def _listen(self):101        channel = self.channel.encode('utf-8')102        await self.pubsub.subscribe(self.channel)103        async for message in self._redis_listen_with_retries():104            if message['channel'] == channel and \105                    message['type'] == 'message' and 'data' in message:106                yield message['data']107        await self.pubsub.unsubscribe(self.channel)108 
codekingpro/portable-devtools · Team Ai