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