From 37315903fa0f84c7d26380ad86f4c681f2b007af Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Wed, 10 Feb 2016 13:52:39 +0100 Subject: [PATCH 01/21] 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 02/21] 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 03/21] 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 04/21] 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 05/21] 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 06/21] 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 07/21] 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 f7d60c61b6cd107c9e81fa9fd05fdb692067e64e Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 17 Feb 2016 14:09:00 +0100 Subject: [PATCH 08/21] remove unused code --- patroni/postgresql.py | 3 +-- tests/test_postgresql.py | 22 +++++++--------------- 2 files changed, 8 insertions(+), 17 deletions(-) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index f5f3d37c..893503e2 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -526,8 +526,7 @@ recovery_target_timeline = 'latest' result[name] = val except IOError: logger.exception('Error when reading postmaster.opts') - finally: - return result + return result def single_user_mode(self, command=None, options=None): """ run a given command in a single-user mode. If the command is empty - then just start and stop """ diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 4c701c60..5558437b 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -2,21 +2,16 @@ import mock # for the mock.call method, importing it without a namespace breaks import os import psycopg2 import shutil +import subprocess import unittest -from six.moves import builtins from mock import Mock, MagicMock, PropertyMock, patch, mock_open from patroni.dcs import Cluster, Leader, Member from patroni.exceptions import PostgresException, PostgresConnectionException from patroni.postgresql import Postgresql from patroni.utils import RetryFailedError +from six.moves import builtins from test_ha import false -import subprocess - - -def is_file_raise_on_backup(*args, **kwargs): - if args[0].endswith('.backup'): - raise Exception("foo") class MockCursor(object): @@ -283,7 +278,7 @@ class TestPostgresql(unittest.TestCase): self.p.pg_rewind = tmp with mock.patch('subprocess.call', MagicMock(return_value=1)): self.assertFalse(self.p.can_rewind) - with mock.patch('subprocess.call', side_effect=OSError("foo")): + with mock.patch('subprocess.call', side_effect=OSError): self.assertFalse(self.p.can_rewind) tmp = self.p.controldata self.p.controldata = lambda: {'wal_log_hints setting': 'on'} @@ -292,7 +287,7 @@ class TestPostgresql(unittest.TestCase): @patch('time.sleep', Mock()) def test_create_replica(self): - self.p.delete_trigger_file = Mock(side_effect=OSError()) + self.p.delete_trigger_file = Mock(side_effect=OSError) with patch('subprocess.call', Mock(side_effect=[1, 0])): self.assertEquals(self.p.create_replica(self.leader, ''), 0) with patch('subprocess.call', Mock(side_effect=[Exception(), 0])): @@ -349,7 +344,7 @@ class TestPostgresql(unittest.TestCase): def test_last_operation(self): self.assertEquals(self.p.last_operation(), '0') - @patch('subprocess.Popen', Mock(side_effect=OSError())) + @patch('subprocess.Popen', Mock(side_effect=OSError)) def test_call_nowait(self): self.assertFalse(self.p.call_nowait('on_start')) @@ -369,7 +364,7 @@ class TestPostgresql(unittest.TestCase): def test_move_data_directory(self): self.p.is_running = false self.p.move_data_directory() - with patch('os.rename', Mock(side_effect=OSError())): + with patch('os.rename', Mock(side_effect=OSError)): self.p.move_data_directory() @patch('patroni.postgresql.Postgresql.write_pgpass', MagicMock(return_value=dict())) @@ -411,13 +406,10 @@ class TestPostgresql(unittest.TestCase): self.assertEquals(int(data['max_replication_slots']), 5) self.assertEqual(data.get('D'), None) - m.side_effect = IOError("foo") + m.side_effect = IOError data = self.p.read_postmaster_opts() self.assertEqual(data, dict()) - m.side_effect = Exception("foo") - self.assertRaises(Exception, self.p.read_postmaster_opts()) - @patch('subprocess.Popen') @patch.object(builtins, 'open', MagicMock(return_value=42)) def test_single_user_mode(self, subprocess_popen_mock): From a210cfd1abde26a56f6fae93c13c7c72d720a4f7 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 17 Feb 2016 14:51:59 +0100 Subject: [PATCH 09/21] Fix more codacy issues --- tests/test_api.py | 16 ++++++++-------- tests/test_ctl.py | 43 ++++++++++++++++++------------------------- tests/test_utils.py | 6 +++--- 3 files changed, 29 insertions(+), 36 deletions(-) diff --git a/tests/test_api.py b/tests/test_api.py index 2faee523..265394e8 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -87,7 +87,7 @@ class MockRestApiServer(RestApiServer): @patch('ssl.wrap_socket', Mock(return_value=0)) class TestRestApiHandler(unittest.TestCase): - def test_do_GET(self): + def test_do_GET(*args): MockRestApiServer(RestApiHandler, b'GET /replica') with patch.object(RestApiHandler, 'get_postgresql_status', Mock(return_value={})): MockRestApiServer(RestApiHandler, b'GET /replica') @@ -103,7 +103,7 @@ class TestRestApiHandler(unittest.TestCase): MockRestApiServer(RestApiHandler, b'GET /master') MockRestApiServer(RestApiHandler, b'GET /master') - def test_do_OPTIONS(self): + def test_do_OPTIONS(*args): MockRestApiServer(RestApiHandler, b'OPTIONS / HTTP/1.0') with patch.object(BaseHTTPRequestHandler, 'handle_one_request') as mock_handle_request: @@ -117,14 +117,14 @@ class TestRestApiHandler(unittest.TestCase): makefile.return_value.flush = Mock(side_effect=socket.error("foo")) MockRestApiServer(RestApiHandler, b'OPTIONS / HTTP/1.0') - def test_do_GET_patroni(self): + def test_do_GET_patroni(*args): MockRestApiServer(RestApiHandler, b'GET /patroni') - def test_basicauth(self): + def test_basicauth(*args): MockRestApiServer(RestApiHandler, b'POST /restart HTTP/1.0') MockRestApiServer(RestApiHandler, b'POST /restart HTTP/1.0\nAuthorization:') - def test_do_POST_restart(self): + def test_do_POST_restart(*args): request = b'POST /restart HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0' MockRestApiServer(RestApiHandler, request) with patch.object(MockHa, 'restart', Mock(side_effect=Exception)): @@ -140,10 +140,10 @@ class TestRestApiHandler(unittest.TestCase): with patch.object(MockHa, 'schedule_reinitialize', Mock(return_value=None)): MockRestApiServer(RestApiHandler, request) cluster.leader.name = 'test' - MockRestApiServer(RestApiHandler, request) + self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) @patch('time.sleep', Mock()) - def test_RestApiServer_query(self): + def test_RestApiServer_query(*args): with patch.object(MockCursor, 'execute', Mock(side_effect=psycopg2.OperationalError)): MockRestApiServer(RestApiHandler, b'GET /patroni') with patch.object(MockPostgresql, 'connection', Mock(side_effect=psycopg2.OperationalError)): @@ -175,4 +175,4 @@ class TestRestApiHandler(unittest.TestCase): MockRestApiServer(RestApiHandler, request) request = b'POST /failover HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0\n' +\ b'Content-Length: 50\n\n{"leader": "postgresql1", "member": "postgresql2"}' - MockRestApiServer(RestApiHandler, request) + self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index fb63aebd..12ab4ad8 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -69,20 +69,16 @@ class TestCtl(unittest.TestCase): @patch('psycopg2.connect', psycopg2_connect) def test_get_cursor(self): - c = get_cursor(get_cluster_initialized_without_leader(), role='master') - assert c is None + self.assertIsNone(get_cursor(get_cluster_initialized_without_leader(), role='master')) - c = get_cursor(get_cluster_initialized_with_leader(), role='master') - assert c is not None + self.assertIsNotNone(get_cursor(get_cluster_initialized_with_leader(), role='master')) - c = get_cursor(get_cluster_initialized_with_leader(), role='replica') - # # MockCursor returns pg_is_in_recovery as false - assert c is None + # MockCursor returns pg_is_in_recovery as false + self.assertIsNone(get_cursor(get_cluster_initialized_with_leader(), role='replica')) - c = get_cursor(get_cluster_initialized_with_leader(), role='any') - assert c is not None + self.assertIsNotNone(get_cursor(get_cluster_initialized_with_leader(), role='any')) - def test_output_members(self): + def test_output_members(*args): cluster = get_cluster_initialized_with_leader() output_members(cluster, name='abc', fmt='pretty') output_members(cluster, name='abc', fmt='json') @@ -224,17 +220,17 @@ y''') @patch('patroni.ctl.get_cursor', Mock(return_value=MockConnect().cursor())) def test_query_member(self): rows = query_member(None, None, None, 'master', 'SELECT pg_is_in_recovery()') - assert 'False' in str(rows) + self.assertTrue('False' in str(rows)) rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()') - assert rows == (None, None) + self.assertEquals(rows, (None, None)) with patch('patroni.ctl.get_cursor', Mock(return_value=None)): rows = query_member(None, None, None, None, 'SELECT pg_is_in_recovery()') - assert 'No connection to' in str(rows) + self.assertTrue('No connection to' in str(rows)) rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()') - assert 'No connection to' in str(rows) + self.assertTrue('No connection to' in str(rows)) with patch('patroni.ctl.get_cursor', Mock(side_effect=psycopg2.OperationalError('bla'))): rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()') @@ -337,26 +333,23 @@ leader''') assert 'Usage:' in result.output def test_get_any_member(self): - m = get_any_member(get_cluster_initialized_without_leader(), role='master') - assert m is None + self.assertIsNone(get_any_member(get_cluster_initialized_without_leader(), role='master')) m = get_any_member(get_cluster_initialized_with_leader(), role='master') - assert m.name == 'leader' + self.assertEquals(m.name, 'leader') def test_get_all_members(self): - r = list(get_all_members(get_cluster_initialized_without_leader(), role='master')) - assert len(r) == 0 + self.assertEquals(list(get_all_members(get_cluster_initialized_without_leader(), role='master')), []) r = list(get_all_members(get_cluster_initialized_with_leader(), role='master')) - assert len(r) == 1 - assert r[0].name == 'leader' + self.assertEquals(len(r), 1) + self.assertEquals(r[0].name, 'leader') r = list(get_all_members(get_cluster_initialized_with_leader(), role='replica')) - assert len(r) == 1 - assert r[0].name == 'other' + self.assertEquals(len(r), 1) + self.assertEquals(r[0].name, 'other') - r = list(get_all_members(get_cluster_initialized_without_leader(), role='replica')) - assert len(r) == 2 + self.assertEquals(len(list(get_all_members(get_cluster_initialized_without_leader(), role='replica'))), 2) @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)) diff --git a/tests/test_utils.py b/tests/test_utils.py index 6f66f4c3..740fef03 100644 --- a/tests/test_utils.py +++ b/tests/test_utils.py @@ -16,14 +16,14 @@ class TestUtils(unittest.TestCase): @patch('time.sleep', Mock()) def test_reap_children(self): - reap_children() + self.assertIsNone(reap_children()) with patch('os.waitpid', Mock(return_value=(0, 0))): sigchld_handler(None, None) - reap_children() + self.assertIsNone(reap_children()) @patch('time.sleep', time_sleep) def test_sleep(self): - sleep(0.01) + self.assertIsNone(sleep(0.01)) @patch('time.sleep', Mock()) From 4038d94c5ac61b345c5248d72a640e7bfd975073 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 17 Feb 2016 14:59:17 +0100 Subject: [PATCH 10/21] Fix more codacy issues --- tests/test_api.py | 26 +++++++++++++------------- tests/test_ctl.py | 8 ++++---- 2 files changed, 17 insertions(+), 17 deletions(-) diff --git a/tests/test_api.py b/tests/test_api.py index 265394e8..dd0f3dfa 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -87,7 +87,7 @@ class MockRestApiServer(RestApiServer): @patch('ssl.wrap_socket', Mock(return_value=0)) class TestRestApiHandler(unittest.TestCase): - def test_do_GET(*args): + def test_do_GET(self): MockRestApiServer(RestApiHandler, b'GET /replica') with patch.object(RestApiHandler, 'get_postgresql_status', Mock(return_value={})): MockRestApiServer(RestApiHandler, b'GET /replica') @@ -101,10 +101,10 @@ class TestRestApiHandler(unittest.TestCase): MockRestApiServer(RestApiHandler, b'GET /master') with patch.object(MockHa, 'restart_scheduled', Mock(return_value=True)): MockRestApiServer(RestApiHandler, b'GET /master') - MockRestApiServer(RestApiHandler, b'GET /master') + self.assertIsNotNone(MockRestApiServer(RestApiHandler, b'GET /master')) - def test_do_OPTIONS(*args): - MockRestApiServer(RestApiHandler, b'OPTIONS / HTTP/1.0') + def test_do_OPTIONS(self): + self.assertIsNotNone(MockRestApiServer(RestApiHandler, b'OPTIONS / HTTP/1.0')) with patch.object(BaseHTTPRequestHandler, 'handle_one_request') as mock_handle_request: mock_handle_request.side_effect = socket.error("foo") @@ -117,16 +117,16 @@ class TestRestApiHandler(unittest.TestCase): makefile.return_value.flush = Mock(side_effect=socket.error("foo")) MockRestApiServer(RestApiHandler, b'OPTIONS / HTTP/1.0') - def test_do_GET_patroni(*args): - MockRestApiServer(RestApiHandler, b'GET /patroni') + def test_do_GET_patroni(self): + self.assertIsNotNone(MockRestApiServer(RestApiHandler, b'GET /patroni')) - def test_basicauth(*args): - MockRestApiServer(RestApiHandler, b'POST /restart HTTP/1.0') + def test_basicauth(self): + self.assertIsNotNone(MockRestApiServer(RestApiHandler, b'POST /restart HTTP/1.0')) MockRestApiServer(RestApiHandler, b'POST /restart HTTP/1.0\nAuthorization:') - def test_do_POST_restart(*args): + def test_do_POST_restart(self): request = b'POST /restart HTTP/1.0\nAuthorization: Basic dGVzdDp0ZXN0' - MockRestApiServer(RestApiHandler, request) + self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) with patch.object(MockHa, 'restart', Mock(side_effect=Exception)): MockRestApiServer(RestApiHandler, request) @@ -143,11 +143,11 @@ class TestRestApiHandler(unittest.TestCase): self.assertIsNotNone(MockRestApiServer(RestApiHandler, request)) @patch('time.sleep', Mock()) - def test_RestApiServer_query(*args): + def test_RestApiServer_query(self): with patch.object(MockCursor, 'execute', Mock(side_effect=psycopg2.OperationalError)): - MockRestApiServer(RestApiHandler, b'GET /patroni') + self.assertIsNotNone(MockRestApiServer(RestApiHandler, b'GET /patroni')) with patch.object(MockPostgresql, 'connection', Mock(side_effect=psycopg2.OperationalError)): - MockRestApiServer(RestApiHandler, b'GET /patroni') + self.assertIsNotNone(MockRestApiServer(RestApiHandler, b'GET /patroni')) @patch('time.sleep', Mock()) @patch.object(MockHa, 'dcs') diff --git a/tests/test_ctl.py b/tests/test_ctl.py index 12ab4ad8..8559e793 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -78,11 +78,11 @@ class TestCtl(unittest.TestCase): self.assertIsNotNone(get_cursor(get_cluster_initialized_with_leader(), role='any')) - def test_output_members(*args): + def test_output_members(self): cluster = get_cluster_initialized_with_leader() - output_members(cluster, name='abc', fmt='pretty') - output_members(cluster, name='abc', fmt='json') - output_members(cluster, name='abc', fmt='tsv') + self.assertIsNone(output_members(cluster, name='abc', fmt='pretty')) + self.assertIsNone(output_members(cluster, name='abc', fmt='json')) + self.assertIsNone(output_members(cluster, name='abc', fmt='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)) From eb1e6788202b771158d6dc127c3b752239a3dfdd Mon Sep 17 00:00:00 2001 From: Jan Keirse Date: Thu, 18 Feb 2016 11:04:41 +0100 Subject: [PATCH 11/21] sample systemd service file --- extras/startup-scripts/patroni.service | 28 ++++++++++++++++++++++++++ 1 file changed, 28 insertions(+) create mode 100644 extras/startup-scripts/patroni.service diff --git a/extras/startup-scripts/patroni.service b/extras/startup-scripts/patroni.service new file mode 100644 index 00000000..67e7ee14 --- /dev/null +++ b/extras/startup-scripts/patroni.service @@ -0,0 +1,28 @@ +# This is an example systemd config file for Patroni +# You can copy it to "/etc/systemd/system/patroni.service", + +[Unit] +Description=Runners to orchestrate a high-availability PostgreSQL +After=syslog.target network.target + +[Service] +Type=simple + +User=postgres +Group=postgres + +# Where to send early-startup messages from the server +# This is normally controlled by the global default set by systemd +# StandardOutput=syslog + +ExecStart=/bin/patroni /etc/patroni.yml + +# Give a reasonable amount of time for the server to start up/shut down +TimeoutSec=10 + +# Always restart the service if it crashes, we want it to continue running +Restart=no + +[Install] +WantedBy=multi-user.target + From e68e253d166cad71125d86bc7a652ef8ef1c79fc Mon Sep 17 00:00:00 2001 From: Jan Keirse Date: Thu, 18 Feb 2016 11:06:40 +0100 Subject: [PATCH 12/21] Add patroni.service file documentation. --- extras/startup-scripts/README.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/extras/startup-scripts/README.md b/extras/startup-scripts/README.md index 244ee65e..7cc0d445 100644 --- a/extras/startup-scripts/README.md +++ b/extras/startup-scripts/README.md @@ -8,3 +8,6 @@ Scripts supplied: ### patroni.upstart.conf Upstart job for Ubuntu 12.04 or 14.04. Requires Upstart > 1.4. Intended for systems where Patroni has been installed on a base system, rather than in Docker. + +### patroni.service +Systemd service file, to be copied to /etc/systemd/system/patroni.service, tested on Centos 7.1 with Patroni installed from pip. From 753ba835f11e771925da3ef7c2b313865bf4fcb0 Mon Sep 17 00:00:00 2001 From: Jan Keirse Date: Thu, 18 Feb 2016 11:08:18 +0100 Subject: [PATCH 13/21] wrong comment about restart --- extras/startup-scripts/patroni.service | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/extras/startup-scripts/patroni.service b/extras/startup-scripts/patroni.service index 67e7ee14..fdd7558b 100644 --- a/extras/startup-scripts/patroni.service +++ b/extras/startup-scripts/patroni.service @@ -20,7 +20,7 @@ ExecStart=/bin/patroni /etc/patroni.yml # Give a reasonable amount of time for the server to start up/shut down TimeoutSec=10 -# Always restart the service if it crashes, we want it to continue running +# Do not restart the service if it crashes, we want to manually inspect database on failure Restart=no [Install] From 26e15862884933fced7694e240605d99e02da4f6 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Mon, 22 Feb 2016 12:20:20 +0100 Subject: [PATCH 14/21] 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 15/21] 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 From 641cc4013e76156dc27b22642499341a8cc703e9 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 23 Feb 2016 11:46:49 +0100 Subject: [PATCH 16/21] Mock a few of methods in Postgresql class instead of the whole class --- tests/test_ha.py | 100 +++++++++++++++++------------------------------ 1 file changed, 36 insertions(+), 64 deletions(-) diff --git a/tests/test_ha.py b/tests/test_ha.py index 3589e952..cdf1156f 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -1,13 +1,14 @@ -import etcd import unittest import datetime import pytz +from etcd import EtcdException from mock import Mock, MagicMock, patch from patroni.dcs import Cluster, Failover, Leader, Member from patroni.etcd import Client, Etcd from patroni.exceptions import DCSError, PostgresException from patroni.ha import Ha +from patroni.postgresql import Postgresql from test_etcd import socket_getaddrinfo, etcd_read, etcd_write, requests_get @@ -45,56 +46,6 @@ def get_cluster_initialized_with_only_leader(failover=None): return get_cluster(True, l, [l], failover) -class MockPostgresql(Mock): - - name = 'postgresql0' - role = 'replica' - state = 'running' - connection_string = 'postgres://foo@bar/postgres' - server_version = '999999' - scope = 'dummy' - - @staticmethod - def is_healthy(): - return True - - @staticmethod - def start(): - return True - - @staticmethod - def is_healthiest_node(members): - return True - - @staticmethod - def is_leader(): - return True - - @staticmethod - def xlog_position(): - return 0 - - @staticmethod - def last_operation(): - return 0 - - @staticmethod - def data_directory_empty(): - return False - - @staticmethod - def bootstrap(*args, **kwargs): - return True - - @staticmethod - def check_replication_lag(last_leader_operation): - return True - - @staticmethod - def check_recovery_conf(leader): - return False - - class MockPatroni(object): def __init__(self, p, d): @@ -111,18 +62,36 @@ def run_async(func, args=()): return func(*args) if args else func() +@patch.object(Postgresql, 'is_running', Mock(return_value=True)) +@patch.object(Postgresql, 'is_leader', Mock(return_value=True)) +@patch.object(Postgresql, 'xlog_position', Mock(return_value=0)) +@patch.object(Postgresql, 'call_nowait', Mock(return_value=True)) +@patch.object(Postgresql, 'data_directory_empty', Mock(return_value=False)) +@patch.object(Postgresql, 'controldata', Mock(return_value={})) +@patch.object(Postgresql, 'sync_replication_slots', Mock()) +@patch.object(Postgresql, 'write_pg_hba', Mock()) +@patch.object(Postgresql, 'write_pgpass', Mock()) +@patch.object(Postgresql, 'write_recovery_conf', Mock()) +@patch.object(Postgresql, 'query', Mock()) +@patch.object(Postgresql, 'checkpoint', Mock()) +@patch('subprocess.call', Mock(return_value=0)) class TestHa(unittest.TestCase): @patch('socket.getaddrinfo', socket_getaddrinfo) def setUp(self): with patch.object(Client, 'machines') as mock_machines: mock_machines.__get__ = Mock(return_value=['http://remotehost:2379']) - self.p = MockPostgresql() + self.p = Postgresql({'name': 'postgresql0', 'scope': 'dummy', 'listen': '127.0.0.1:5432', + 'data_dir': 'data/postgresql0', 'superuser': {}, 'admin': {}, + 'replication': {'username': '', 'password': '', 'network': ''}}) + self.p._state = 'running' + self.p._sysid = '1234567890' + self.p.check_replication_lag = true self.p.can_create_replica_without_leader = MagicMock(return_value=False) self.e = Etcd('foo', {'ttl': 30, 'host': 'ok:2379', 'scope': 'test'}) self.e.client.read = etcd_read self.e.client.write = etcd_write - self.e.client.delete = Mock(side_effect=etcd.EtcdException()) + self.e.client.delete = Mock(side_effect=EtcdException()) self.ha = Ha(MockPatroni(self.p, self.e)) self.ha._async_executor.run_async = run_async self.ha.old_cluster = self.e.get_cluster() @@ -154,7 +123,7 @@ class TestHa(unittest.TestCase): self.p.is_healthy = false self.p.is_running = false self.ha.has_lock = true - self.p.role = 'master' + self.p._role = 'master' self.p.controldata = lambda: {'Database cluster state': 'in production'} self.assertEquals(self.ha.run_cycle(), 'started as readonly because i had the session lock') self.assertEquals(self.ha.run_cycle(), 'removed leader key after trying and failing to start postgres') @@ -302,57 +271,60 @@ class TestHa(unittest.TestCase): self.ha.has_lock = true 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, None)) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', self.p.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', None)) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') - f = Failover(0, MockPostgresql.name, '', None) + f = Failover(0, self.p.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, None)) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', self.p.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.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', self.p.name, scheduled)) 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.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', self.p.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.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', self.p.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.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', self.p.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.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', self.p.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, None)) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', self.p.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', None)) + self.p._role = 'replica' 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, '', None)) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, self.p.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.p._role = 'replica' 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', None)) self.ha.fetch_node_status = lambda e: (e, True, True, 0, {'nofailover': 'True'}) # accessible, in_recovery + self.p._role = 'replica' 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 self.ha.patroni.nofailover = True From dd20fc7e71ad05423b65fa6d2c24e77a25888d44 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 23 Feb 2016 11:47:47 +0100 Subject: [PATCH 17/21] Refactor follow method --- patroni/postgresql.py | 76 +++++++++++++++++++++---------------------- 1 file changed, 37 insertions(+), 39 deletions(-) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 893503e2..a1cedb58 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -559,47 +559,45 @@ recovery_target_timeline = 'latest' logger.exception("Unable to list %s", status_dir) def follow(self, leader, recovery=False): - if not self.check_recovery_conf(leader) or recovery: - change_role = (self.role == 'master') - - self._need_rewind = (self._need_rewind or change_role) and self.can_rewind - if self._need_rewind: - logger.info("set the rewind flag after demote") - self.write_recovery_conf(leader) - if not leader or not self._need_rewind: # do not rewind until the leader becomes available - ret = self.restart() - else: # we have a leader and need to rewind - if self.is_running(): - self.stop() - # at present, pg_rewind only runs when the cluster is shut down cleanly - # and not shutdown in recovery. We have to remove the recovery.conf if present - # and start/shutdown in a single user mode to emulate this. - # XXX: if recovery.conf is linked, it will be written anew as a normal file. - if os.path.islink(self.recovery_conf): - os.unlink(self.recovery_conf) - else: - os.remove(self.recovery_conf) - # Archived segments might be useful to pg_rewind, - # clean the flags that tell we should remove them. - self.cleanup_archive_status() - # Start in a single user mode and stop to produce a clean shutdown - opts = self.read_postmaster_opts() - opts['archive_mode'] = 'on' - opts['archive_command'] = 'false' - self.single_user_mode(options=opts) - if self.rewind(leader): - ret = self.start() - else: - logger.error("unable to rewind the former master") - self.remove_data_directory() - ret = True - self._need_rewind = False - if change_role and ret: - self.call_nowait(ACTION_ON_ROLE_CHANGE) - return ret - else: + if self.check_recovery_conf(leader) and not recovery: return True + change_role = self.role == 'master' + self._need_rewind = (self._need_rewind or change_role) and self.can_rewind + if self._need_rewind: + logger.info("set the rewind flag after demote") + self.write_recovery_conf(leader) + if leader and self._need_rewind: # we have a leader and need to rewind + if self.is_running(): + self.stop() + # at present, pg_rewind only runs when the cluster is shut down cleanly + # and not shutdown in recovery. We have to remove the recovery.conf if present + # and start/shutdown in a single user mode to emulate this. + # XXX: if recovery.conf is linked, it will be written anew as a normal file. + if os.path.islink(self.recovery_conf): + os.unlink(self.recovery_conf) + else: + os.remove(self.recovery_conf) + # Archived segments might be useful to pg_rewind, + # clean the flags that tell we should remove them. + self.cleanup_archive_status() + # Start in a single user mode and stop to produce a clean shutdown + opts = self.read_postmaster_opts() + opts.update({'archive_mode': 'on', 'archive_command': 'false'}) + self.single_user_mode(options=opts) + if self.rewind(leader): + ret = self.start() + else: + logger.error("unable to rewind the former master") + self.remove_data_directory() + ret = True + self._need_rewind = False + else: # do not rewind until the leader becomes available + ret = self.restart() + if change_role and ret: + self.call_nowait(ACTION_ON_ROLE_CHANGE) + return ret + def save_configuration_files(self): """ copy postgresql.conf to postgresql.conf.backup to be able to retrive configuration files From ce33090c0d4d0665ab80788e106b4059dec5bd99 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 23 Feb 2016 11:48:52 +0100 Subject: [PATCH 18/21] Mock dcs.watch directly instead of using wraper --- tests/test_patroni.py | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/tests/test_patroni.py b/tests/test_patroni.py index aaefdf1b..2698de58 100644 --- a/tests/test_patroni.py +++ b/tests/test_patroni.py @@ -15,10 +15,6 @@ from test_postgresql import Postgresql, psycopg2_connect from test_zookeeper import MockKazooClient -def time_sleep(*args): - raise SleepException() - - @patch('time.sleep', Mock()) @patch('subprocess.call', Mock(return_value=0)) @patch('psycopg2.connect', psycopg2_connect) @@ -62,7 +58,7 @@ class TestPatroni(unittest.TestCase): @patch('time.sleep', Mock(side_effect=SleepException())) def test_run(self): - self.p.ha.dcs.watch = time_sleep + self.p.ha.dcs.watch = Mock(side_effect=SleepException()) self.assertRaises(SleepException, self.p.run) self.p.ha.state_handler.is_leader = Mock(return_value=False) From 6b3c4697fc36409280911e378d7c146854bbb08e Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 23 Feb 2016 11:49:22 +0100 Subject: [PATCH 19/21] Remove unused code --- tests/test_ctl.py | 40 ++++++++++++++-------------------------- 1 file changed, 14 insertions(+), 26 deletions(-) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index a877ed4b..b2590da5 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -1,25 +1,19 @@ -#!/usr/bin/env python -# -*- coding: utf-8 -*- - import os import pytest import unittest -import psycopg2 -import requests -import patroni.exceptions -import etcd -from mock import patch, Mock, MagicMock - from click.testing import CliRunner +from etcd import EtcdException +from mock import patch, Mock, MagicMock from patroni.ctl import ctl, members, store_config, load_config, output_members, post_patroni, get_dcs, \ wait_for_leader, get_all_members, get_any_member, get_cursor, query_member, configure -from patroni.ha import Ha from patroni.etcd import Etcd, Client -from test_ha import get_cluster_initialized_without_leader, get_cluster_initialized_with_leader, \ - get_cluster_initialized_with_only_leader, MockPostgresql, MockPatroni, run_async, \ - get_cluster_not_initialized_without_leader +from patroni.exceptions import PatroniCtlException +from psycopg2 import OperationalError +from requests.exceptions import ConnectionError from test_etcd import etcd_read, etcd_write, requests_get, socket_getaddrinfo, MockResponse +from test_ha import get_cluster_initialized_without_leader, get_cluster_initialized_with_leader, \ + get_cluster_initialized_with_only_leader from test_postgresql import MockConnect, psycopg2_connect CONFIG_FILE_PATH = './test-ctl.yaml' @@ -56,16 +50,10 @@ class TestCtl(unittest.TestCase): self.runner = CliRunner() with patch.object(Client, 'machines') as mock_machines: mock_machines.__get__ = Mock(return_value=['http://remotehost:2379']) - self.p = MockPostgresql() self.e = Etcd('foo', {'ttl': 30, 'host': 'ok:2379', 'scope': 'test'}) self.e.client.read = etcd_read self.e.client.write = etcd_write - self.e.client.delete = Mock(side_effect=etcd.EtcdException()) - self.ha = Ha(MockPatroni(self.p, self.e)) - self.ha._async_executor.run_async = run_async - self.ha.old_cluster = self.e.get_cluster() - self.ha.cluster = get_cluster_not_initialized_without_leader() - self.ha.load_cluster_from_dcs = Mock() + self.e.client.delete = Mock(side_effect=EtcdException) @patch('psycopg2.connect', psycopg2_connect) def test_get_cursor(self): @@ -186,7 +174,7 @@ y''') assert 'Failover failed' in result.output def test_(self): - self.assertRaises(patroni.exceptions.PatroniCtlException, get_dcs, {'scheme': 'dummy'}, 'dummy') + self.assertRaises(PatroniCtlException, get_dcs, {'scheme': 'dummy'}, 'dummy') @patch('psycopg2.connect', psycopg2_connect) @patch('patroni.ctl.query_member', Mock(return_value=([['mock column']], None))) @@ -248,10 +236,10 @@ y''') rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()') self.assertTrue('No connection to' in str(rows)) - with patch('patroni.ctl.get_cursor', Mock(side_effect=psycopg2.OperationalError('bla'))): + with patch('patroni.ctl.get_cursor', Mock(side_effect=OperationalError('bla'))): rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()') - with patch('test_postgresql.MockCursor.execute', Mock(side_effect=psycopg2.OperationalError('bla'))): + with patch('test_postgresql.MockCursor.execute', Mock(side_effect=OperationalError('bla'))): rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()') @patch('patroni.dcs.AbstractDCS.get_cluster', Mock(return_value=get_cluster_initialized_with_leader())) @@ -341,15 +329,15 @@ leader''') @patch('patroni.etcd.Etcd.get_cluster', Mock(return_value=get_cluster_initialized_with_leader())) def test_wait_for_leader(self): dcs = self.e - self.assertRaises(patroni.exceptions.PatroniCtlException, wait_for_leader, dcs, 0) + self.assertRaises(PatroniCtlException, wait_for_leader, dcs, 0) cluster = wait_for_leader(dcs=dcs, timeout=2) assert cluster.leader.member.name == 'leader' def test_post_patroni(self): - with patch('requests.post', MagicMock(side_effect=requests.exceptions.ConnectionError('foo'))): + with patch('requests.post', MagicMock(side_effect=ConnectionError('foo'))): member = get_cluster_initialized_with_leader().leader.member - self.assertRaises(requests.exceptions.ConnectionError, post_patroni, member, 'dummy', {}) + self.assertRaises(ConnectionError, post_patroni, member, 'dummy', {}) def test_ctl(self): self.runner.invoke(ctl, ['list']) From 756158a735efd3cbb127e5749fc82b8b232a1d74 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 23 Feb 2016 11:59:02 +0100 Subject: [PATCH 20/21] make codacy and quantifiedcode happier --- tests/test_ctl.py | 6 +++--- tests/test_ha.py | 10 +++++----- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index b2590da5..3987619d 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -1,5 +1,6 @@ import os import pytest +import requests.exceptions import unittest from click.testing import CliRunner @@ -10,7 +11,6 @@ from patroni.ctl import ctl, members, store_config, load_config, output_members, from patroni.etcd import Etcd, Client from patroni.exceptions import PatroniCtlException from psycopg2 import OperationalError -from requests.exceptions import ConnectionError from test_etcd import etcd_read, etcd_write, requests_get, socket_getaddrinfo, MockResponse from test_ha import get_cluster_initialized_without_leader, get_cluster_initialized_with_leader, \ get_cluster_initialized_with_only_leader @@ -335,9 +335,9 @@ leader''') assert cluster.leader.member.name == 'leader' def test_post_patroni(self): - with patch('requests.post', MagicMock(side_effect=ConnectionError('foo'))): + with patch('requests.post', MagicMock(side_effect=requests.exceptions.ConnectionError('foo'))): member = get_cluster_initialized_with_leader().leader.member - self.assertRaises(ConnectionError, post_patroni, member, 'dummy', {}) + self.assertRaises(requests.exceptions.ConnectionError, post_patroni, member, 'dummy', {}) def test_ctl(self): self.runner.invoke(ctl, ['list']) diff --git a/tests/test_ha.py b/tests/test_ha.py index cdf1156f..57754aa6 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -84,7 +84,7 @@ class TestHa(unittest.TestCase): self.p = Postgresql({'name': 'postgresql0', 'scope': 'dummy', 'listen': '127.0.0.1:5432', 'data_dir': 'data/postgresql0', 'superuser': {}, 'admin': {}, 'replication': {'username': '', 'password': '', 'network': ''}}) - self.p._state = 'running' + self.p.set_state('running') self.p._sysid = '1234567890' self.p.check_replication_lag = true self.p.can_create_replica_without_leader = MagicMock(return_value=False) @@ -123,7 +123,7 @@ class TestHa(unittest.TestCase): self.p.is_healthy = false self.p.is_running = false self.ha.has_lock = true - self.p._role = 'master' + self.p.set_role('master') self.p.controldata = lambda: {'Database cluster state': 'in production'} self.assertEquals(self.ha.run_cycle(), 'started as readonly because i had the session lock') self.assertEquals(self.ha.run_cycle(), 'removed leader key after trying and failing to start postgres') @@ -311,20 +311,20 @@ class TestHa(unittest.TestCase): self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', self.p.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', None)) - self.p._role = 'replica' + self.p.set_role('replica') 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, self.p.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.p._role = 'replica' + self.p.set_role('replica') 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', None)) self.ha.fetch_node_status = lambda e: (e, True, True, 0, {'nofailover': 'True'}) # accessible, in_recovery - self.p._role = 'replica' + self.p.set_role('replica') 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 self.ha.patroni.nofailover = True From ec85e2eb4908a7fa1b50257e1de8cfd2d81f9a23 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 23 Feb 2016 12:05:02 +0100 Subject: [PATCH 21/21] make quantifiedcode happier --- tests/test_ha.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/tests/test_ha.py b/tests/test_ha.py index 57754aa6..a65a444a 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -67,7 +67,7 @@ def run_async(func, args=()): @patch.object(Postgresql, 'xlog_position', Mock(return_value=0)) @patch.object(Postgresql, 'call_nowait', Mock(return_value=True)) @patch.object(Postgresql, 'data_directory_empty', Mock(return_value=False)) -@patch.object(Postgresql, 'controldata', Mock(return_value={})) +@patch.object(Postgresql, 'controldata', Mock(return_value={'Database system identifier': '1234567890'})) @patch.object(Postgresql, 'sync_replication_slots', Mock()) @patch.object(Postgresql, 'write_pg_hba', Mock()) @patch.object(Postgresql, 'write_pgpass', Mock()) @@ -85,7 +85,6 @@ class TestHa(unittest.TestCase): 'data_dir': 'data/postgresql0', 'superuser': {}, 'admin': {}, 'replication': {'username': '', 'password': '', 'network': ''}}) self.p.set_state('running') - self.p._sysid = '1234567890' self.p.check_replication_lag = true self.p.can_create_replica_without_leader = MagicMock(return_value=False) self.e = Etcd('foo', {'ttl': 30, 'host': 'ok:2379', 'scope': 'test'})