From 5bdb18761b19a9d041c3e4d968f9407d55f994e1 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 9 Sep 2015 15:10:45 +0200 Subject: [PATCH] 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. --- patroni/dcs.py | 28 +++++++++++++++++++++++++++- patroni/etcd.py | 24 +++++++++++------------- patroni/zookeeper.py | 26 +++++++++++++------------- tests/test_etcd.py | 4 ++-- 4 files changed, 53 insertions(+), 29 deletions(-) diff --git a/patroni/dcs.py b/patroni/dcs.py index 6fb7aea8..a33dea29 100644 --- a/patroni/dcs.py +++ b/patroni/dcs.py @@ -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): diff --git a/patroni/etcd.py b/patroni/etcd.py index 97bd2cf7..7a527f6e 100644 --- a/patroni/etcd.py +++ b/patroni/etcd.py @@ -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 diff --git a/patroni/zookeeper.py b/patroni/zookeeper.py index 29b8c7f1..3901e89d 100644 --- a/patroni/zookeeper.py +++ b/patroni/zookeeper.py @@ -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) diff --git a/tests/test_etcd.py b/tests/test_etcd.py index ab5717ca..94b7f807 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -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": [