refactor: connexion
This commit is contained in:
+10
-8
@@ -15,21 +15,21 @@ STATE_ACTIVE = 2
|
||||
|
||||
|
||||
class Client(object):
|
||||
def __init__(self, factory=None, id='default'):
|
||||
def __init__(self, factory=None, supervisor=False):
|
||||
assert(factory)
|
||||
|
||||
self._rep_store = ReplicationGraph()
|
||||
self._net_client = ClientNetService(
|
||||
store_reference=self._rep_store,
|
||||
factory=factory,
|
||||
id=id)
|
||||
factory=factory)
|
||||
self._factory = factory
|
||||
self._is_supervisor = supervisor
|
||||
|
||||
def connect(self, address="127.0.0.1", port=5560):
|
||||
def connect(self, id="Default", address="127.0.0.1", port=5560):
|
||||
"""
|
||||
Connect to the server
|
||||
"""
|
||||
self._net_client.connect(address=address, port=port)
|
||||
self._net_client.connect(id=id, address=address, port=port)
|
||||
|
||||
def disconnect(self):
|
||||
"""
|
||||
@@ -96,7 +96,7 @@ class Client(object):
|
||||
|
||||
|
||||
class ClientNetService(threading.Thread):
|
||||
def __init__(self, store_reference=None, factory=None, id="default"):
|
||||
def __init__(self, store_reference=None, factory=None):
|
||||
|
||||
# Threading
|
||||
threading.Thread.__init__(self)
|
||||
@@ -106,7 +106,7 @@ class ClientNetService(threading.Thread):
|
||||
self._exit_event = threading.Event()
|
||||
self._factory = factory
|
||||
self._store_reference = store_reference
|
||||
self._id = id
|
||||
self._id = "None"
|
||||
|
||||
assert(self._factory)
|
||||
|
||||
@@ -114,11 +114,13 @@ class ClientNetService(threading.Thread):
|
||||
self.context = zmq.Context.instance()
|
||||
self.state = STATE_INITIAL
|
||||
|
||||
def connect(self, address='127.0.0.1', port=5560):
|
||||
def connect(self, id=None, address='127.0.0.1', port=5560):
|
||||
"""
|
||||
Network socket setup
|
||||
"""
|
||||
assert(id)
|
||||
if self.state == STATE_INITIAL:
|
||||
self._id = id
|
||||
logger.debug("connecting on {}:{}".format(address, port))
|
||||
self.command = self.context.socket(zmq.DEALER)
|
||||
self.command.setsockopt(zmq.IDENTITY, self._id.encode())
|
||||
|
||||
+22
-22
@@ -66,10 +66,10 @@ class TestClient(unittest.TestCase):
|
||||
factory.register_type(SampleData, RepSampleData)
|
||||
|
||||
server = Server(factory=factory)
|
||||
client = Client(factory=factory, id="client_test_callback")
|
||||
client = Client(factory=factory)
|
||||
|
||||
server.serve(port=5570)
|
||||
client.connect(port=5570)
|
||||
client.connect(port=5570, id="client_test_callback")
|
||||
|
||||
test_state = client.state
|
||||
|
||||
@@ -84,16 +84,16 @@ class TestClient(unittest.TestCase):
|
||||
factory.register_type(SampleData, RepSampleData)
|
||||
|
||||
server = Server(factory=factory)
|
||||
client = Client(factory=factory, id="cli_test_filled_snapshot")
|
||||
client2 = Client(factory=factory, id="client_2")
|
||||
client = Client(factory=factory)
|
||||
client2 = Client(factory=factory)
|
||||
|
||||
server.serve(port=5575)
|
||||
client.connect(port=5575)
|
||||
client.connect(port=5575,id="cli_test_filled_snapshot")
|
||||
|
||||
# Test the key registering
|
||||
data_sample_key = client.register(SampleData())
|
||||
|
||||
client2.connect(port=5575)
|
||||
client2.connect(port=5575, id="client_2")
|
||||
time.sleep(0.2)
|
||||
rep_test_key = client2._rep_store[data_sample_key].uuid
|
||||
|
||||
@@ -112,11 +112,11 @@ class TestClient(unittest.TestCase):
|
||||
server = Server(factory=factory)
|
||||
server.serve(port=5560)
|
||||
|
||||
client = Client(factory=factory, id="cli_test_register_client_data")
|
||||
client.connect(port=5560)
|
||||
client = Client(factory=factory)
|
||||
client.connect(port=5560,id="cli_test_register_client_data")
|
||||
|
||||
client2 = Client(factory=factory, id="cli2_test_register_client_data")
|
||||
client2.connect(port=5560)
|
||||
client2 = Client(factory=factory)
|
||||
client2.connect(port=5560, id="cli2_test_register_client_data")
|
||||
|
||||
# Test the key registering
|
||||
data_sample_key = client.register(SampleData())
|
||||
@@ -139,11 +139,11 @@ class TestClient(unittest.TestCase):
|
||||
server = Server(factory=factory)
|
||||
server.serve(port=5560)
|
||||
|
||||
client = Client(factory=factory, id="cli_test_client_data_intergity")
|
||||
client.connect(port=5560)
|
||||
client = Client(factory=factory)
|
||||
client.connect(port=5560, id="cli_test_client_data_intergity")
|
||||
|
||||
client2 = Client(factory=factory, id="cli2_test_client_data_intergity")
|
||||
client2.connect(port=5560)
|
||||
client2 = Client(factory=factory)
|
||||
client2.connect(port=5560, id="cli2_test_client_data_intergity")
|
||||
|
||||
test_map = {"toto": "test"}
|
||||
# Test the key registering
|
||||
@@ -169,11 +169,11 @@ class TestClient(unittest.TestCase):
|
||||
server = Server(factory=factory)
|
||||
server.serve(port=5560)
|
||||
|
||||
client = Client(factory=factory, id="cli_test_client_data_intergity")
|
||||
client.connect(port=5560)
|
||||
client = Client(factory=factory)
|
||||
client.connect(port=5560, id="cli_test_client_data_intergity")
|
||||
|
||||
client2 = Client(factory=factory, id="cli2_test_client_data_intergity")
|
||||
client2.connect(port=5560)
|
||||
client2 = Client(factory=factory)
|
||||
client2.connect(port=5560, id="cli2_test_client_data_intergity")
|
||||
|
||||
test_map = {"toto": "test"}
|
||||
# Test the key registering
|
||||
@@ -217,12 +217,12 @@ class TestStressClient(unittest.TestCase):
|
||||
factory.register_type(SampleData, RepSampleData)
|
||||
|
||||
server = Server(factory=factory)
|
||||
client = Client(factory=factory, id="cli_test_filled_snapshot")
|
||||
client2 = Client(factory=factory, id="client_2")
|
||||
client = Client(factory=factory)
|
||||
client2 = Client(factory=factory)
|
||||
|
||||
server.serve(port=5575)
|
||||
client.connect(port=5575)
|
||||
client2.connect(port=5575)
|
||||
client.connect(port=5575,id="cli_test_filled_snapshot")
|
||||
client2.connect(port=5575,id="client_2")
|
||||
|
||||
# Test the key registering
|
||||
for i in range(10000):
|
||||
|
||||
Reference in New Issue
Block a user