Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
async_aiopika_manager.py127 linesDownload Raw Back to socketio
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 
codekingpro/portable-devtools · Team Ai