Created
February 4, 2013 20:15
-
-
Save michelp/4709363 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| from gevent import monkey | |
| monkey.patch_all() | |
| import logging | |
| import socket | |
| import struct | |
| import errno | |
| import uuid | |
| import time | |
| import gevent | |
| from zmq import green as zmq | |
| beacon = struct.Struct('3sB16sH') | |
| BEACON_INTERVAL = 1 | |
| log = logging.getLogger(__name__) | |
| class Beaconer(object): | |
| """ZRE beacon emmiter. http://rfc.zeromq.org/spec:20 | |
| This implements only the base UDP beaconing and 0mq socket | |
| interconnection layer. | |
| """ | |
| def __init__(self, callback, address='*', broadcast_port=35713): | |
| self.callback = callback | |
| self.endpoint = 'tcp://%s' % address | |
| self.broadcast_port = broadcast_port | |
| self.peers = {} | |
| def start(self): | |
| """Greenlet to start the beaconer. This sets up zmq context, | |
| sockets, and spawns worker greenlets. | |
| """ | |
| self.context = zmq.Context() | |
| self.poller = zmq.Poller() | |
| self.router = self.context.socket(zmq.ROUTER) | |
| self.port = self.router.bind_to_random_port(self.endpoint) | |
| self.me = uuid.uuid4().bytes | |
| self.broadcaster = socket.socket( | |
| socket.AF_INET, | |
| socket.SOCK_DGRAM, | |
| socket.IPPROTO_UDP) | |
| self.broadcaster.setsockopt( | |
| socket.SOL_SOCKET, | |
| socket.SO_BROADCAST, | |
| 2) | |
| self.broadcaster.setsockopt( | |
| socket.SOL_SOCKET, | |
| socket.SO_REUSEADDR, | |
| 1) | |
| self.broadcaster.bind(('', self.broadcast_port)) | |
| gevent.joinall( | |
| [gevent.spawn(self._send_beacon), | |
| gevent.spawn(self._recv_beacon), | |
| gevent.spawn(self._recv_msg)]) | |
| def _recv_beacon(self): | |
| """Greenlet that receives udp beacons, maintaining connecitons | |
| with a dictionary of peers. | |
| """ | |
| while True: | |
| try: | |
| data, addr = self.broadcaster.recvfrom(beacon.size) | |
| except socket.error: | |
| log.exception('Error recving beacon:') | |
| gevent.sleep(BEACON_INTERVAL) # don't busy error loop | |
| continue | |
| if len(data) != beacon.size: | |
| continue | |
| try: | |
| greet, ver, peer_id, peer_port = beacon.unpack(data) | |
| except Exception: | |
| continue | |
| if greet != 'ZRE': | |
| continue | |
| if peer_id == self.me: | |
| continue | |
| self.handle_peer(peer_id, addr[0], peer_port) | |
| def _send_beacon(self): | |
| """Greenlet that sends a udp beacon at intervals. | |
| """ | |
| while True: | |
| try: | |
| self.broadcaster.sendto( | |
| beacon.pack('ZRE', 1, self.me, self.port), | |
| ('<broadcast>', self.broadcast_port)) | |
| except socket.error: | |
| log.exception('Error sending beacon:') | |
| gevent.sleep(BEACON_INTERVAL) | |
| def _recv_msg(self): | |
| """Greenlet that receives messages from the local ROUTER | |
| socket. | |
| """ | |
| while True: | |
| events = self.poller.poll() | |
| if self.router in events and events[self.router] == zmq.POLLIN: | |
| self.handle_msg(self.router.recv_multipart()) | |
| def handle_peer(self, peer_id, addr, peer_port): | |
| """ Handle a new peer. | |
| Overide this method to handle new peers. By default, connects | |
| a DEALER socket to the new peers broadcast endpoint and | |
| registers it in a dictionary. | |
| """ | |
| if peer_id not in self.peers: | |
| peer = self.context.socket(zmq.DEALER) | |
| peer.setsockopt(zmq.IDENTITY, self.me) | |
| peer_addr = 'tcp://%s:%s' % (addr, peer_port) | |
| log.info('conecting to: %s at %s' % | |
| (str(uuid.UUID(bytes=peer_id)), peer_addr)) | |
| peer.connect(peer_addr) | |
| self.peers[peer_id] = (peer, peer_addr, time.time()) | |
| def handle_msg(self, msg): | |
| """Override this method to customize message handling. | |
| Defaults to calling the callback, if it exists. | |
| """ | |
| gevent.spawn(self.callback, msg) | |
| if __name__ == '__main__': | |
| logging.basicConfig(level=logging.DEBUG) | |
| def test_backer(): | |
| count = 0 | |
| def my_callback(pyre, peer_id, msg): | |
| time.sleep(1) | |
| print msg | |
| count += 1 | |
| pyre.shout(str(count)) | |
| return my_callback | |
| p = Beaconer(test_backer()) | |
| gevent.spawn(p.start).join() |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment