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