Define initialize, leader, optime and members string constansts in AbstractDCS

Also define following properties:
* initialize_path
* members_path
* member_path
* leader_path
* leader_optime_path

And replace any occurrences of these strings or client_path calls in a
etcd and zookeeper implementations with given constants and properties.
This commit is contained in:
Alexander Kukushkin
2015-09-09 15:10:45 +02:00
parent e90b14cd3b
commit 5bdb18761b
4 changed files with 53 additions and 29 deletions
+27 -1
View File
@@ -74,6 +74,12 @@ class AbstractDCS:
__metaclass__ = abc.ABCMeta
_INITIALIZE = 'initialize'
_LEADER = 'leader'
_MEMBERS = 'members/'
_OPTIME = 'optime'
_LEADER_OPTIME = _OPTIME + '/' + _LEADER
def __init__(self, name, config):
"""
:param name: name of current instance (the same value as `~Postgresql.name`)
@@ -85,7 +91,27 @@ class AbstractDCS:
self._base_path = '/service/' + self._scope
def client_path(self, path):
return self._base_path + path
return '/'.join([self._base_path, path.lstrip('/')])
@property
def initialize_path(self):
return self.client_path(self._INITIALIZE)
@property
def members_path(self):
return self.client_path(self._MEMBERS)
@property
def member_path(self):
return self.client_path(self._MEMBERS + self._name)
@property
def leader_path(self):
return self.client_path(self._LEADER)
@property
def leader_optime_path(self):
return self.client_path(self._LEADER_OPTIME)
@abc.abstractmethod
def get_cluster(self):
+11 -13
View File
@@ -179,17 +179,17 @@ class Etcd(AbstractDCS):
nodes = {os.path.relpath(node.key, result.key): node for node in result.leaves}
# get initialize flag
initialize = bool(nodes.get('initialize', False))
initialize = bool(nodes.get(self._INITIALIZE, False))
# get last leader operation
last_leader_operation = nodes.get('optime/leader', None)
last_leader_operation = nodes.get(self._LEADER_OPTIME, None)
last_leader_operation = 0 if last_leader_operation is None else int(last_leader_operation.value)
# get list of members
members = [self.member(n) for k, n in nodes.items() if k.startswith('members/') and len(k.split('/')) == 2]
members = [self.member(n) for k, n in nodes.items() if k.startswith(self._MEMBERS) and k.count('/') == 1]
# get leader
leader = nodes.get('leader', None)
leader = nodes.get(self._LEADER, None)
if leader:
member = Member(-1, leader.value, None, None, None, None)
member = ([m for m in members if m.name == leader.value] or [member])[0]
@@ -206,17 +206,15 @@ class Etcd(AbstractDCS):
@catch_etcd_errors
def touch_member(self, connection_string, ttl=None):
return self.retry(self.client.set, self.client_path('/members/' + self._name),
connection_string, ttl or self.member_ttl)
return self.retry(self.client.set, self.member_path, connection_string, ttl or self.member_ttl)
@catch_etcd_errors
def take_leader(self):
return self.retry(self.client.set, self.client_path('/leader'), self._name, self.ttl)
return self.retry(self.client.set, self.leader_path, self._name, self.ttl)
def attempt_to_acquire_leader(self):
try:
return not self.retry(self.client.write, self.client_path('/leader'),
self._name, ttl=self.ttl, prevExist=False) is None
return bool(self.retry(self.client.write, self.leader_path, self._name, ttl=self.ttl, prevExist=False))
except etcd.EtcdAlreadyExist:
logger.info('Could not take out TTL lock')
except (RetryFailedError, etcd.EtcdException):
@@ -225,11 +223,11 @@ class Etcd(AbstractDCS):
@catch_etcd_errors
def write_leader_optime(self, state_handler):
return self.client.set(self.client_path('/optime/leader'), state_handler.last_operation())
return self.client.set(self.leader_optime_path, state_handler.last_operation())
@catch_etcd_errors
def update_leader(self, state_handler):
ret = self.retry(self.client.test_and_set, self.client_path('/leader'), self._name, self._name, self.ttl)
ret = self.retry(self.client.test_and_set, self.leader_path, self._name, self._name, self.ttl)
ret and self.write_leader_optime(state_handler)
return ret
@@ -239,7 +237,7 @@ class Etcd(AbstractDCS):
@catch_etcd_errors
def delete_leader(self):
return self.client.delete(self.client_path('/leader'), prevValue=self._name)
return self.client.delete(self.leader_path, prevValue=self._name)
def watch(self, timeout):
# watch on leader key changes if it is defined and current node is not lock owner
@@ -249,7 +247,7 @@ class Etcd(AbstractDCS):
while index and timeout >= 1: # when timeout is too small urllib3 doesn't have enough time to connect
try:
res = self.client.watch(self.client_path('/leader'), index=index + 1, timeout=timeout)
res = self.client.watch(self.leader_path, index=index + 1, timeout=timeout)
if res.action not in ['set', 'compareAndSwap'] or res.value != self.cluster.leader.name:
return
index = res.modifiedIndex
+13 -13
View File
@@ -107,9 +107,9 @@ class ZooKeeper(AbstractDCS):
self.fetch_cluster = True
self.cluster_event.set()
def get_node(self, name, watch=None):
def get_node(self, key, watch=None):
try:
return self.client.get(self.client_path(name), watch)
return self.client.get(key, watch)
except NoNodeError:
pass
except:
@@ -123,21 +123,21 @@ class ZooKeeper(AbstractDCS):
def load_members(self):
members = []
for member in self.client.get_children(self.client_path('/members'), self.cluster_watcher):
data = self.get_node('/members/' + member)
for member in self.client.get_children(self.members_path, self.cluster_watcher):
data = self.get_node(self.member_path)
if data is not None:
members.append(self.member(member, *data))
return members
def _inner_load_cluster(self):
self.cluster_event.clear()
leader = self.get_node('/leader', self.cluster_watcher)
leader = self.get_node(self.leader_path, self.cluster_watcher)
self.members = self.load_members()
if leader:
client_id = self.client.client_id
if leader[0] == self._name and client_id is not None and client_id[0] != leader[1].ephemeralOwner:
logger.info('I am leader but not owner of the session. Removing leader node')
self.client.delete(self.client_path('/leader'))
self.client.delete(self.leader_path)
leader = None
if leader:
@@ -148,7 +148,7 @@ class ZooKeeper(AbstractDCS):
self.leader = leader
if self.fetch_cluster:
last_leader_operation = self.get_node('/optime/leader')
last_leader_operation = self.get_node(self.leader_optime_path)
if last_leader_operation:
self.last_leader_operation = int(last_leader_operation[0])
@@ -167,24 +167,24 @@ class ZooKeeper(AbstractDCS):
def _create(self, path, value, **kwargs):
try:
self.client.retry(self.client.create, self.client_path(path), value, **kwargs)
self.client.retry(self.client.create, path, value, **kwargs)
return True
except:
return False
def attempt_to_acquire_leader(self):
ret = self._create('/leader', self._name, makepath=True, ephemeral=True)
ret = self._create(self.leader_path, self._name, makepath=True, ephemeral=True)
ret or logger.info('Could not take out TTL lock')
return ret
def race(self, path):
return self._create(path, self._name, makepath=True)
return self._create(self.client_path(path), self._name, makepath=True)
def touch_member(self, connection_string, ttl=None):
for m in self.members:
if m.name == self._name:
return True
path = self.client_path('/members/' + self._name)
path = self.member_path
try:
self.client.retry(self.client.create, path, connection_string, makepath=True, ephemeral=True)
return True
@@ -204,7 +204,7 @@ class ZooKeeper(AbstractDCS):
last_operation = state_handler.last_operation()
if last_operation != self.last_leader_operation:
self.last_leader_operation = last_operation
path = self.client_path('/optime/leader')
path = self.leader_optime_path
try:
self.client.retry(self.client.set, path, last_operation)
except NoNodeError:
@@ -218,7 +218,7 @@ class ZooKeeper(AbstractDCS):
def delete_leader(self):
if isinstance(self.leader, Leader) and self.leader.name == self._name:
self.client.delete(self.client_path('/leader'))
self.client.delete(self.leader_path)
def watch(self, timeout):
self.cluster_event.wait(timeout)
+2 -2
View File
@@ -93,9 +93,9 @@ def etcd_delete(key, **kwargs):
def etcd_read(key, **kwargs):
if key == '/service/noleader':
if key == '/service/noleader/':
raise DCSError('noleader')
elif key == '/service/nocluster':
elif key == '/service/nocluster/':
raise etcd.EtcdKeyNotFound
response = {"action": "get", "node": {"key": "/service/batman5", "dir": True, "nodes": [