From 6dc1d9c88eae163ece5292fa4b8696f99bf9e754 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 29 Aug 2016 15:37:20 +0200 Subject: [PATCH] Trigger reinitialize from api and make it possible to reinitialize in a pause state --- patroni/api.py | 25 +++++---------------- patroni/ha.py | 55 ++++++++++++++++------------------------------- tests/test_api.py | 11 ++-------- tests/test_ha.py | 25 ++++++++------------- 4 files changed, 35 insertions(+), 81 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index 7e066289..26c15cea 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -255,27 +255,12 @@ class RestApiHandler(BaseHTTPRequestHandler): @check_auth def do_POST_reinitialize(self): - patroni = self.server.patroni - cluster = patroni.dcs.get_cluster() - status_code = 500 - if cluster.is_paused(): - self._write_response(status_code, "Can't do reinitialize in the paused state") - return - - if cluster.is_unlocked(): - status_code = 503 - data = 'Cluster has no leader, can not reinitialize' - elif cluster.leader.name == patroni.ha.state_handler.name: - status_code = 503 - data = 'I am the leader, can not reinitialize' + data = self.server.patroni.ha.reinitialize() + if data is None: + status_code = 200 + data = 'reinitialize started' else: - action = patroni.ha.schedule_reinitialize() - if action is not None: - status_code = 503 - data = action + ' already in progress' - else: - status_code = 200 - data = 'reinitialize scheduled' + status_code = 503 self._write_response(status_code, data) def poll_failover_result(self, leader, candidate): diff --git a/patroni/ha.py b/patroni/ha.py index 75be52a9..7ef196f6 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -454,10 +454,6 @@ class Ha(object): logger.info("not proceeding with the restart: %s", reason_to_cancel) return False - def schedule(self, action, immediate=False): - with self._async_executor: - return self._async_executor.schedule(action, immediate) - def schedule_future_restart(self, restart_data): with self._async_executor: if not self.patroni.scheduled_restart: @@ -479,15 +475,6 @@ class Ha(object): return self.patroni.scheduled_restart.copy() if (self.patroni.scheduled_restart and isinstance(self.patroni.scheduled_restart, dict)) else None - def schedule_reinitialize(self): - return self.schedule('reinitialize') - - def reinitialize_scheduled(self): - return self._async_executor.scheduled_action == 'reinitialize' - - def schedule_restart(self, immediate=False): - return self.schedule('restart', immediate) - def restart_scheduled(self): return self._async_executor.scheduled_action == 'restart' @@ -500,7 +487,7 @@ class Ha(object): return (False, "restart conditions are not satisfied") with self._async_executor: - prev = self.schedule_restart(immediate=(not run_async)) + prev = self._async_executor.schedule('restart', not run_async) if prev is not None: return (False, prev + ' already in progress') if not run_async: @@ -512,28 +499,29 @@ class Ha(object): self._async_executor.run_async(self.state_handler.restart) return (True, "restart initiated") - def reinitialize(self, cluster): + def _do_reinitialize(self, cluster): self.state_handler.stop('immediate') self.state_handler.remove_data_directory() - clone_member = cluster.get_clone_member() - member_role = 'leader' if clone_member == cluster.leader else 'replica' + clone_member = self.cluster.get_clone_member() + member_role = 'leader' if clone_member == self.cluster.leader else 'replica' self.clone(clone_member, "from {0} '{1}'".format(member_role, clone_member.name)) - def process_scheduled_action(self): - if self.reinitialize_scheduled(): - if self.is_paused(): - logger.warning('Cluster is in a pause state, can not reinitialize') - self._async_executor.reset_scheduled_action() - elif self.cluster.is_unlocked(): - logger.error('Cluster has no leader, can not reinitialize') - self._async_executor.reset_scheduled_action() - elif self.has_lock(): - logger.error('I am the leader, can not reinitialize') - self._async_executor.reset_scheduled_action() - else: - self._async_executor.run_async(self.reinitialize, args=(self.cluster, )) - return 'reinitialize started' + def reinitialize(self): + with self._async_executor: + self.load_cluster_from_dcs() + + if self.cluster.is_unlocked(): + return 'Cluster has no leader, can not reinitialize' + + if self.cluster.leader.name == self.state_handler.name: + return 'I am the leader, can not reinitialize' + + action = self._async_executor.schedule('reinitialize', immediately=True) + if action is not None: + return '{0} already in progress'.format(action) + + self._async_executor.run_async(self._do_reinitialize, args=(self.cluster, )) def handle_long_action_in_progress(self): if self.has_lock(): @@ -585,11 +573,6 @@ class Ha(object): if msg is not None: return msg - # currently it can trigger only reinitialize - msg = self.process_scheduled_action() - if msg is not None: - return msg - # is data directory empty? if self.state_handler.data_directory_empty(): return self.bootstrap() # new node diff --git a/tests/test_api.py b/tests/test_api.py index ad96ed87..1c025b15 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -40,7 +40,7 @@ class MockHa(object): state_handler = MockPostgresql() @staticmethod - def schedule_reinitialize(): + def reinitialize(): return 'reinitialize' @staticmethod @@ -238,15 +238,8 @@ class TestRestApiHandler(unittest.TestCase): cluster.is_paused.return_value = False request = 'POST /reinitialize HTTP/1.0' + self._authorization MockRestApiServer(RestApiHandler, request) - cluster.is_unlocked.return_value = False - MockRestApiServer(RestApiHandler, request) - with patch.object(MockHa, 'schedule_reinitialize', Mock(return_value=None)): + with patch.object(MockHa, 'reinitialize', Mock(return_value=None)): MockRestApiServer(RestApiHandler, request) - cluster.leader.name = 'test' - self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) - - cluster.is_paused.return_value = True - self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) @patch('time.sleep', Mock()) def test_RestApiServer_query(self): diff --git a/tests/test_ha.py b/tests/test_ha.py index 69293cec..52b7af18 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -277,31 +277,24 @@ class TestHa(unittest.TestCase): self.assertRaises(PostgresException, self.ha.bootstrap) def test_reinitialize(self): - self.ha.schedule_reinitialize() - self.ha.schedule_reinitialize() - self.ha.run_cycle() + self.assertIsNotNone(self.ha.reinitialize()) self.assertIsNone(self.ha._async_executor.scheduled_action) - with patch.object(Ha, 'is_paused', true): - self.ha.schedule_reinitialize() - self.ha.run_cycle() - self.assertIsNone(self.ha._async_executor.scheduled_action) - self.ha.cluster = get_cluster_initialized_with_leader() - self.ha.has_lock = true - self.ha.schedule_reinitialize() - self.ha.run_cycle() - self.assertIsNone(self.ha._async_executor.scheduled_action) + self.assertIsNone(self.ha.reinitialize()) + self.assertIsNotNone(self.ha._async_executor.scheduled_action) - self.ha.has_lock = false - self.ha.schedule_reinitialize() - self.ha.run_cycle() + self.assertIsNotNone(self.ha.reinitialize()) + + self.ha.state_handler.name = self.ha.cluster.leader.name + self.assertIsNotNone(self.ha.reinitialize()) def test_restart(self): self.assertEquals(self.ha.restart(), (True, 'restarted successfully')) self.p.restart = false self.assertEquals(self.ha.restart(), (False, 'restart failed')) - self.ha.schedule_reinitialize() + self.ha.cluster = get_cluster_initialized_with_leader() + self.ha.reinitialize() self.assertEquals(self.ha.restart(), (False, 'reinitialize already in progress')) with patch.object(self.ha, "restart_matches", return_value=False): self.assertEquals(self.ha.restart({'foo': 'bar'}), (False, "restart conditions are not satisfied"))