diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index 747962c7..b80c0387 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -200,7 +200,7 @@ class Consul(AbstractDCS): def touch_member(self, data, **kwargs): cluster = self.cluster - member = cluster and ([m for m in cluster.members if m.name == self._name] or [None])[0] + member = cluster and cluster.get_member(self._name, fallback_to_leader=False) create_member = self.refresh_session() if member and (create_member or member.session != self._session): try: @@ -223,6 +223,9 @@ class Consul(AbstractDCS): @catch_consul_errors def attempt_to_acquire_leader(self, permanent=False): + if not self._session and not permanent: + self.refresh_session() + args = {} if permanent else {'acquire': self._session} ret = self.retry(self._client.kv.put, self.leader_path, self._name, **args) if not ret: diff --git a/patroni/dcs/zookeeper.py b/patroni/dcs/zookeeper.py index b8d79672..201602a9 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -232,7 +232,7 @@ class ZooKeeper(AbstractDCS): def touch_member(self, data, ttl=None, permanent=False): cluster = self.cluster - member = cluster and ([m for m in cluster.members if m.name == self._name] or [None])[0] + member = cluster and cluster.get_member(self._name, fallback_to_leader=False) data = data.encode('utf-8') if member and self._client.client_id is not None and member.session != self._client.client_id[0]: try: diff --git a/patroni/ha.py b/patroni/ha.py index c77fcd78..a552e800 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -165,7 +165,6 @@ class Ha(object): return message else: self.state_handler.promote() - self.touch_member() return promote_message @staticmethod @@ -309,7 +308,6 @@ class Ha(object): self.state_handler.stop() self.state_handler.set_role('demoted') self.dcs.delete_leader() - self.touch_member() self.dcs.reset_cluster() sleep(2) # Give a time to somebody to take the leader lock cluster = self.dcs.get_cluster() @@ -581,10 +579,12 @@ class Ha(object): return None def _run_cycle(self): + dcs_failed = False try: self.load_cluster_from_dcs() - self.touch_member() + if not self.cluster.has_member(self.state_handler.name): + self.touch_member() # cluster has leader key but not initialize key if not (self.cluster.is_unlocked() or self.sysid_valid(self.cluster.initialize)) and self.has_lock(): @@ -645,6 +645,7 @@ class Ha(object): self.state_handler.call_nowait(ACTION_ON_START) self.state_handler.sync_replication_slots(self.cluster) except DCSError: + dcs_failed = True logger.error('Error communicating with DCS') if not self.is_paused() and self.state_handler.is_running() and self.state_handler.is_leader(): self.demote(delete_leader=False) @@ -652,6 +653,9 @@ class Ha(object): return 'DCS is not accessible' except (psycopg2.Error, PostgresConnectionException): return 'Error communicating with PostgreSQL. Will try again later' + finally: + if not dcs_failed: + self.touch_member() def run_cycle(self): with self._async_executor: diff --git a/tests/test_consul.py b/tests/test_consul.py index 5d4d85a9..d176a6c8 100644 --- a/tests/test_consul.py +++ b/tests/test_consul.py @@ -99,6 +99,8 @@ class TestConsul(unittest.TestCase): @patch.object(consul.Consul.KV, 'put', Mock(return_value=False)) def test_take_leader(self): + self.c.set_ttl(20) + self.c.refresh_session = Mock() self.c.take_leader() @patch.object(consul.Consul.KV, 'put', Mock(return_value=True))