codekingpro/portable-devtools
114k
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 