Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
kombu_manager.py135 linesDownload Raw Back to socketio
1import pickle2import time3import uuid4 5try:6    import kombu7except ImportError:8    kombu = None9 10from .pubsub_manager import PubSubManager11 12 13class KombuManager(PubSubManager):  # pragma: no cover14    """Client manager that uses kombu for inter-process messaging.15 16    This class implements a client manager backend for event sharing across17    multiple processes, using RabbitMQ, Redis or any other messaging mechanism18    supported by `kombu <http://kombu.readthedocs.org/en/latest/>`_.19 20    To use a kombu backend, initialize the :class:`Server` instance as21    follows::22 23        url = 'amqp://user:password@hostname:port//'24        server = socketio.Server(client_manager=socketio.KombuManager(url))25 26    :param url: The connection URL for the backend messaging queue. Example27                connection URLs are ``'amqp://guest:guest@localhost:5672//'``28                and ``'redis://localhost:6379/'`` for RabbitMQ and Redis29                respectively. Consult the `kombu documentation30                <http://kombu.readthedocs.org/en/latest/userguide\31                /connections.html#urls>`_ for more on how to construct32                connection URLs.33    :param channel: The channel name on which the server sends and receives34                    notifications. Must be the same in all the servers.35    :param write_only: If set to ``True``, only initialize to emit events. The36                       default of ``False`` initializes the class for emitting37                       and receiving.38    :param connection_options: additional keyword arguments to be passed to39                               ``kombu.Connection()``.40    :param exchange_options: additional keyword arguments to be passed to41                             ``kombu.Exchange()``.42    :param queue_options: additional keyword arguments to be passed to43                          ``kombu.Queue()``.44    :param producer_options: additional keyword arguments to be passed to45                             ``kombu.Producer()``.46    """47    name = 'kombu'48 49    def __init__(self, url='amqp://guest:guest@localhost:5672//',50                 channel='socketio', write_only=False, logger=None,51                 connection_options=None, exchange_options=None,52                 queue_options=None, producer_options=None):53        if kombu is None:54            raise RuntimeError('Kombu package is not installed '55                               '(Run "pip install kombu" in your '56                               'virtualenv).')57        super().__init__(channel=channel, write_only=write_only, logger=logger)58        self.url = url59        self.connection_options = connection_options or {}60        self.exchange_options = exchange_options or {}61        self.queue_options = queue_options or {}62        self.producer_options = producer_options or {}63        self.publisher_connection = self._connection()64 65    def initialize(self):66        super().initialize()67 68        monkey_patched = True69        if self.server.async_mode == 'eventlet':70            from eventlet.patcher import is_monkey_patched71            monkey_patched = is_monkey_patched('socket')72        elif 'gevent' in self.server.async_mode:73            from gevent.monkey import is_module_patched74            monkey_patched = is_module_patched('socket')75        if not monkey_patched:76            raise RuntimeError(77                'Kombu requires a monkey patched socket library to work '78                'with ' + self.server.async_mode)79 80    def _connection(self):81        return kombu.Connection(self.url, **self.connection_options)82 83    def _exchange(self):84        options = {'type': 'fanout', 'durable': False}85        options.update(self.exchange_options)86        return kombu.Exchange(self.channel, **options)87 88    def _queue(self):89        queue_name = 'flask-socketio.' + str(uuid.uuid4())90        options = {'durable': False, 'queue_arguments': {'x-expires': 300000}}91        options.update(self.queue_options)92        return kombu.Queue(queue_name, self._exchange(), **options)93 94    def _producer_publish(self, connection):95        producer = connection.Producer(exchange=self._exchange(),96                                       **self.producer_options)97        return connection.ensure(producer, producer.publish)98 99    def _publish(self, data):100        retry = True101        while True:102            try:103                producer_publish = self._producer_publish(104                    self.publisher_connection)105                producer_publish(pickle.dumps(data))106                break107            except (OSError, kombu.exceptions.KombuError):108                if retry:109                    self._get_logger().error('Cannot publish to rabbitmq... '110                                             'retrying')111                    retry = False112                else:113                    self._get_logger().error(114                        'Cannot publish to rabbitmq... giving up')115                    break116 117    def _listen(self):118        reader_queue = self._queue()119        retry_sleep = 1120        while True:121            try:122                with self._connection() as connection:123                    with connection.SimpleQueue(reader_queue) as queue:124                        while True:125                            message = queue.get(block=True)126                            message.ack()127                            yield message.payload128                            retry_sleep = 1129            except (OSError, kombu.exceptions.KombuError):130                self._get_logger().error(131                    'Cannot receive from rabbitmq... '132                    'retrying in {} secs'.format(retry_sleep))133                time.sleep(retry_sleep)134                retry_sleep = min(retry_sleep * 2, 60)135