From 33ff372ef6b1e56410f008a463329da1b32482be Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 1 Sep 2016 11:08:26 +0200 Subject: [PATCH] Always try to rewind on manual failover --- features/patroni_api.feature | 6 +++--- patroni/__init__.py | 4 ++-- patroni/dcs/__init__.py | 3 +++ patroni/ha.py | 32 +++++++++++++++++++++----------- patroni/postgresql.py | 35 +++++++++++++++++------------------ tests/test_ha.py | 4 ++++ tests/test_postgresql.py | 8 ++------ 7 files changed, 52 insertions(+), 40 deletions(-) diff --git a/features/patroni_api.feature b/features/patroni_api.feature index e201613c..67b886e7 100644 --- a/features/patroni_api.feature +++ b/features/patroni_api.feature @@ -34,13 +34,13 @@ Scenario: check local configuration reload Then I receive a response code 202 Scenario: check dynamic configuration change via DCS - Given I issue a PATCH request to http://127.0.0.1:8008/config with {"ttl": 20, "loop_wait": 1, "postgresql": {"parameters": {"max_connections": 101}}} + Given I issue a PATCH request to http://127.0.0.1:8008/config with {"ttl": 20, "loop_wait": 2, "postgresql": {"parameters": {"max_connections": 101}}} Then I receive a response code 200 - And I receive a response loop_wait 1 + And I receive a response loop_wait 2 And Response on GET http://127.0.0.1:8008/patroni contains pending_restart after 11 seconds When I issue a GET request to http://127.0.0.1:8008/config Then I receive a response code 200 - And I receive a response loop_wait 1 + And I receive a response loop_wait 2 When I issue a GET request to http://127.0.0.1:8008/patroni Then I receive a response code 200 And I receive a response tags {'tag': 'new_value'} diff --git a/patroni/__init__.py b/patroni/__init__.py index d713969b..7d7a38c2 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -53,7 +53,7 @@ class Patroni(object): @property def nofailover(self): - return self.tags.get('nofailover', False) + return bool(self.tags.get('nofailover', False)) def reload_config(self): try: @@ -78,7 +78,7 @@ class Patroni(object): @property def noloadbalance(self): - return self.tags.get('noloadbalance', False) + return bool(self.tags.get('noloadbalance', False)) def schedule_next_run(self): self.next_run += self.dcs.loop_wait diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index 0c403a1c..6b8aaaec 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -186,6 +186,9 @@ class Failover(namedtuple('Failover', 'index,leader,candidate,scheduled_at')): return Failover(index, data.get('leader'), data.get('member'), data.get('scheduled_at')) + def __len__(self): + return int(bool(self.leader)) + int(bool(self.candidate)) + class ClusterConfig(namedtuple('ClusterConfig', 'index,data,modify_index')): diff --git a/patroni/ha.py b/patroni/ha.py index cc9bca0f..d974951b 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -134,7 +134,7 @@ class Ha(object): return node_to_follow if node_to_follow and node_to_follow.name != self.state_handler.name else None - def follow(self, demote_reason, follow_reason, refresh=True, recovery=False): + def follow(self, demote_reason, follow_reason, refresh=True, recovery=False, need_rewind=None): if refresh: self.load_cluster_from_dcs() @@ -146,14 +146,14 @@ class Ha(object): node_to_follow = self._get_node_to_follow(self.cluster) - if self.is_paused(): + if self.is_paused() and not self.state_handler.need_rewind: self.state_handler.set_role('master' if is_leader else 'replica') if is_leader: return 'continue to run as master without lock' elif not node_to_follow: return 'no action' - self.state_handler.follow(node_to_follow, self.cluster.leader, recovery, self._async_executor) + self.state_handler.follow(node_to_follow, self.cluster.leader, recovery, self._async_executor, need_rewind) return ret @@ -307,14 +307,14 @@ class Ha(object): def demote(self, delete_leader=True): if delete_leader: self.state_handler.stop() - self.state_handler.set_role('unknown') + 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 promote + sleep(2) # Give a time to somebody to take the leader lock cluster = self.dcs.get_cluster() node_to_follow = self._get_node_to_follow(cluster) - self.state_handler.follow(node_to_follow, cluster.leader, True) + self.state_handler.follow(node_to_follow, cluster.leader, recovery=True, need_rewind=True) else: self.state_handler.follow(None, None) @@ -399,11 +399,18 @@ class Ha(object): return self.follow('demoted self after trying and failing to obtain lock', 'following new leader after trying and failing to obtain lock') else: + # when we are doing manual failover there is no guaranty that new leader is ahead of any other node + need_rewind = bool(self.cluster.failover) or self.patroni.nofailover + if need_rewind: + sleep(2) # Give a time to somebody to take the leader lock + if self.patroni.nofailover: return self.follow('demoting self because I am not allowed to become master', - 'following a different leader because I am not allowed to promote') - return self.follow('demoting self because i am not the healthiest node', # should not happen in real life - 'following a different leader because i am not the healthiest node') + 'following a different leader because I am not allowed to promote', + need_rewind=need_rewind) + return self.follow('demoting self because i am not the healthiest node', + 'following a different leader because i am not the healthiest node', + need_rewind=need_rewind) def process_healthy_cluster(self): if self.has_lock(): @@ -413,6 +420,9 @@ class Ha(object): return msg if self.is_paused() and not self.state_handler.is_leader(): + if self.cluster.failover and self.cluster.failover.candidate == self.state_handler.name: + return 'waiting to become master after promote...' + self.dcs.delete_leader() self.dcs.reset_cluster() return 'removed leader lock because postgres is not running as master' @@ -585,7 +595,7 @@ class Ha(object): return self.handle_long_action_in_progress() # we've got here, so any async action has finished. Check if we tried to recover and failed - if self.recovering: + if self.recovering and not self.state_handler.need_rewind: self.recovering = False msg = self.post_recover() if msg is not None: @@ -610,7 +620,7 @@ class Ha(object): self.dcs.delete_leader() self.dcs.reset_cluster() return 'removed leader lock because postgres is not running' - else: + elif not self.state_handler.need_rewind: return 'postgres is not running' # try to start dead postgres diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 84960b67..99d8cde0 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -729,7 +729,14 @@ class Postgresql(object): except OSError: logger.exception("Unable to list %s", status_dir) - def follow(self, member, leader, recovery=False, async_executor=None): + @property + def need_rewind(self): + return self._need_rewind + + def follow(self, member, leader, recovery=False, async_executor=None, need_rewind=None): + if need_rewind is not None: + self._need_rewind = need_rewind + primary_conninfo = self.primary_conninfo(member) if self.check_recovery_conf(primary_conninfo) and not recovery: @@ -742,31 +749,23 @@ class Postgresql(object): self._do_follow(primary_conninfo, leader, recovery) def _do_follow(self, primary_conninfo, leader, recovery=False): - change_role = self.role == 'master' + change_role = self.role in ('master', 'demoted') - if change_role: - if leader: - if leader.name == self.name: - self._need_rewind = False - primary_conninfo = None - if self.is_running(): - return - else: - self._need_rewind = bool(leader.conn_url) and self.can_rewind - else: - self._need_rewind = False - primary_conninfo = None + if leader and leader.name == self.name: + primary_conninfo = None + self._need_rewind = False + if self.is_running(): + return + + self._need_rewind &= bool(leader and leader.conn_url) and self.can_rewind if self._need_rewind: - logger.info("set the rewind flag after demote") + logger.info("rewind flag is set") self.set_role('unknown') if self.is_running() and not self.stop(): return logger.warning('Can not run pg_rewind because postgres is still running') - if not (leader and leader.conn_url): - return logger.info('Leader unknown, can not rewind') - # prepare pg_rewind connection r = leader.conn_kwargs(self._superuser) diff --git a/tests/test_ha.py b/tests/test_ha.py index 754a5788..3cf29d9f 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -364,6 +364,7 @@ class TestHa(unittest.TestCase): self.assertEquals('PAUSE: no action. i am the leader with the lock', self.ha.run_cycle()) @patch('requests.get', requests_get) + @patch('time.sleep', Mock()) def test_manual_failover_process_no_leader(self): self.p.is_leader = false self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', self.p.name, None)) @@ -388,6 +389,7 @@ class TestHa(unittest.TestCase): self.ha.patroni.nofailover = True self.assertEquals(self.ha.run_cycle(), 'following a different leader because I am not allowed to promote') + @patch('time.sleep', Mock()) def test_manual_failover_process_no_leader_in_pause(self): self.ha.is_paused = true self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'other', None)) @@ -491,6 +493,8 @@ class TestHa(unittest.TestCase): self.p.name = 'leader' self.ha.cluster = get_cluster_initialized_with_leader() self.assertEquals(self.ha.run_cycle(), 'PAUSE: removed leader lock because postgres is not running as master') + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', self.p.name, None)) + self.assertEquals(self.ha.run_cycle(), 'PAUSE: waiting to become master after promote...') def test_postgres_unhealthy_in_pause(self): self.ha.is_paused = true diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 2ab29124..9fb7bb7a 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -263,20 +263,16 @@ class TestPostgresql(unittest.TestCase): with patch.object(Postgresql, 'restart', Mock(return_value=False)): self.p.set_role('replica') self.p.follow(None, None) # restart without rewind - self.p.set_role('master') with patch.object(Postgresql, 'stop', Mock(return_value=False)): - self.p.follow(self.leader, self.leader) # failed to stop postgres - - self.p.follow(self.leader, None) # Leader unknown, can not rewind + self.p.follow(self.leader, self.leader, need_rewind=True) # failed to stop postgres self.p.follow(self.leader, self.leader) # "leader" is not accessible or is_in_recovery with patch.object(Postgresql, 'checkpoint', Mock(return_value=None)): self.p.follow(self.leader, self.leader) - self.p.set_role('master') mock_pg_rewind.return_value = True - self.p.follow(self.leader, self.leader) + self.p.follow(self.leader, self.leader, need_rewind=True) self.p.follow(None, None) # check_recovery_conf...