Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
zmq_manager.py106 linesDownload Raw Back to socketio
1import pickle2import re3 4from .pubsub_manager import PubSubManager5 6 7class ZmqManager(PubSubManager):  # pragma: no cover8    """zmq based client manager.9 10    NOTE: this zmq implementation should be considered experimental at this11    time. At this time, eventlet is required to use zmq.12 13    This class implements a zmq backend for event sharing across multiple14    processes. To use a zmq backend, initialize the :class:`Server` instance as15    follows::16 17        url = 'zmq+tcp://hostname:port1+port2'18        server = socketio.Server(client_manager=socketio.ZmqManager(url))19 20    :param url: The connection URL for the zmq message broker,21                which will need to be provided and running.22    :param channel: The channel name on which the server sends and receives23                    notifications. Must be the same in all the servers.24    :param write_only: If set to ``True``, only initialize to emit events. The25                       default of ``False`` initializes the class for emitting26                       and receiving.27 28    A zmq message broker must be running for the zmq_manager to work.29    you can write your own or adapt one from the following simple broker30    below::31 32        import zmq33 34        receiver = zmq.Context().socket(zmq.PULL)35        receiver.bind("tcp://*:5555")36 37        publisher = zmq.Context().socket(zmq.PUB)38        publisher.bind("tcp://*:5556")39 40        while True:41            publisher.send(receiver.recv())42    """43    name = 'zmq'44 45    def __init__(self, url='zmq+tcp://localhost:5555+5556',46                 channel='socketio',47                 write_only=False,48                 logger=None):49        try:50            from eventlet.green import zmq51        except ImportError:52            raise RuntimeError('zmq package is not installed '53                               '(Run "pip install pyzmq" in your '54                               'virtualenv).')55 56        r = re.compile(r':\d+\+\d+$')57        if not (url.startswith('zmq+tcp://') and r.search(url)):58            raise RuntimeError('unexpected connection string: ' + url)59 60        url = url.replace('zmq+', '')61        (sink_url, sub_port) = url.split('+')62        sink_port = sink_url.split(':')[-1]63        sub_url = sink_url.replace(sink_port, sub_port)64 65        sink = zmq.Context().socket(zmq.PUSH)66        sink.connect(sink_url)67 68        sub = zmq.Context().socket(zmq.SUB)69        sub.setsockopt_string(zmq.SUBSCRIBE, u'')70        sub.connect(sub_url)71 72        self.sink = sink73        self.sub = sub74        self.channel = channel75        super().__init__(channel=channel, write_only=write_only, logger=logger)76 77    def _publish(self, data):78        pickled_data = pickle.dumps(79            {80                'type': 'message',81                'channel': self.channel,82                'data': data83            }84        )85        return self.sink.send(pickled_data)86 87    def zmq_listen(self):88        while True:89            response = self.sub.recv()90            if response is not None:91                yield response92 93    def _listen(self):94        for message in self.zmq_listen():95            if isinstance(message, bytes):96                try:97                    message = pickle.loads(message)98                except Exception:99                    pass100            if isinstance(message, dict) and \101                    message['type'] == 'message' and \102                    message['channel'] == self.channel and \103                    'data' in message:104                yield message['data']105        return106 
codekingpro/portable-devtools · Team Ai