diff --git a/patroni/dcs.py b/patroni/dcs.py index 942c0610..0c2bb9f7 100644 --- a/patroni/dcs.py +++ b/patroni/dcs.py @@ -75,6 +75,12 @@ class AbstractDCS: __metaclass__ = abc.ABCMeta initialize_key = '/initialize' + _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`) @@ -86,7 +92,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 1ff7e317..d5730ffe 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) @catch_etcd_errors def cancel_initialization(self): @@ -252,7 +250,7 @@ class Etcd(AbstractDCS): index = self.cluster.leader.index 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 ccbc5073..7245b1a3 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 initialize(self): - return self._create(self.initialize_key, self._name, makepath=True) + return self._create(self.initialize_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,12 +218,12 @@ 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 cancel_initialization(self): - node = self.get_node(self.initialize_key) + node = self.get_node(self.initialize_path) if node and node == self._name: - self.client.delete(self.client_path(self.initialize_key)) + self.client.delete(self.initialize_path) def watch(self, timeout): self.cluster_event.wait(timeout) diff --git a/postgres0.yml b/postgres0.yml index ce90da14..659a4db2 100644 --- a/postgres0.yml +++ b/postgres0.yml @@ -47,7 +47,7 @@ postgresql: env_dir: /home/postgres/etc/wal-e.d/env threshold_megabytes: 10240 threshold_backup_size_percentage: 30 - restore: scripts/restore.py + restore: patroni/scripts/restore.py #recovery_conf: #restore_command: cp ../wal_archive/%f %p parameters: diff --git a/postgres1.yml b/postgres1.yml index 763444e8..bc8b6fd1 100644 --- a/postgres1.yml +++ b/postgres1.yml @@ -49,7 +49,7 @@ postgresql: env_dir: /home/postgres/etc/wal-e.d/env threshold_megabytes: 10240 threshold_backup_size_percentage: 30 - restore: scripts/restore.py + restore: patroni/scripts/restore.py parameters: archive_mode: "on" wal_level: hot_standby diff --git a/tests/test_etcd.py b/tests/test_etcd.py index dbcd4b08..66261eeb 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": [ diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py index 83c97b7f..3da71f02 100644 --- a/tests/test_zookeeper.py +++ b/tests/test_zookeeper.py @@ -56,9 +56,9 @@ class MockKazooClient: func(*args, **kwargs) def get(self, path, watch=None): - if path == '/service/test/no_node': + if path == '/no_node': raise NoNodeError - elif path == '/service/test/other_exception': + elif path == '/other_exception': raise Exception() elif '/members/' in path: return (