From 10c7fa41f32bac474cf32cdfba24eba3128b9eb6 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 21 Sep 2016 09:42:48 +0200 Subject: [PATCH] Exclude unhealthy nodes when choosing where to clone from (#313) Node MUST have tag clonefrom: true, be in the 'running' state and also we should not try to clone from itself. --- features/cascading_replication.feature | 1 + features/steps/cascading_replication.py | 17 +++++++++++++++++ patroni/dcs/__init__.py | 13 +++++++++++-- patroni/ha.py | 4 ++-- tests/test_ha.py | 1 + 5 files changed, 32 insertions(+), 4 deletions(-) diff --git a/features/cascading_replication.feature b/features/cascading_replication.feature index 4a1672df..a68acd6e 100644 --- a/features/cascading_replication.feature +++ b/features/cascading_replication.feature @@ -8,6 +8,7 @@ Scenario: check a base backup and streaming replication from a replica And replication works from postgres0 to postgres1 after 20 seconds And I create label with "postgres0" in postgres0 data directory And I create label with "postgres1" in postgres1 data directory + And postgres1 has state=running in dcs after 12 seconds And I configure and start postgres2 with a tag replicatefrom postgres1 Then replication works from postgres0 to postgres2 after 30 seconds And there is a label with "postgres1" in postgres2 data directory diff --git a/features/steps/cascading_replication.py b/features/steps/cascading_replication.py index 59399c97..b80ae659 100644 --- a/features/steps/cascading_replication.py +++ b/features/steps/cascading_replication.py @@ -1,3 +1,6 @@ +import json +import time + from behave import step, then @@ -15,3 +18,17 @@ def check_label(context, content, name): @step('I create label with "{content:w}" in {name:w} data directory') def write_label(context, content, name): context.pctl.write_label(name, content) + + +@step('{name:w} has {key:w}={value:w} in dcs after {time_limit:d} seconds') +def check_member(context, name, key, value, time_limit): + max_time = time.time() + int(time_limit) + while time.time() < max_time: + try: + response = json.loads(context.dcs_ctl.query('members/' + name)) + if response.get(key) == value: + return + except Exception: + pass + time.sleep(1) + assert False, "{0} does not have {1}={2} in dcs after {3} seconds".format(name, key, value, time_limit) diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index 9bcc7038..3c2c4429 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -142,6 +142,14 @@ class Member(namedtuple('Member', 'index,name,session,data')): def clonefrom(self): return self.tags.get('clonefrom', False) and bool(self.conn_url) + @property + def state(self): + return self.data.get('state', 'unknown') + + @property + def is_running(self): + return self.state == 'running' + class Leader(namedtuple('Leader', 'index,session,member')): @@ -243,8 +251,9 @@ class Cluster(namedtuple('Cluster', 'initialize,config,leader,last_leader_operat def get_member(self, member_name, fallback_to_leader=True): return ([m for m in self.members if m.name == member_name] or [self.leader if fallback_to_leader else None])[0] - def get_clone_member(self): - candidates = [m for m in self.members if m.clonefrom and (not self.leader or m.name != self.leader.name)] + def get_clone_member(self, exclude): + exclude = [exclude] + [self.leader.name] if self.leader else [] + candidates = [m for m in self.members if m.clonefrom and m.is_running and m.name not in exclude] return candidates[randint(0, len(candidates) - 1)] if candidates else self.leader def is_paused(self): diff --git a/patroni/ha.py b/patroni/ha.py index 5f3a00cc..c77fcd78 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -89,7 +89,7 @@ class Ha(object): def bootstrap(self): if not self.cluster.is_unlocked(): # cluster already has leader - clone_member = self.cluster.get_clone_member() + clone_member = self.cluster.get_clone_member(self.state_handler.name) member_role = 'leader' if clone_member == self.cluster.leader else 'replica' msg = "from {0} '{1}'".format(member_role, clone_member.name) self._async_executor.schedule('bootstrap {0}'.format(msg)) @@ -533,7 +533,7 @@ class Ha(object): self.state_handler.stop('immediate') self.state_handler.remove_data_directory() - clone_member = self.cluster.get_clone_member() + clone_member = self.cluster.get_clone_member(self.state_handler.name) member_role = 'leader' if clone_member == self.cluster.leader else 'replica' self.clone(clone_member, "from {0} '{1}'".format(member_role, clone_member.name)) diff --git a/tests/test_ha.py b/tests/test_ha.py index f62a6f16..7a87d6c1 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -37,6 +37,7 @@ def get_cluster_initialized_without_leader(leader=False, failover=None): l = Leader(0, 0, m1) if leader else None m2 = Member(0, 'other', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5436/postgres', 'api_url': 'http://127.0.0.1:8011/patroni', + 'state': 'running', 'tags': {'clonefrom': True}, 'scheduled_restart': {'schedule': "2100-01-01 10:53:07.560445+00:00", 'postgres_version': '99.0.0'}})