From cbff544b9c4877f3372d89288fbc4f225e1024d3 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 25 Jun 2020 16:27:57 +0200 Subject: [PATCH] Implement patronictl flush switchover (#1554) It includes implementing the `DELETE /switchover` REST API endpoint. Close https://github.com/zalando/patroni/issues/1376 --- docs/rest_api.rst | 7 +++++-- patroni/api.py | 14 ++++++++++++++ patroni/ctl.py | 48 +++++++++++++++++++++++++++++++++++------------ tests/test_api.py | 9 +++++++++ tests/test_ctl.py | 22 +++++++++++++++++++++- 5 files changed, 85 insertions(+), 15 deletions(-) diff --git a/docs/rest_api.rst b/docs/rest_api.rst index 93d5b52b..f0f9b6ab 100644 --- a/docs/rest_api.rst +++ b/docs/rest_api.rst @@ -299,7 +299,10 @@ Example: schedule a switchover from the leader to any other healthy replica in t Depending on the situation the request might finish with a different HTTP status code and body. The status code **200** is returned when the switchover or failover successfully completed. If the switchover was successfully scheduled, Patroni will return HTTP status code **202**. In case something went wrong, the error status code (one of **400**, **412** or **503**) will be returned with some details in the response body. For more information please check the source code of ``patroni/api.py:do_POST_failover()`` method. -The switchover and failover endpoints are used by ``patronictl switchover`` and ``patronictl failover``, respectively. +- ``DELETE /switchover``: delete the scheduled switchover + +The ``POST /switchover`` and ``POST failover`` endpoints are used by ``patronictl switchover`` and ``patronictl failover``, respectively. +The ``DELETE /switchover`` is used by ``patronictl flush switchover``. Restart endpoint @@ -315,7 +318,7 @@ Restart endpoint - ``DELETE /restart``: delete the scheduled restart -``POST /restart`` and ``DELETE /restart`` endpoints are used by ``patronictl restart`` and ``patronictl flush`` respectively. +``POST /restart`` and ``DELETE /restart`` endpoints are used by ``patronictl restart`` and ``patronictl flush restart`` respectively. Reload endpoint diff --git a/patroni/api.py b/patroni/api.py index 4616595e..23c4e6bf 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -292,6 +292,20 @@ class RestApiHandler(BaseHTTPRequestHandler): code = 404 self._write_response(code, data) + @check_auth + def do_DELETE_switchover(self): + failover = self.server.patroni.dcs.get_cluster().failover + if failover and failover.scheduled_at: + if not self.server.patroni.dcs.manual_failover('', '', index=failover.index): + return self.send_error(409) + else: + data = "scheduled switchover deleted" + code = 200 + else: + data = "no switchover is scheduled" + code = 404 + self._write_response(code, data) + @check_auth def do_POST_reinitialize(self): request = self._read_json_content(body_is_optional=True) diff --git a/patroni/ctl.py b/patroni/ctl.py index efdda0fe..94df4a5b 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -242,6 +242,15 @@ def get_any_member(cluster, role='master', member=None): return m +def get_all_members_leader_first(cluster): + leader_name = cluster.leader.member.name if cluster.leader and cluster.leader.member.api_url else None + if leader_name: + yield cluster.leader.member + for member in cluster.members: + if member.api_url and member.name != leader_name: + yield member + + def get_cursor(cluster, connect_parameters, role='master', member=None): member = get_any_member(cluster, role=role, member=member) if member is None: @@ -935,25 +944,45 @@ def scaffold(obj, cluster_name, sysid): click.echo("Cluster {0} has been created successfully".format(cluster_name)) -@ctl.command('flush', help='Discard scheduled events (restarts only currently)') +@ctl.command('flush', help='Discard scheduled events') @click.argument('cluster_name') @click.argument('member_names', nargs=-1) -@click.argument('target', type=click.Choice(['restart'])) +@click.argument('target', type=click.Choice(['restart', 'switchover'])) @click.option('--role', '-r', help='Flush only members with this role', default='any', type=click.Choice(['master', 'replica', 'any'])) @option_force @click.pass_obj def flush(obj, cluster_name, member_names, force, role, target): - cluster = get_dcs(obj, cluster_name).get_cluster() + dcs = get_dcs(obj, cluster_name) + cluster = dcs.get_cluster() - members = get_members(cluster, cluster_name, member_names, role, force, 'flush') - for member in members: - if target == 'restart': + if target == 'restart': + for member in get_members(cluster, cluster_name, member_names, role, force, 'flush'): if member.data.get('scheduled_restart'): r = request_patroni(member, 'delete', 'restart') check_response(r, member.name, 'flush scheduled restart') else: click.echo('No scheduled restart for member {0}'.format(member.name)) + elif target == 'switchover': + failover = cluster.failover + if not failover or not failover.scheduled_at: + return click.echo('No pending scheduled switchover') + for member in get_all_members_leader_first(cluster): + try: + r = request_patroni(member, 'delete', 'switchover') + if r.status in (200, 404): + prefix = 'Success' if r.status == 200 else 'Failed' + return click.echo('{0}: {1}'.format(prefix, r.data.decode('utf-8'))) + except Exception as err: + logging.warning(str(err)) + logging.warning('Member %s is not accessible', member.name) + + click.echo('Failed: member={0}, status_code={1}, ({2})'.format( + member.name, r.status, r.data.decode('utf-8'))) + + logging.warning('Failing over to DCS') + click.echo('{0} Could not find any accessible member of cluster {1}'.format(timestamp(), cluster_name)) + dcs.manual_failover('', '', index=failover.index) def wait_until_pause_is_applied(dcs, paused, old_cluster): @@ -980,12 +1009,7 @@ def toggle_pause(config, cluster_name, paused, wait): if cluster.is_paused() == paused: raise PatroniCtlException('Cluster is {0} paused'.format(paused and 'already' or 'not')) - members = [] - if cluster.leader and cluster.leader.member.api_url: - members.append(cluster.leader.member) - members.extend([m for m in cluster.members if m.api_url and (not members or members[0].name != m.name)]) - - for member in members: + for member in get_all_members_leader_first(cluster): try: r = request_patroni(member, 'patch', 'config', {'pause': paused or None}) except Exception as err: diff --git a/tests/test_api.py b/tests/test_api.py index 92ea89c2..32de0869 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -311,6 +311,15 @@ class TestRestApiHandler(unittest.TestCase): request = 'DELETE /restart HTTP/1.0' + self._authorization self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) + @patch.object(MockPatroni, 'dcs') + def test_do_DELETE_switchover(self, mock_dcs): + request = 'DELETE /switchover HTTP/1.0' + self._authorization + self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) + mock_dcs.manual_failover.return_value = False + self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) + mock_dcs.get_cluster.return_value.failover = None + self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) + @patch.object(MockPatroni, 'dcs') def test_do_POST_reinitialize(self, mock_dcs): cluster = mock_dcs.get_cluster.return_value diff --git a/tests/test_ctl.py b/tests/test_ctl.py index e2cb11de..02883af3 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -449,7 +449,7 @@ class TestCtl(unittest.TestCase): @patch('patroni.ctl.get_dcs') @patch.object(PoolManager, 'request', Mock(return_value=MockResponse())) - def test_flush(self, mock_get_dcs): + def test_flush_restart(self, mock_get_dcs): mock_get_dcs.return_value = self.e mock_get_dcs.return_value.get_cluster = get_cluster_initialized_with_leader @@ -462,6 +462,26 @@ class TestCtl(unittest.TestCase): result = self.runner.invoke(ctl, ['flush', 'dummy', 'restart', '--force']) assert 'Failed: flush scheduled restart' in result.output + @patch('patroni.ctl.get_dcs') + @patch.object(PoolManager, 'request', Mock(return_value=MockResponse())) + def test_flush_switchover(self, mock_get_dcs): + mock_get_dcs.return_value = self.e + + mock_get_dcs.return_value.get_cluster = get_cluster_initialized_with_leader + result = self.runner.invoke(ctl, ['flush', 'dummy', 'switchover']) + assert 'No pending scheduled switchover' in result.output + + scheduled_at = datetime.now(tzutc) + timedelta(seconds=600) + mock_get_dcs.return_value.get_cluster = Mock( + return_value=get_cluster_initialized_with_leader(Failover(1, 'a', 'b', scheduled_at))) + result = self.runner.invoke(ctl, ['flush', 'dummy', 'switchover']) + assert result.output.startswith('Success: ') + + mock_get_dcs.return_value.manual_failover = Mock() + with patch.object(PoolManager, 'request', side_effect=[MockResponse(409), Exception]): + result = self.runner.invoke(ctl, ['flush', 'dummy', 'switchover']) + assert 'Could not find any accessible member of cluster' in result.output + @patch.object(PoolManager, 'request') @patch('patroni.ctl.get_dcs') @patch('patroni.ctl.polling_loop', Mock(return_value=[1]))