codekingpro/portable-devtools
114k
1import asyncio2import pickle3 4from .async_pubsub_manager import AsyncPubSubManager5 6try:7 import aio_pika8except ImportError:9 aio_pika = None10 11 12class AsyncAioPikaManager(AsyncPubSubManager): # pragma: no cover13 """Client manager that uses aio_pika for inter-process messaging under14 asyncio.15 16 This class implements a client manager backend for event sharing across17 multiple processes, using RabbitMQ18 19 To use a aio_pika backend, initialize the :class:`Server` instance as20 follows::21 22 url = 'amqp://user:password@hostname:port//'23 server = socketio.Server(client_manager=socketio.AsyncAioPikaManager(24 url))25 26 :param url: The connection URL for the backend messaging queue. Example27 connection URLs are ``'amqp://guest:guest@localhost:5672//'``28 for RabbitMQ.29 :param channel: The channel name on which the server sends and receives30 notifications. Must be the same in all the servers.31 With this manager, the channel name is the exchange name32 in rabbitmq33 :param write_only: If set to ``True``, only initialize to emit events. The34 default of ``False`` initializes the class for emitting35 and receiving.36 """37 38 name = 'asyncaiopika'39 40 def __init__(self, url='amqp://guest:guest@localhost:5672//',41 channel='socketio', write_only=False, logger=None):42 if aio_pika is None:43 raise RuntimeError('aio_pika package is not installed '44 '(Run "pip install aio_pika" in your '45 'virtualenv).')46 self.url = url47 self._lock = asyncio.Lock()48 self.publisher_connection = None49 self.publisher_channel = None50 self.publisher_exchange = None51 super().__init__(channel=channel, write_only=write_only, logger=logger)52 53 async def _connection(self):54 return await aio_pika.connect_robust(self.url)55 56 async def _channel(self, connection):57 return await connection.channel()58 59 async def _exchange(self, channel):60 return await channel.declare_exchange(self.channel,61 aio_pika.ExchangeType.FANOUT)62 63 async def _queue(self, channel, exchange):64 queue = await channel.declare_queue(durable=False,65 arguments={'x-expires': 300000})66 await queue.bind(exchange)67 return queue68 69 async def _publish(self, data):70 if self.publisher_connection is None:71 async with self._lock:72 if self.publisher_connection is None:73 self.publisher_connection = await self._connection()74 self.publisher_channel = await self._channel(75 self.publisher_connection76 )77 self.publisher_exchange = await self._exchange(78 self.publisher_channel79 )80 retry = True81 while True:82 try:83 await self.publisher_exchange.publish(84 aio_pika.Message(85 body=pickle.dumps(data),86 delivery_mode=aio_pika.DeliveryMode.PERSISTENT87 ), routing_key='*',88 )89 break90 except aio_pika.AMQPException:91 if retry:92 self._get_logger().error('Cannot publish to rabbitmq... '93 'retrying')94 retry = False95 else:96 self._get_logger().error(97 'Cannot publish to rabbitmq... giving up')98 break99 except aio_pika.exceptions.ChannelInvalidStateError:100 # aio_pika raises this exception when the task is cancelled101 raise asyncio.CancelledError()102 103 async def _listen(self):104 async with (await self._connection()) as connection:105 channel = await self._channel(connection)106 await channel.set_qos(prefetch_count=1)107 exchange = await self._exchange(channel)108 queue = await self._queue(channel, exchange)109 110 retry_sleep = 1111 while True:112 try:113 async with queue.iterator() as queue_iter:114 async for message in queue_iter:115 async with message.process():116 yield pickle.loads(message.body)117 retry_sleep = 1118 except aio_pika.AMQPException:119 self._get_logger().error(120 'Cannot receive from rabbitmq... '121 'retrying in {} secs'.format(retry_sleep))122 await asyncio.sleep(retry_sleep)123 retry_sleep = min(retry_sleep * 2, 60)124 except aio_pika.exceptions.ChannelInvalidStateError:125 # aio_pika raises this exception when the task is cancelled126 raise asyncio.CancelledError()127 