Implement patronictl flush switchover (#1554)

It includes implementing the `DELETE /switchover` REST API endpoint.

Close https://github.com/zalando/patroni/issues/1376
This commit is contained in:
Alexander Kukushkin
2020-06-25 16:27:57 +02:00
committed by GitHub
parent 7f343c2c57
commit cbff544b9c
5 changed files with 85 additions and 15 deletions
+5 -2
View File
@@ -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 <cluster-name> 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 <cluster-name> restart`` respectively.
Reload endpoint
+14
View File
@@ -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)
+36 -12
View File
@@ -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:
+9
View File
@@ -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
+21 -1
View File
@@ -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]))