Register control socket in every component.

This commit is contained in:
Stephen Diehl
2012-02-24 20:06:12 -05:00
parent a63a7c6dea
commit 2deb6ba254
2 changed files with 22 additions and 23 deletions
+6
View File
@@ -233,6 +233,12 @@ class Component(object):
"""
assert self.controller
self.control_out = self.controller.message_sender()
self.control_in = self.controller.message_listener()
self.poll.register(self.control_in, self.zmq.POLLIN)
self.sockets.extend([self.control_in, self.control_out])
def setup_sync(self):
qutil.LOGGER.debug("Connecting sync client for {id}".format(id=self.get_id))
+16 -23
View File
@@ -10,6 +10,7 @@ class Controller(object):
def __init__(self, pull_socket, pub_socket, context=None, logging = None):
self.associated = []
if not context:
self._ctx = zmq.Context()
@@ -19,11 +20,6 @@ class Controller(object):
self.pull_socket = pull_socket
self.pub_socket = pub_socket
self.pull = self._ctx.socket(zmq.PULL)
self.pub = self._ctx.socket(zmq.PUB)
self.associated = [self.pull, self.pub]
if logging:
self.logging = logging
self.dologging = True
@@ -34,23 +30,13 @@ class Controller(object):
self.success = 0
self.failed = 0
try:
self.pull.bind(pull_socket)
except zmq.ZMQError:
raise Exception('Cannot not bind on %s' % pull_socket)
try:
self.pub.bind(pub_socket)
except zmq.ZMQError:
raise Exception('Cannot not bind on %s' % pub_socket)
def run(self, debug=False):
self.polling = True
if debug:
return self._poll(False)
else:
return self._poll_fast()
#if debug:
return self._poll()
#else:
#return self._poll_fast()
def _poll_fast(self):
"""
@@ -64,11 +50,17 @@ class Controller(object):
mostly used for debugging.
"""
self.pull = self._ctx.socket(zmq.PULL)
self.pub = self._ctx.socket(zmq.PUB)
self.associated.extend([self.pull, self.pub])
self.pull.bind(self.pull_socket)
self.pub.bind(self.pub_socket)
while self.polling:
try:
self.logging.info('msg')
self.pub.send(self.pull.recv())
#self.pub.send(self.pull.recv(copy=False))
except KeyboardInterrupt:
self.polling = False
break
@@ -78,7 +70,8 @@ class Controller(object):
except Exception as e:
# Its common to wrap these in wildcard exceptions so
# that we don't loose messages, ever
self.logging.error(str(e))
if self.logging:
self.logging.error(str(e))
self.failed += 1
continue
@@ -89,7 +82,6 @@ class Controller(object):
"""
s = self._ctx.socket(zmq.PUSH)
s.connect(self.pull_socket)
s.setsockopt(zmq.LINGER, -1)
self.associated.append(s)
return s
@@ -103,6 +95,7 @@ class Controller(object):
s.setsockopt(zmq.SUBSCRIBE, '')
self.associated.append(s)
return s
def destroy(self):
"""
Manual cleanup.