From 37315903fa0f84c7d26380ad86f4c681f2b007af Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Wed, 10 Feb 2016 13:52:39 +0100 Subject: [PATCH 1/9] Implement scheduled failover. Scheduled failover allows scheduling of a failover in the future. It does this by writing a failover key in the DCS which contains the scheduled failover time. The reason to allow a scheduled failover, is that it does not require one to use a scheduler (e.g. cron) to schedule such a failover. One of the issues with using a scheduler is that it may need to authenticate itself. With scheduled failover the authentication takes place during the scheduling, not during the actual failover. To allow the time of failover to be expressed, the failover key has changed its format; the old format however can still be used. The new format expects the failover key to be a json-document with relevant keys set. We need the timestamp specified to be time zone aware and to be expressed unambigiously, e.g. ISO 8601. --- patroni/api.py | 35 +++++++++++++++++++++++++------- patroni/ctl.py | 53 +++++++++++++++++++++++++++++------------------- patroni/dcs.py | 53 +++++++++++++++++++++++++++++++++++++++++++----- patroni/ha.py | 25 +++++++++++++++++++++++ patroni/utils.py | 33 ++++++++---------------------- 5 files changed, 142 insertions(+), 57 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index 2318fbf0..677c8f85 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -5,6 +5,9 @@ import logging import psycopg2 import socket import time +import dateutil +import datetime +import pytz from patroni.exceptions import PostgresConnectionException from patroni.utils import Retry, RetryFailedError @@ -176,13 +179,31 @@ class RestApiHandler(BaseHTTPRequestHandler): member = request.get('member', None) cluster = self.server.patroni.ha.dcs.get_cluster() status_code = 503 - data = self.is_failover_possible(cluster, leader, member) - if not data: - if not self.server.patroni.dcs.manual_failover(leader, member): - data = b'failed to write failover key into DCS' - else: - self.server.patroni.dcs.event.set() - status_code, data = self.poll_failover_result(cluster.leader and cluster.leader.name, member) + + data = b'' + if request.get('scheduled_at'): + try: + scheduled_at = dateutil.parser.parse(request['scheduled_at']) + if scheduled_at.tzinfo is None: + data = b'Timezone information is mandatory for scheduled_at' + status_code = 400 + elif scheduled_at < datetime.datetime.now(pytz.utc): + data = b'Cannot schedule failover in the past' + status_code = 422 + elif self.server.patroni.dcs.manual_failover(leader, member, scheduled_at): + data = b'Failover scheduled' + status_code = 200 + except (ValueError, TypeError): + logger.exception('Invalid scheduled failover time: {}'.format(request['scheduled_at'])) + data = b'Unable to parse scheduled timestamp. It should be in an unambiguous format, e.g. ISO 8601' + else: + data = self.is_failover_possible(cluster, leader, member) + if not data: + if not self.server.patroni.dcs.manual_failover(leader, member): + data = b'failed to write failover key into DCS' + else: + self.server.patroni.dcs.event.set() + status_code, data = self.poll_failover_result(cluster.leader and cluster.leader.name, member) self.send_response(status_code) self.send_header('Content-Type', 'text/html') diff --git a/patroni/ctl.py b/patroni/ctl.py index 635b0a89..1d923420 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -14,6 +14,8 @@ import datetime from prettytable import PrettyTable from six.moves.urllib_parse import urlparse import logging +import dateutil +import tzlocal from .etcd import Etcd from .exceptions import PatroniCtlException @@ -473,10 +475,12 @@ def reinit(cluster_name, member_names, config_file, dcs, force): @click.argument('cluster_name') @click.option('--master', help='The name of the current master', default=None) @click.option('--candidate', help='The name of the candidate', default=None) +@click.option('--scheduled', help='Timestamp of a scheduled failover in unambiguous format (e.g. ISO 8601)', + default=None) @click.option('--force', is_flag=True) @option_config_file @option_dcs -def failover(config_file, cluster_name, master, candidate, force, dcs): +def failover(config_file, cluster_name, master, candidate, force, dcs, scheduled): """ We want to trigger a failover for the specified cluster name. @@ -514,6 +518,25 @@ def failover(config_file, cluster_name, master, candidate, force, dcs): if candidate and candidate not in candidate_names: raise PatroniCtlException('Member {} does not exist in cluster {}'.format(candidate, cluster_name)) + if scheduled is None and not force: + scheduled = click.prompt('When should the failover take place (e.g. 2015-10-01T14:30) ', type=str, + default='now') + + if (scheduled or 'now') == 'now': + scheduled_at = None + else: + try: + scheduled_at = dateutil.parser.parse(scheduled) + if scheduled_at.tzinfo is None: + scheduled_at = tzlocal.get_localzone().localize(scheduled_at) + except (ValueError, TypeError): + message = 'Unable to parse scheduled timestamp ({}). It should be in an unambiguous format (e.g. ISO 8601)' + raise PatroniCtlException(message.format(scheduled)) + scheduled_at = scheduled_at.isoformat() + + failover_value = {'leader': master, 'member': candidate, 'scheduled_at': scheduled_at} + logging.debug(failover_value) + # By now we have established that the leader exists and the candidate exists click.echo('Current cluster topology') output_members(dcs.get_cluster(), name=cluster_name) @@ -525,17 +548,14 @@ def failover(config_file, cluster_name, master, candidate, force, dcs): if not a: raise PatroniCtlException('Aborting failover') - failover_value = '{}:{}'.format(master, candidate or '') - - t_started = time.time() r = None try: - r = post_patroni(cluster.leader.member, 'failover', {'leader': master, 'member': candidate or ''}) + r = post_patroni(cluster.leader.member, 'failover', failover_value) if r.status_code == 200: logging.debug(r) - logging.debug(r.text) cluster = dcs.get_cluster() - click.echo(timestamp() + ' Failing over to new leader: {}'.format(cluster.leader.member.name)) + logging.debug(cluster) + click.echo('{} {}'.format(timestamp(), r.text)) else: click.echo('Failover failed, details: {}, {}'.format(r.status_code, r.text)) return @@ -543,17 +563,9 @@ def failover(config_file, cluster_name, master, candidate, force, dcs): logging.exception(r) logging.warning('Failing over to DCS') click.echo(timestamp() + ' Could not failover using Patroni api, falling back to DCS') - dcs.set_failover_value(failover_value) - click.echo(timestamp() + ' Initialized failover from master {}'.format(master)) - # The failover process should within a minute update the failover key, we will keep watching it until it changes - # or we timeout - cluster = wait_for_leader(dcs, timeout=60) - if cluster.leader.member.name == master: - click.echo('Failover failed, master did not change after {:0.1f} seconds'.format(time.time() - t_started)) - return + click.echo(timestamp() + ' Initializing failover from master {}'.format(master)) + dcs.manual_failover(leader=master, member=candidate, scheduled_at=failover_value) - click.echo(timestamp() + ' Failover completed in {:0.1f} seconds, new leader is {}'.format(time.time() - t_started, - str(cluster.leader.member.name))) output_members(cluster, name=cluster_name) @@ -577,10 +589,9 @@ def output_members(cluster, name=None, format='pretty'): host = build_connect_parameters(m.conn_url)['host'] - xlog_location = m.data.get('xlog_location') - if xlog_location is None or (xlog_location_cluster < xlog_location): - lag = '' - else: + xlog_location = m.data.get('xlog_location') or 0 + lag = '' + if (xlog_location_cluster >= xlog_location): lag = round((xlog_location_cluster - xlog_location)/1024/1024) rows.append([ diff --git a/patroni/dcs.py b/patroni/dcs.py index a6b9d856..c689c4a9 100644 --- a/patroni/dcs.py +++ b/patroni/dcs.py @@ -1,5 +1,6 @@ import abc import json +import dateutil from collections import namedtuple from patroni.exceptions import DCSError @@ -89,12 +90,44 @@ class Leader(namedtuple('Leader', 'index,session,member')): return self.member.conn_url -class Failover(namedtuple('Failover', 'index,leader,member')): +class Failover(namedtuple('Failover', 'index,leader,member,scheduled_at')): + """ + >>> 'Failover' in str(Failover.from_node(1, '{"leader": "cluster_leader"}')) + True + >>> 'Failover' in str(Failover.from_node(1, '{"leader": "cluster_leader", "member": "cluster:member"}')) + True + >>> Failover.from_node(1, 'null') is None + True + >>> n = '{"leader": "cluster_leader", "member": "cluster:member", "scheduled_at": "2016-01-14T10:09:57.1394Z"}' + >>> 'tzinfo=' in str(Failover.from_node(1, n)) + True + >>> Failover.from_node(1, None) is None + True + >>> Failover.from_node(1, '{}') is None + True + >>> 'abc' in Failover.from_node(1, 'abc:def') + True + """ @staticmethod def from_node(index, value): - t = [a.strip() for a in value.split(':')] + [''] - return Failover(index, t[0], t[1]) if t[0] or t[1] else None + if not value: + return None + + try: + data = json.loads(value) + if not data: + return None + except ValueError: + t = [a.strip() for a in value.split(':')] + leader = t[0] + candidate = t[1] if len(t) > 1 else None + return Failover(index, leader, candidate, None) if leader or candidate else None + + if data.get('scheduled_at'): + data['scheduled_at'] = dateutil.parser.parse(data['scheduled_at']) + + return Failover(index, data.get('leader'), data.get('member'), data.get('scheduled_at')) class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members,failover')): @@ -223,8 +256,18 @@ class AbstractDCS: def set_failover_value(self, value, index=None): """Create or update `/failover` key""" - def manual_failover(self, leader, member, index=None): - return self.set_failover_value(leader + (':' + member if member else ''), index) + def manual_failover(self, leader, member, scheduled_at=None, index=None): + failover_value = dict() + if leader: + failover_value['leader'] = leader + + if member: + failover_value['member'] = member + + if scheduled_at: + failover_value['scheduled_at'] = scheduled_at.isoformat() + + return self.set_failover_value(json.dumps(failover_value), index) def current_leader(self): try: diff --git a/patroni/ha.py b/patroni/ha.py index 7e78eb77..dc038da9 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -3,6 +3,9 @@ import logging import psycopg2 import requests import sys +import time +import datetime +import pytz from patroni.async_executor import AsyncExecutor from patroni.exceptions import DCSError, PostgresConnectionException @@ -268,6 +271,28 @@ class Ha: def process_manual_failover_from_leader(self): failover = self.cluster.failover + + if failover.scheduled_at: + # If the failover is in the far future, we shouldn't do anything and just return. + # If the failover is in the past, we consider the value to be stale and we remove + # the value. + # If the value is close to now, we initiate the failover + now = datetime.datetime.now(pytz.utc) + delta = (failover.scheduled_at - now).total_seconds() + + if delta > 10: + logging.info('Awaiting failover at {0} (in {1:.0f} seconds)'.format(failover.scheduled_at.isoformat(), + delta)) + return + elif delta < -15: + logger.warning('Found a stale failover value, cleaning up: {}'.format(failover.scheduled_at)) + self.dcs.manual_failover('', '', self.cluster.failover.index) + return + + # The value is very close to now + time.sleep(max(delta, 0)) + logger.info('Manual scheduled failover at {}'.format(failover.scheduled_at.isoformat())) + if not failover.leader or failover.leader == self.state_handler.name: if not failover.member or failover.member != self.state_handler.name: members = [m for m in self.cluster.members if not failover.member or m.name == failover.member] diff --git a/patroni/utils.py b/patroni/utils.py index 9b040294..7e827222 100644 --- a/patroni/utils.py +++ b/patroni/utils.py @@ -1,10 +1,11 @@ import datetime import os import random -import re import signal import sys import time +import pytz +import dateutil.parser from patroni.exceptions import PatroniException @@ -12,39 +13,23 @@ ignore_sigterm = False interrupted_sleep = False reap_children = False -_DATE_TIME_RE = re.compile(r'''^ -(?P\d{4})\-(?P\d{2})\-(?P\d{2}) # date -T -(?P\d{2}):(?P\d{2}):(?P\d{2})\.(?P\d{6}) # time -\d*Z$''', re.X) - - -def parse_datetime(time_str): - """ - >>> parse_datetime('2015-06-10T12:56:30.552539016Z') - datetime.datetime(2015, 6, 10, 12, 56, 30, 552539) - >>> parse_datetime('2015-06-10 12:56:30.552539016Z') - """ - m = _DATE_TIME_RE.match(time_str) - if not m: - return None - p = dict((n, int(m.group(n))) for n in 'year month day hour minute second microsecond'.split(' ')) - return datetime.datetime(**p) - def calculate_ttl(expiration): """ >>> calculate_ttl(None) - >>> calculate_ttl('2015-06-10 12:56:30.552539016Z') + >>> calculate_ttl('2015-06-10 12:56:30.552539016Z') < 0 + True >>> calculate_ttl('2015-06-10T12:56:30.552539016Z') < 0 True + >>> calculate_ttl('fail-06-10T12:56:30.552539016Z') """ if not expiration: return None - expiration = parse_datetime(expiration) - if not expiration: + try: + expiration = dateutil.parser.parse(expiration) + except (ValueError, TypeError): return None - now = datetime.datetime.utcnow() + now = datetime.datetime.now(pytz.utc) return int((expiration - now).total_seconds()) From 1e2fdac8919ea78224921759d81b82f54084b652 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Wed, 10 Feb 2016 14:19:41 +0100 Subject: [PATCH 2/9] Scheduled Failover tests Add tests for the scheduled failover feature, also add more and better tests for patronictl. --- tests/test_api.py | 20 +++++++ tests/test_ctl.py | 134 +++++++++++++++++++++++++++++---------------- tests/test_etcd.py | 2 + tests/test_ha.py | 43 ++++++++++++--- 4 files changed, 142 insertions(+), 57 deletions(-) diff --git a/tests/test_api.py b/tests/test_api.py index 4b24704d..a9e9398b 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -169,3 +169,23 @@ class TestRestApiHandler(unittest.TestCase): request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ b'Content-Length: 50\n\n{"leader": "postgresql1", "member": "postgresql2"}' MockRestApiServer(RestApiHandler, request) + + ## Valid future date + request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ + b'Content-Length: 103\n\n{"leader": "postgresql1", "member": "postgresql2", "scheduled_at": "6016-02-15T18:13:30.568224+01:00"}' + MockRestApiServer(RestApiHandler, request) + + ## Exception: No timezone specified + request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ + b'Content-Length: 97\n\n{"leader": "postgresql1", "member": "postgresql2", "scheduled_at": "6016-02-15T18:13:30.568224"}' + MockRestApiServer(RestApiHandler, request) + + ## Exception: Scheduled in the past + request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ + b'Content-Length: 103\n\n{"leader": "postgresql1", "member": "postgresql2", "scheduled_at": "1016-02-15T18:13:30.568224+01:00"}' + MockRestApiServer(RestApiHandler, request) + + ## Invalid date + request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ + b'Content-Length: 103\n\n{"leader": "postgresql1", "member": "postgresql2", "scheduled_at": "2010-02-29T18:13:30.568224+01:00"}' + MockRestApiServer(RestApiHandler, request) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index 81f825aa..985b526f 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -84,6 +84,7 @@ class TestCtl(unittest.TestCase): output_members(cluster, name='abc', format='json') output_members(cluster, name='abc', format='tsv') + @patch('patroni.etcd.Etcd.get_cluster', Mock(return_value=get_cluster_initialized_with_leader())) @patch('patroni.etcd.Etcd.get_etcd_client', Mock(return_value=None)) @patch('patroni.etcd.Etcd.set_failover_value', Mock(return_value=None)) @@ -97,73 +98,99 @@ class TestCtl(unittest.TestCase): with patch('patroni.etcd.Etcd.get_cluster', Mock(return_value=get_cluster_initialized_with_leader())): result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader other -y''') - assert 'Failing over to new leader' in result.output +y''') + assert 'leader' in result.output + + result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader +other +2100-01-01T12:23:00 +y''') + assert result.exit_code == 0 + + result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader +other +2030-01-01T12:23:00 +y''') + assert result.exit_code == 0 + + ## Aborting failover,as we anser NO to the confirmation + result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader +other +2030-01-01T12:23:00 +y''') + assert result.exit_code == 0 + + ## Aborting failover,as we anser NO to the confirmation result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader other N''') - assert 'Aborting failover' in str(result.output) + assert result.exit_code == 1 + ## Target and source are equal result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader leader -y''') - assert 'target and source are the same' in str(result.output) +y''') + assert result.exit_code == 1 + + ## Reality is not part of this cluster result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader Reality + y''') - assert 'Reality does not exist' in str(result.output) + assert result.exit_code == 1 result = runner.invoke(ctl, ['failover', 'dummy', '--force']) - assert 'Failing over to new leader' in result.output + assert 'Member' in result.output + result = runner.invoke(ctl, ['failover', 'dummy', '--force', '--scheduled', '2015-01-01T12:00:00+01:00']) + assert result.exit_code == 0 + + ## Invalid timestamp + result = runner.invoke(ctl, ['failover', 'dummy', '--force', '--scheduled', 'invalid']) + assert result.exit_code != 0 + + ## Invalid timestamp + result = runner.invoke(ctl, ['failover', 'dummy', '--force', '--scheduled', '2115-02-30T12:00:00+01:00']) + assert result.exit_code != 0 + + ## Specifying wrong leader result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='dummy') - assert 'is not the leader of cluster' in str(result.output) + assert result.exit_code == 1 with patch('patroni.etcd.Etcd.get_cluster', Mock(return_value=get_cluster_initialized_with_only_leader())): + ## No members available result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader other y''') - assert 'No candidates found to failover to' in str(result.output) + assert result.exit_code == 1 with patch('patroni.etcd.Etcd.get_cluster', Mock(return_value=get_cluster_initialized_without_leader())): + ## No master available result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader other + y''') - assert 'This cluster has no master' in str(result.output) + assert result.exit_code == 1 with patch('patroni.ctl.post_patroni', Mock(side_effect=Exception())): + ## Non-responding patroni result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader other + y''') assert 'falling back to DCS' in result.output - assert 'Failover failed' in result.output mocked = Mock() mocked.return_value.status_code = 500 with patch('patroni.ctl.post_patroni', Mock(return_value=mocked)): result = runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader other + y''') - assert 'Failover failed, details' in result.output + assert 'Failover failed' in result.output -# with patch('patroni.dcs.AbstractDCS.get_cluster', Mock(return_value=get_cluster_initialized_with_leader())): -# result = runner.invoke(ctl, ['failover', 'alpha', '--dcs', '8.8.8.8'], input='nonsense') -# assert 'is not the leader of cluster' in str(result.output) - - # result = runner.invoke(ctl, ['failover', 'alpha', '--dcs', '8.8.8.8', '--master', 'nonsense']) - # assert 'is not the leader of cluster' in str(result.output) - - # result = runner.invoke(ctl, ['failover', 'alpha', '--dcs', '8.8.8.8'], input='leader\nother\nn') - # assert 'Aborting failover' in str(result.output) - - # with patch('patroni.ctl.wait_for_leader', Mock(return_value = get_cluster_initialized_with_leader())): - # result = runner.invoke(ctl, ['failover', 'alpha', '--dcs', '8.8.8.8'], input='leader\nother\nY') - # assert 'master did not change after' in result.output - - # result = runner.invoke(ctl, ['failover', 'alpha', '--dcs', '8.8.8.8'], input='leader\nother\nY') - # assert 'Failover failed' in result.output def test_(self): self.assertRaises(patroni.exceptions.PatroniCtlException, get_dcs, {'scheme': 'dummy'}, 'dummy') @@ -174,6 +201,7 @@ y''') runner = CliRunner() with patch('patroni.ctl.get_dcs', Mock(return_value=self.e)): + ## Mutually exclusive result = runner.invoke(ctl, [ 'query', 'alpha', @@ -182,19 +210,14 @@ y''') '--role', 'master', ]) - assert 'mutually exclusive' in str(result.output) + assert result.exit_code == 1 with runner.isolated_filesystem(): dummy_file = open('dummy', 'w') dummy_file.write('SELECT 1') dummy_file.close() - result = runner.invoke(ctl, [ - 'query', - 'alpha' - ]) - assert 'You need to specify' in str(result.output) - + ## Mutually exclusive result = runner.invoke(ctl, [ 'query', 'alpha', @@ -203,7 +226,7 @@ y''') '--command', 'dummy', ]) - assert 'mutually exclusive' in str(result.output) + assert result.exit_code == 1 result = runner.invoke(ctl, ['query', 'alpha', '--file', 'dummy']) @@ -212,7 +235,12 @@ y''') result = runner.invoke(ctl, ['query', 'alpha', '--command', 'SELECT 1']) assert 'mock column' in result.output - result = runner.invoke(ctl, ['query', 'alpha', '--command', 'SELECT 1', '--dbname', 'dummy', '--password', '--username', 'dummy'], input='password\n') + ## --command or --file is mandatory + result = runner.invoke(ctl, ['query', 'alpha']) + assert result.exit_code == 1 + + result = runner.invoke(ctl, ['query', 'alpha', '--command', 'SELECT 1', + '--username', 'root', '--password', '--dbname', 'postgres'], input='ab\nab') assert 'mock column' in result.output @patch('patroni.ctl.get_cursor', Mock(return_value=MockConnect().cursor())) @@ -244,6 +272,7 @@ y''') result = runner.invoke(ctl, ['dsn', 'alpha', '--dcs', '8.8.8.8']) assert 'host=127.0.0.1 port=5435' in result.output + ## Mutually exclusive options result = runner.invoke(ctl, [ 'dsn', 'alpha', @@ -252,13 +281,11 @@ y''') '--member', 'dummy', ]) - assert 'mutually exclusive' in str(result.output) + assert result.exit_code == 1 + ## Non-existing member result = runner.invoke(ctl, ['dsn', 'alpha', '--member', 'dummy']) - assert 'Can not find' in str(result.output) - - # result = runner.invoke(ctl, ['dsn', 'alpha', '--dcs', '8.8.8.8', '--role', 'replica']) - # assert 'host=127.0.0.1 port=5436' in result.output + assert result.exit_code == 1 @patch('patroni.etcd.Etcd.get_cluster', Mock(return_value=get_cluster_initialized_with_leader())) @patch('patroni.etcd.Etcd.get_etcd_client', Mock(return_value=None)) @@ -268,9 +295,16 @@ y''') runner = CliRunner() result = runner.invoke(ctl, ['restart', 'alpha', '--dcs', '8.8.8.8'], input='y') - result = runner.invoke(ctl, ['reinit', 'alpha', '--dcs', '8.8.8.8'], input='y') + assert result.exit_code == 0 + result = runner.invoke(ctl, ['reinit', 'alpha', '--dcs', '8.8.8.8'], input='y') + assert result.exit_code == 1 + + # Aborted restart result = runner.invoke(ctl, ['restart', 'alpha', '--dcs', '8.8.8.8'], input='N') + assert result.exit_code == 1 + + ## Not a member result = runner.invoke(ctl, [ 'restart', 'alpha', @@ -279,7 +313,7 @@ y''') 'dummy', '--any', ], input='y') - assert 'not a member' in str(result.output) + assert result.exit_code == 1 with patch('requests.post', Mock(return_value=MockResponse())): result = runner.invoke(ctl, ['restart', 'alpha', '--dcs', '8.8.8.8'], input='y') @@ -292,15 +326,18 @@ y''') result = runner.invoke(ctl, ['remove', 'alpha', '--dcs', '8.8.8.8'], input='alpha\nslave') assert 'Please confirm' in result.output assert 'You are about to remove all' in result.output - assert 'You did not exactly type' in str(result.output) + ## Not typing an exact confirmation + assert result.exit_code == 1 + ## master specified does not match master of cluster result = runner.invoke(ctl, ['remove', 'alpha', '--dcs', '8.8.8.8'], input='''alpha Yes I am aware slave''') - assert 'You did not specify the current master of the cluster' in str(result.output) + assert result.exit_code == 1 + ## cluster specified on cmdline does not match verification prompt result = runner.invoke(ctl, ['remove', 'alpha', '--dcs', '8.8.8.8'], input='beta\nleader') - assert 'Cluster names specified do not match' in str(result.output) + assert result.exit_code == 1 with patch('patroni.etcd.Etcd.get_cluster', get_cluster_initialized_with_leader): result = runner.invoke(ctl, ['remove', 'alpha', '--dcs', '8.8.8.8'], @@ -310,11 +347,12 @@ leader''') assert 'object has no attribute' in str(result.exception) with patch('patroni.ctl.get_dcs', Mock(return_value=Mock())): + ## Not implemented DCS result = runner.invoke(ctl, ['remove', 'alpha', '--dcs', '8.8.8.8'], input='''alpha Yes I am aware leader''') - assert 'We have not implemented this for DCS of type' in str(result.output) + assert result.exit_code == 1 @patch('patroni.etcd.Etcd.watch', Mock(return_value=None)) @patch('patroni.etcd.Etcd.get_cluster', Mock(return_value=get_cluster_initialized_with_leader())) diff --git a/tests/test_etcd.py b/tests/test_etcd.py index 443dc2d8..33043936 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -60,6 +60,8 @@ def requests_get(url, **kwargs): response.content = '[{}]' else: response.content = members + elif url.endswith('/members'): + response.content = '{"action":"set","node":{"key":"/service/alpha/failover","value":"{\"leader\": \"f1410e163b6a\"}","modifiedIndex":257,"createdIndex":257},"prevNode":{"key":"/service/alpha/failover","value":"{\"scheduled_at\": \"2016-01-15T17:50:00+01:00\", \"leader\": \"f1410e163b6a\"}","modifiedIndex":241,"createdIndex":241}}' elif url.startswith('http://exhibitor'): response.content = '{"servers":["127.0.0.1","127.0.0.2","127.0.0.3"],"port":2181}' else: diff --git a/tests/test_ha.py b/tests/test_ha.py index 1ac164cf..5b4104b9 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -1,5 +1,7 @@ import etcd import unittest +import datetime +import pytz from mock import Mock, MagicMock, patch from patroni.dcs import Cluster, Failover, Leader, Member @@ -284,40 +286,63 @@ class TestHa(unittest.TestCase): self.ha.update_lock = false self.assertEquals(self.ha.run_cycle(), 'failed to update leader lock during restart') + @patch('requests.get', requests_get) def test_manual_failover_from_leader(self): self.ha.has_lock = true - self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', '')) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', '', None)) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') - self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', MockPostgresql.name)) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', MockPostgresql.name, None)) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') - self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', 'blabla')) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', 'blabla', None)) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') - f = Failover(0, MockPostgresql.name, '') + f = Failover(0, MockPostgresql.name, '', None) self.ha.cluster = get_cluster_initialized_with_leader(f) self.assertEquals(self.ha.run_cycle(), 'manual failover: demoting myself') self.ha.fetch_node_status = lambda e: (e, True, True, 0, {'nofailover': 'True'}) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') # manual failover from the previous leader to us won't happen if we hold the nofailover flag - self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name)) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, None)) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') + ## Failover scheduled time must include timezone + scheduled = datetime.datetime.now() + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, scheduled)) + + self.assertRaises(TypeError, self.ha.run_cycle) + + scheduled = datetime.datetime.utcnow().replace(tzinfo=pytz.UTC) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, scheduled)) + self.assertEquals('no action. i am the leader with the lock', self.ha.run_cycle()) + + scheduled = scheduled + datetime.timedelta(seconds=30) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, scheduled)) + self.assertEquals('no action. i am the leader with the lock', self.ha.run_cycle()) + + scheduled = scheduled + datetime.timedelta(seconds=-600) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, scheduled)) + self.assertEquals('no action. i am the leader with the lock', self.ha.run_cycle()) + + scheduled = None + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, scheduled)) + self.assertEquals('no action. i am the leader with the lock', self.ha.run_cycle()) + @patch('requests.get', requests_get) def test_manual_failover_process_no_leader(self): self.p.is_leader = false - self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', MockPostgresql.name)) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', MockPostgresql.name, None)) self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') - self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'leader')) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'leader', None)) self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') self.ha.fetch_node_status = lambda e: (e, True, True, 0, {}) # accessible, in_recovery self.assertEquals(self.ha.run_cycle(), 'following a different leader because i am not the healthiest node') - self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, MockPostgresql.name, '')) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, MockPostgresql.name, '', None)) self.assertEquals(self.ha.run_cycle(), 'following a different leader because i am not the healthiest node') self.ha.fetch_node_status = lambda e: (e, False, True, 0, {}) # inaccessible, in_recovery self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') # set failover flag to True for all members of the cluster # this should elect the current member, as we are not going to call the API for it. - self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'other')) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'other', None)) self.ha.fetch_node_status = lambda e: (e, True, True, 0, {'nofailover': 'True'}) # accessible, in_recovery self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') # same as previous, but set the current member to nofailover. In no case it should be elected as a leader From 854ad293c57d8daadd832a5f9f7e1bfa02cf1164 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Wed, 10 Feb 2016 14:26:25 +0100 Subject: [PATCH 3/9] Scheduled failover: Add requirements --- requirements-py2.txt | 2 ++ requirements-py3.txt | 2 ++ 2 files changed, 4 insertions(+) diff --git a/requirements-py2.txt b/requirements-py2.txt index 1e194bc0..26c0892d 100644 --- a/requirements-py2.txt +++ b/requirements-py2.txt @@ -9,3 +9,5 @@ kazoo>=2.2.1 python-etcd==0.4.2 click>=4.1 prettytable>=0.7 +tzlocal +python-dateutil diff --git a/requirements-py3.txt b/requirements-py3.txt index 13c3010c..a827efb3 100644 --- a/requirements-py3.txt +++ b/requirements-py3.txt @@ -9,3 +9,5 @@ kazoo>=2.2.1 python-etcd==0.4.2 click>=4.1 prettytable>=0.7 +tzlocal +python-dateutil From 0c2efeb7a7c1b571a654d907c7a8e74f4c19f5c2 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Thu, 11 Feb 2016 09:01:44 +0100 Subject: [PATCH 4/9] Change default http status code to 500. Instead of returning 503 (Service Unavailable) we no default to returning 500 (Internal Server Error). --- patroni/api.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index 677c8f85..9f715882 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -104,7 +104,7 @@ class RestApiHandler(BaseHTTPRequestHandler): @check_auth def do_POST_restart(self): - status_code = 503 + status_code = 500 data = b'restart failed' try: status, msg = self.server.patroni.ha.restart() @@ -178,7 +178,7 @@ class RestApiHandler(BaseHTTPRequestHandler): leader = request.get('leader', None) member = request.get('member', None) cluster = self.server.patroni.ha.dcs.get_cluster() - status_code = 503 + status_code = 500 data = b'' if request.get('scheduled_at'): @@ -196,11 +196,13 @@ class RestApiHandler(BaseHTTPRequestHandler): except (ValueError, TypeError): logger.exception('Invalid scheduled failover time: {}'.format(request['scheduled_at'])) data = b'Unable to parse scheduled timestamp. It should be in an unambiguous format, e.g. ISO 8601' + status_code = 422 else: data = self.is_failover_possible(cluster, leader, member) if not data: if not self.server.patroni.dcs.manual_failover(leader, member): data = b'failed to write failover key into DCS' + status_code = 503 else: self.server.patroni.dcs.event.set() status_code, data = self.poll_failover_result(cluster.leader and cluster.leader.name, member) From 1b14229da480e36ea36d662b1f1b92e5a46f8269 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 17 Feb 2016 12:18:50 +0100 Subject: [PATCH 5/9] Catch TypeError within ha loop not in the unit test In addition to that use sleep function from patroni.utils instead of time.sleep which is interruptable --- patroni/ha.py | 30 ++++++++++++++++-------------- tests/test_ha.py | 6 ++---- 2 files changed, 18 insertions(+), 18 deletions(-) diff --git a/patroni/ha.py b/patroni/ha.py index bf002a51..4aead262 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -3,13 +3,13 @@ import logging import psycopg2 import requests import sys -import time import datetime import pytz +from multiprocessing.pool import ThreadPool from patroni.async_executor import AsyncExecutor from patroni.exceptions import DCSError, PostgresConnectionException -from multiprocessing.pool import ThreadPool +from patroni.utils import sleep logger = logging.getLogger(__name__) @@ -283,20 +283,22 @@ class Ha(object): # the value. # If the value is close to now, we initiate the failover now = datetime.datetime.now(pytz.utc) - delta = (failover.scheduled_at - now).total_seconds() + try: + delta = (failover.scheduled_at - now).total_seconds() - if delta > 10: - logging.info('Awaiting failover at {0} (in {1:.0f} seconds)'.format(failover.scheduled_at.isoformat(), - delta)) - return - elif delta < -15: - logger.warning('Found a stale failover value, cleaning up: {}'.format(failover.scheduled_at)) - self.dcs.manual_failover('', '', self.cluster.failover.index) - return + if delta > 10: + logging.info('Awaiting failover at %s (in %.0f seconds)', failover.scheduled_at.isoformat(), delta) + return + elif delta < -15: + logger.warning('Found a stale failover value, cleaning up: %s', failover.scheduled_at) + self.dcs.manual_failover('', '', self.cluster.failover.index) + return - # The value is very close to now - time.sleep(max(delta, 0)) - logger.info('Manual scheduled failover at {}'.format(failover.scheduled_at.isoformat())) + # The value is very close to now + sleep(max(delta, 0)) + logger.info('Manual scheduled failover at {}'.format(failover.scheduled_at.isoformat())) + except TypeError: + logger.warning('Incorrect value in of scheduled_at: %s', failover.scheduled_at) if not failover.leader or failover.leader == self.state_handler.name: if not failover.member or failover.member != self.state_handler.name: diff --git a/tests/test_ha.py b/tests/test_ha.py index c6b98712..3589e952 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -297,7 +297,6 @@ class TestHa(unittest.TestCase): self.ha.update_lock = false self.assertEquals(self.ha.run_cycle(), 'failed to update leader lock during restart') - @patch('requests.get', requests_get) def test_manual_failover_from_leader(self): self.ha.has_lock = true @@ -316,11 +315,10 @@ class TestHa(unittest.TestCase): self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, None)) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') - ## Failover scheduled time must include timezone + # Failover scheduled time must include timezone scheduled = datetime.datetime.now() self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, scheduled)) - - self.assertRaises(TypeError, self.ha.run_cycle) + self.ha.run_cycle() scheduled = datetime.datetime.utcnow().replace(tzinfo=pytz.UTC) self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', MockPostgresql.name, scheduled)) From 1b9e77fe8320f375367a0b7eac22e70d689b8399 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 17 Feb 2016 12:34:04 +0100 Subject: [PATCH 6/9] pep8 formatting --- tests/test_api.py | 24 ++++++++++++------------ 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/tests/test_api.py b/tests/test_api.py index f06b0bb0..82a76f34 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -177,22 +177,22 @@ class TestRestApiHandler(unittest.TestCase): b'Content-Length: 50\n\n{"leader": "postgresql1", "member": "postgresql2"}' MockRestApiServer(RestApiHandler, request) - ## Valid future date - request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ - b'Content-Length: 103\n\n{"leader": "postgresql1", "member": "postgresql2", "scheduled_at": "6016-02-15T18:13:30.568224+01:00"}' + # Valid future date + request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\nContent-Length: 103\n\n{"leader": ' +\ + b'"postgresql1", "member": "postgresql2", "scheduled_at": "6016-02-15T18:13:30.568224+01:00"}' MockRestApiServer(RestApiHandler, request) - ## Exception: No timezone specified - request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ - b'Content-Length: 97\n\n{"leader": "postgresql1", "member": "postgresql2", "scheduled_at": "6016-02-15T18:13:30.568224"}' + # Exception: No timezone specified + request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\nContent-Length: 97\n\n{"leader": ' +\ + b'"postgresql1", "member": "postgresql2", "scheduled_at": "6016-02-15T18:13:30.568224"}' MockRestApiServer(RestApiHandler, request) - ## Exception: Scheduled in the past - request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ - b'Content-Length: 103\n\n{"leader": "postgresql1", "member": "postgresql2", "scheduled_at": "1016-02-15T18:13:30.568224+01:00"}' + # Exception: Scheduled in the past + request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\nContent-Length: 103\n\n{"leader": ' +\ + b'"postgresql1", "member": "postgresql2", "scheduled_at": "1016-02-15T18:13:30.568224+01:00"}' MockRestApiServer(RestApiHandler, request) - ## Invalid date - request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ - b'Content-Length: 103\n\n{"leader": "postgresql1", "member": "postgresql2", "scheduled_at": "2010-02-29T18:13:30.568224+01:00"}' + # Invalid date + request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\nContent-Length: 103\n\n{"leader": ' +\ + b'"postgresql1", "member": "postgresql2", "scheduled_at": "2010-02-29T18:13:30.568224+01:00"}' MockRestApiServer(RestApiHandler, request) From f079a9f308a7159c9c85c5d56ef3595835c854be Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 17 Feb 2016 12:34:22 +0100 Subject: [PATCH 7/9] remove unused code --- tests/test_etcd.py | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/tests/test_etcd.py b/tests/test_etcd.py index f81885dd..601a0cd6 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -58,12 +58,7 @@ def requests_get(url, **kwargs): elif ':8011/patroni' in url: response.content = '{"role": "replica", "xlog": {"replayed_location": 0}, "tags": {}}' elif url.endswith('/members'): - if url.startswith('http://error'): - response.content = '[{}]' - else: - response.content = members - elif url.endswith('/members'): - response.content = '{"action":"set","node":{"key":"/service/alpha/failover","value":"{\"leader\": \"f1410e163b6a\"}","modifiedIndex":257,"createdIndex":257},"prevNode":{"key":"/service/alpha/failover","value":"{\"scheduled_at\": \"2016-01-15T17:50:00+01:00\", \"leader\": \"f1410e163b6a\"}","modifiedIndex":241,"createdIndex":241}}' + response.content = '[{}]' if url.startswith('http://error') else members elif url.startswith('http://exhibitor'): response.content = '{"servers":["127.0.0.1","127.0.0.2","127.0.0.3"],"port":2181}' else: From 26e15862884933fced7694e240605d99e02da4f6 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Mon, 22 Feb 2016 12:20:20 +0100 Subject: [PATCH 8/9] Make the patronictl test provide an input for the schedule, even if it's empty. --- tests/test_ctl.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index 16814839..9ae7e3f0 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -118,6 +118,7 @@ y''') # Aborting failover,as we anser NO to the confirmation result = self.runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader other + N''') assert result.exit_code == 1 @@ -159,6 +160,7 @@ y''') # No members available result = self.runner.invoke(ctl, ['failover', 'dummy', '--dcs', '8.8.8.8'], input='''leader other + y''') assert result.exit_code == 1 From 287c0b312522e4dd25b613311e6c34cb61245d9d Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Mon, 22 Feb 2016 14:31:29 +0100 Subject: [PATCH 9/9] Fix the call to the function that was forgotten to be renamed. --- patroni/ha.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/patroni/ha.py b/patroni/ha.py index 4aead262..72414d89 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -272,7 +272,7 @@ class Ha(object): self.dcs.delete_leader() self.touch_member() self.dcs.reset_cluster() - self.state_handler.follow_the_leader(None) + self.state_handler.follow(None) def process_manual_failover_from_leader(self): failover = self.cluster.failover