Skip to content

Instantly share code, notes, and snippets.

@michelp
Created February 4, 2013 20:15
Show Gist options
  • Select an option

  • Save michelp/4709363 to your computer and use it in GitHub Desktop.

Select an option

Save michelp/4709363 to your computer and use it in GitHub Desktop.
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