diff --git a/.zappr.yml b/.zappr.yml new file mode 100644 index 00000000..4d628636 --- /dev/null +++ b/.zappr.yml @@ -0,0 +1,12 @@ +approvals: + # PR needs at least 4 approvals + minimum: 1 + # approval = comment that matches this regex + pattern: "^:?\\+1:?$" + from: + # commenter must be either one of: + # a public zalando org member + orgs: + - zalando + # a collaborator of the repo + collaborators: true diff --git a/patroni/__init__.py b/patroni/__init__.py index 36c8feac..85b6c5df 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -35,6 +35,10 @@ class Patroni(object): def replicatefrom(self): return self.tags.get('replicatefrom') + @property + def clonefrom(self): + return self.tags.get('clonefrom') + @staticmethod def get_dcs(name, config): if 'etcd' in config: diff --git a/patroni/api.py b/patroni/api.py index 1186321f..342e1bf6 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -143,7 +143,7 @@ class RestApiHandler(BaseHTTPRequestHandler): self.wfile.write(data) def poll_failover_result(self, leader, member): - for a in range(0, 15): + for _ in range(0, 15): time.sleep(1) try: cluster = self.server.patroni.dcs.get_cluster() @@ -166,7 +166,7 @@ class RestApiHandler(BaseHTTPRequestHandler): members = [m for m in cluster.members if m.name != cluster.leader.name and m.api_url] if not members: return b'failover is not possible: cluster does not have members except leader' - for member, reachable, in_recovery, xlog_location, tags in self.server.patroni.ha.fetch_nodes_statuses(members): + for member, reachable, _, xlog_location, tags in self.server.patroni.ha.fetch_nodes_statuses(members): if reachable and not tags.get('nofailover', False): return None return b'failover is not possible: no good candidates have been found' @@ -262,6 +262,7 @@ class RestApiHandler(BaseHTTPRequestHandler): END, pg_xlog_location_diff(pg_last_xlog_receive_location(), '0/0')::bigint, pg_xlog_location_diff(pg_last_xlog_replay_location(), '0/0')::bigint, + to_char(pg_last_xact_replay_timestamp(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'), pg_is_in_recovery() AND pg_is_xlog_replay_paused()""", retry=retry)[0] return { 'state': self.server.patroni.postgresql.state, @@ -271,7 +272,8 @@ class RestApiHandler(BaseHTTPRequestHandler): 'xlog': ({ 'received_location': row[3], 'replayed_location': row[4], - 'paused': row[5]} if row[1] else { + 'replayed_timestamp': row[5], + 'paused': row[6]} if row[1] else { 'location': row[2] }) } diff --git a/patroni/ctl.py b/patroni/ctl.py index c4cc9fe0..0cfd6039 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -18,6 +18,7 @@ import dateutil import tzlocal from .etcd import Etcd +from .zookeeper import ZooKeeper from .exceptions import PatroniCtlException from .postgresql import parseurl @@ -46,12 +47,12 @@ def parse_dcs(dcs): parsed = urlparse('//' + dcs) if scheme == '': - default_schemes = {'2181': 'zookeeper', '8500': 'consul'} + default_schemes = {'2181': 'zookeeper', '8181': 'exhibitor', '8500': 'consul'} scheme = default_schemes.get(str(parsed.port), 'etcd') port = parsed.port if port is None: - default_ports = {'consul': 8500, 'zookeeper': 2181} + default_ports = {'consul': 8500, 'zookeeper': 2181, 'exhibitor': 8181} port = default_ports.get(str(scheme), 4001) return {'scheme': str(scheme), 'hostname': str(parsed.hostname), 'port': int(port)} @@ -105,6 +106,12 @@ def get_dcs(config, scope): if scheme == 'etcd': return Etcd(name=scope, config={'scope': scope, 'host': '{0}:{1}'.format(hostname, port)}) + if scheme == 'zookeeper': + return ZooKeeper(name=scope, config={'scope': scope, 'hosts': [hostname], 'port': port}) + + if scheme == 'exhibitor': + return ZooKeeper(name=scope, config={'scope': scope, 'exhibitor': {'hosts': [hostname], 'port': port}}) + raise PatroniCtlException('Can not find suitable configuration of distributed configuration store') @@ -240,7 +247,7 @@ def dsn(cluster_name, config_file, dcs, role, member): if member is None and role is None: role = 'master' - config, dcs, cluster = ctl_load_config(cluster_name, config_file, dcs) + _, dcs, cluster = ctl_load_config(cluster_name, config_file, dcs) m = get_any_member(cluster=cluster, role=role, member=member) if m is None: raise PatroniCtlException('Can not find a suitable member') diff --git a/patroni/dcs.py b/patroni/dcs.py index 4c2dddad..2768b13e 100644 --- a/patroni/dcs.py +++ b/patroni/dcs.py @@ -3,7 +3,6 @@ import json import dateutil from collections import namedtuple -from patroni.exceptions import DCSError from six.moves.urllib_parse import urlparse, urlunparse, parse_qsl from threading import Event, Lock @@ -147,6 +146,9 @@ class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,mem def has_member(self, member_name): return any(m for m in self.members if m.name == member_name) + def get_member(self, member_name): + return ([m for m in self.members if m.name == member_name] or [None])[0] + class AbstractDCS(object): @@ -269,13 +271,6 @@ class AbstractDCS(object): return self.set_failover_value(json.dumps(failover_value), index) - def current_leader(self): - try: - cluster = self.get_cluster() - return None if cluster.is_unlocked() else cluster.leader - except DCSError: - return None - @abc.abstractmethod def touch_member(self, connection_string, ttl=None): """Update member key in DCS. diff --git a/patroni/ha.py b/patroni/ha.py index f9550e81..8c841290 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -65,19 +65,25 @@ class Ha(object): pass self.dcs.touch_member(json.dumps(data, separators=(',', ':'))) - def clone(self, leader): - if self.state_handler.bootstrap(cluster_initialized=True, current_leader=leader): - logger.info('bootstrapped from leader' if leader else 'bootstrapped without leader') + def clone(self, clone_member, clone_member_name="leader"): + if self.state_handler.bootstrap(cluster_initialized=True, clone_member=clone_member): + logger.info('bootstrapped from {0}'.format(clone_member_name) + if clone_member else 'bootstrapped without leader') else: self.state_handler.stop('immediate') self.state_handler.remove_data_directory() - logger.error('failed to bootstrap from leader' if leader else 'failed to bootstrap (without leader)') + logger.error('failed to bootstrap from {0}'.format(clone_member_name) + if clone_member else 'failed to bootstrap (without leader)') def bootstrap(self): if not self.cluster.is_unlocked(): # cluster already has leader - self._async_executor.schedule('bootstrap from leader') - self._async_executor.run_async(self.clone, args=(self.cluster.leader, )) - return 'trying to bootstrap from leader' + clonefrom = self.patroni.clonefrom + clone_member = self.cluster.get_member(clonefrom)\ + if self.cluster.has_member(clonefrom) else self.cluster.leader + clone_member_name = 'leader' if clone_member == self.cluster.leader else 'replica \'{0}\''.format(clonefrom) + self._async_executor.schedule('bootstrap from {0}'.format(clone_member_name)) + self._async_executor.run_async(self.clone, args=(clone_member, clone_member_name)) + return 'trying to bootstrap from {0}'.format(clone_member_name) elif not self.cluster.initialize and not self.patroni.nofailover: # no initialize key if self.dcs.initialize(create_new=True): # race for initialization try: @@ -96,7 +102,7 @@ class Ha(object): else: return 'failed to acquire initialize lock' else: - if self.state_handler.can_create_replica_without_leader(): + if self.state_handler.can_create_replica_without_replication_connection(): self._async_executor.run_async(self.clone, args=(None, )) return "trying to bootstrap without leader" return 'waiting for leader to bootstrap' @@ -203,7 +209,7 @@ class Ha(object): ret = False members = [m for m in members if m.name != self.state_handler.name and not m.nofailover and m.api_url] if members: - for member, reachable, in_recovery, xlog_location, tags in self.fetch_nodes_statuses(members): + for member, reachable, _, _, tags in self.fetch_nodes_statuses(members): if reachable and not tags.get('nofailover', False): ret = True # TODO: check xlog_location elif not reachable: @@ -223,7 +229,7 @@ class Ha(object): # find specific node and check that it is healthy members = [m for m in self.cluster.members if m.name == failover.member] if members: - member, reachable, in_recovery, xlog_location, tags = self.fetch_node_status(members[0]) + member, reachable, _, _, tags = self.fetch_node_status(members[0]) if reachable and not tags.get('nofailover', False): # node is healthy logger.info('manual failover: to %s, i am %s', member.name, self.state_handler.name) return False diff --git a/patroni/postgresql.py b/patroni/postgresql.py index d875057b..df17200d 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -233,9 +233,10 @@ class Postgresql(object): env['PGPASSFILE'] = self.pgpass return env - def sync_replica(self, leader): - env = self.write_pgpass(parseurl(leader.conn_url)) if leader else os.environ.copy() - if self.create_replica(leader, env) == 0: + def sync_replica(self, clone_member): + # add the credentials to connect to the replica origin to pgpass. + env = self.write_pgpass(parseurl(clone_member.conn_url)) if clone_member else os.environ.copy() + if self.create_replica(clone_member, env) == 0: self.delete_trigger_file() return True return False @@ -248,33 +249,35 @@ class Postgresql(object): """ return ' '.join('{0}={1}'.format(param, val) for param, val in sorted(conn.items())) - def replica_method_can_work_without_leader(self, method): + def replica_method_can_work_without_replication_connection(self, method): return method != 'basebackup' and self.config and self.config.get(method, {}).get('no_master') - def can_create_replica_without_leader(self): + def can_create_replica_without_replication_connection(self): """ go through the replication methods to see if there are ones - that does not require a running leader to create the replica. + that does not require a working replication connection. """ replica_methods = self.config.get('create_replica_method', []) - return any(self.replica_method_can_work_without_leader(replica_method) for replica_method in replica_methods) + return any(self.replica_method_can_work_without_replication_connection(replica_method) + for replica_method in replica_methods) - def create_replica(self, leader, env): + def create_replica(self, clone_member, env): # create the replica according to the replica_method # defined by the user. this is a list, so we need to # loop through all methods the user supplies - connstring = leader.conn_url if leader else "" + connstring = clone_member.conn_url if clone_member else "" # get list of replica methods from config. # If there is no configuration key, or no value is specified, use basebackup replica_methods = self.config.get('create_replica_method') or ['basebackup'] - # if we don't have any leader, leave only replica methods that work without it - replica_methods = [r for r in replica_methods if self.replica_method_can_work_without_leader(r)] if not leader \ - else replica_methods + # if we don't have any source, leave only replica methods that work without it + replica_methods = \ + [r for r in replica_methods if self.replica_method_can_work_without_replication_connection(r)]\ + if not clone_member else replica_methods # go through them in priority order ret = 1 for replica_method in replica_methods: # if the method is basebackup, then use the built-in if replica_method == "basebackup": - ret = self.basebackup(leader, env) + ret = self.basebackup(clone_member, env) if ret == 0: logger.info("replica has been created using basebackup") # if basebackup succeeds, exit with success @@ -698,18 +701,19 @@ $$""".format(name, options), name, password, password) def last_operation(self): return str(self.xlog_position()) - def bootstrap(self, cluster_initialized=False, current_leader=None): + def bootstrap(self, cluster_initialized=False, clone_member=None): """ Populate PostgreSQL data directory by doing one of the following: - create with initdb if there is no master. - - initialize the replica from an existing master + - initialize the replica from an existing member (master or replica) - initialize the replica using the replica creation method that - works without the master (i.e. restore from on-disk base backup) + works without the replication connection (i.e. restore from on-disk + base backup) The choice between the last 2 is triggered by the initialize flag. We should never try to initdb an already initialized cluster, nor - try to bootstrap the cluster that lacks the initialize key from from - the master-less replica creation method (in the latter case, there is + try to bootstrap the cluster that lacks the initialize key using the + master-less replica creation method (in the latter case, there is no clear inidicator of the moment we should abandon our attempts and swich to initdb). @@ -719,7 +723,7 @@ $$""".format(name, options), name, password, password) that should be retried in the future. """ ret = False - if not (cluster_initialized or current_leader): + if not (cluster_initialized or clone_member): ret = self.initialize() and self.start() if ret: self.create_replication_user() @@ -727,9 +731,9 @@ $$""".format(name, options), name, password, password) else: raise PostgresException("Could not bootstrap master PostgreSQL") else: - if self.sync_replica(current_leader): + if self.sync_replica(clone_member): self.restore_configuration_files() - self.write_recovery_conf(current_leader, True) + self.write_recovery_conf(clone_member, True) ret = self.start() return ret @@ -757,12 +761,12 @@ $$""".format(name, options), name, password, password) logger.exception('Could not remove data directory %s', self.data_dir) self.move_data_directory() - def basebackup(self, leader, env): + def basebackup(self, clone_member, env): # creates a replica data dir using pg_basebackup. # this is the default, built-in create_replica_method # tries twice, then returns failure (as 1) # uses "stream" as the xlog-method to avoid sync issues - master_connection = leader.conn_url + master_connection = clone_member.conn_url maxfailures = 2 ret = 1 for bbfailures in range(0, maxfailures): diff --git a/patroni/scripts/wale_restore.py b/patroni/scripts/wale_restore.py index c80cdfae..f3691907 100755 --- a/patroni/scripts/wale_restore.py +++ b/patroni/scripts/wale_restore.py @@ -154,7 +154,7 @@ def main(): args = parser.parse_args() # retry cloning in a loop - for retry in range(0, args.retries + 1): + for _ in range(0, args.retries + 1): restore = WALERestore(scope=args.scope, datadir=args.datadir, connstring=args.connstring, env_dir=args.envdir, threshold_mb=args.threshold_megabytes, threshold_pct=args.threshold_backup_size_percentage, use_iam=args.use_iam, diff --git a/tests/test_ctl.py b/tests/test_ctl.py index 3987619d..ed3aa1b9 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -12,6 +12,7 @@ from patroni.etcd import Etcd, Client from patroni.exceptions import PatroniCtlException from psycopg2 import OperationalError from test_etcd import etcd_read, etcd_write, requests_get, socket_getaddrinfo, MockResponse +from test_zookeeper import MockKazooClient 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 @@ -173,7 +174,11 @@ other y''') assert 'Failover failed' in result.output - def test_(self): + @patch('patroni.zookeeper.KazooClient', MockKazooClient) + @patch('requests.get', requests_get) + def test_get_dcs(self): + self.assertIsNotNone(get_dcs({'dcs': {'scheme': 'zookeeper', 'hostname': 'foo', 'port': 2181}}, 'dummy')) + self.assertIsNotNone(get_dcs({'dcs': {'scheme': 'exhibitor', 'hostname': 'exhibitor', 'port': 8181}}, 'dummy')) self.assertRaises(PatroniCtlException, get_dcs, {'scheme': 'dummy'}, 'dummy') @patch('psycopg2.connect', psycopg2_connect) diff --git a/tests/test_etcd.py b/tests/test_etcd.py index 601a0cd6..c14d3cac 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -7,8 +7,9 @@ import unittest from dns.exception import DNSException from mock import Mock, patch -from patroni.dcs import Cluster, DCSError, Leader +from patroni.dcs import Cluster from patroni.etcd import Client, Etcd, EtcdError +from patroni.exceptions import DCSError class MockResponse(object): @@ -229,11 +230,8 @@ class TestEtcd(unittest.TestCase): cluster = self.etcd.get_cluster() self.assertIsInstance(cluster, Cluster) self.assertIsNone(cluster.leader) - - def test_current_leader(self): - self.assertIsInstance(self.etcd.current_leader(), Leader) self.etcd._base_path = '/service/noleader' - self.assertIsNone(self.etcd.current_leader()) + self.assertRaises(EtcdError, self.etcd.get_cluster) def test_touch_member(self): self.assertFalse(self.etcd.touch_member('', '')) diff --git a/tests/test_ha.py b/tests/test_ha.py index 98b1ad1c..2f9e5478 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -29,12 +29,12 @@ def get_cluster_not_initialized_without_leader(): def get_cluster_initialized_without_leader(leader=False, failover=None): - m = Member(0, 'leader', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', - 'api_url': 'http://127.0.0.1:8008/patroni', 'xlog_location': 4}) - l = Leader(0, 0, m) if leader else None - o = Member(0, 'other', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5436/postgres', - 'api_url': 'http://127.0.0.1:8011/patroni'}) - return get_cluster(True, l, [m, o], failover) + m1 = Member(0, 'leader', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', + 'api_url': 'http://127.0.0.1:8008/patroni', 'xlog_location': 4}) + l = Leader(0, 0, m1) if leader else None + m2 = Member(0, 'other', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5436/postgres', + 'api_url': 'http://127.0.0.1:8011/patroni'}) + return get_cluster(True, l, [m1, m2], failover) def get_cluster_initialized_with_leader(failover=None): @@ -57,6 +57,7 @@ class MockPatroni(object): self.nap_time = 10 self.replicatefrom = None self.api.connection_string = 'http://127.0.0.1:8008' + self.clonefrom = None def run_async(func, args=()): @@ -87,7 +88,7 @@ class TestHa(unittest.TestCase): 'replication': {'username': '', 'password': '', 'network': ''}}) self.p.set_state('running') self.p.check_replication_lag = true - self.p.can_create_replica_without_leader = MagicMock(return_value=False) + self.p.can_create_replica_without_replication_connection = 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 @@ -205,13 +206,18 @@ class TestHa(unittest.TestCase): self.p.bootstrap = false self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap from leader') + def test_bootstrap_from_another_member(self): + self.ha.cluster = get_cluster_initialized_with_leader() + self.ha.patroni.clonefrom = 'other' + self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap from replica \'other\'') + def test_bootstrap_waiting_for_leader(self): self.ha.cluster = get_cluster_initialized_without_leader() self.assertEquals(self.ha.bootstrap(), 'waiting for leader to bootstrap') def test_bootstrap_without_leader(self): self.ha.cluster = get_cluster_initialized_without_leader() - self.p.can_create_replica_without_leader = MagicMock(return_value=True) + self.p.can_create_replica_without_replication_connection = MagicMock(return_value=True) self.assertEquals(self.ha.bootstrap(), "trying to bootstrap without leader") def test_bootstrap_initialize_lock_failed(self): diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 5558437b..fdd321df 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -33,7 +33,7 @@ class MockCursor(object): elif sql == 'SELECT pg_is_in_recovery()': self.results = [(False, )] elif sql.startswith('SELECT to_char(pg_postmaster_start_time'): - self.results = [('', True, '', '', '', False)] + self.results = [('', True, '', '', '', '', False)] else: self.results = [( None, @@ -471,17 +471,17 @@ class TestPostgresql(unittest.TestCase): def test_restore_configuration_files(self): self.p.restore_configuration_files() - def test_can_create_replica_without_leader(self): + def test_can_create_replica_without_replication_connection(self): self.p.config['create_replica_method'] = [] - self.assertFalse(self.p.can_create_replica_without_leader()) + self.assertFalse(self.p.can_create_replica_without_replication_connection()) self.p.config['create_replica_method'] = ['wale', 'basebackup'] self.p.config['wale'] = {'command': 'foo', 'no_master': 1} - self.assertTrue(self.p.can_create_replica_without_leader()) + self.assertTrue(self.p.can_create_replica_without_replication_connection()) - def test_replica_method_can_work_without_leader(self): - self.assertFalse(self.p.replica_method_can_work_without_leader('basebackup')) - self.assertFalse(self.p.replica_method_can_work_without_leader('foobar')) + def test_replica_method_can_work_without_replication_connection(self): + self.assertFalse(self.p.replica_method_can_work_without_replication_connection('basebackup')) + self.assertFalse(self.p.replica_method_can_work_without_replication_connection('foobar')) self.p.config['foo'] = {'command': 'bar', 'no_master': 1} - self.assertTrue(self.p.replica_method_can_work_without_leader('foo')) + self.assertTrue(self.p.replica_method_can_work_without_replication_connection('foo')) self.p.config['foo'] = {'command': 'bar'} - self.assertFalse(self.p.replica_method_can_work_without_leader('foo')) + self.assertFalse(self.p.replica_method_can_work_without_replication_connection('foo'))