mirror of
https://github.com/wassname/catalyst.git
synced 2026-08-03 12:40:47 +08:00
switching to use push/pull for traffic btw merge and client
This commit is contained in:
@@ -210,7 +210,6 @@ class FinanceTestCase(TestCase):
|
||||
"The merge should be drained of all messages, found {n} remaining." \
|
||||
.format(n=zipline.sim.merge.pending_messages())
|
||||
)
|
||||
|
||||
self.assertEqual(
|
||||
zipline.algorithm.count,
|
||||
zipline.algorithm.incr,
|
||||
|
||||
@@ -49,6 +49,7 @@ class TradeSimulationClient(Component):
|
||||
self.perf.open(self.context)
|
||||
|
||||
def do_work(self):
|
||||
from rdb import set_trace; set_trace()
|
||||
# poll all the sockets
|
||||
socks = dict(self.poll.poll(self.heartbeat_timeout))
|
||||
|
||||
|
||||
@@ -486,10 +486,27 @@ class Component(object):
|
||||
return self.connect_push_socket(self.addresses['merge_address'])
|
||||
|
||||
def bind_result(self):
|
||||
return self.bind_pub_socket(self.addresses['result_address'])
|
||||
return self.bind_push_socket(self.addresses['result_address'])
|
||||
|
||||
def connect_result(self):
|
||||
return self.connect_sub_socket(self.addresses['result_address'])
|
||||
return self.connect_pull_socket(self.addresses['result_address'])
|
||||
|
||||
def bind_push_socket(self, addr):
|
||||
push_socket = self.context.socket(self.zmq.PUSH)
|
||||
push_socket.bind(addr)
|
||||
self.out_socket = push_socket
|
||||
self.sockets.append(push_socket)
|
||||
|
||||
return push_socket
|
||||
|
||||
def connect_pull_socket(self, addr):
|
||||
pull_socket = self.context.socket(self.zmq.PULL)
|
||||
pull_socket.connect(addr)
|
||||
self.sockets.append(pull_socket)
|
||||
self.poll.register(pull_socket, self.zmq.POLLIN)
|
||||
|
||||
return pull_socket
|
||||
|
||||
|
||||
def bind_pull_socket(self, addr):
|
||||
pull_socket = self.context.socket(self.zmq.PULL)
|
||||
|
||||
Reference in New Issue
Block a user