From c2d46a084efcbf1bc3b1b48f6ba306f5c1df8e7e Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Fri, 5 Feb 2016 09:22:38 +0100 Subject: [PATCH 01/15] Include timestamp of last replayed location in api call. --- patroni/api.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/patroni/api.py b/patroni/api.py index 2318fbf0..3137a829 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -229,6 +229,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, @@ -238,7 +239,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] }) } From de129b733d27f2704371f054eb90350ba17adb30 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 17 Feb 2016 12:46:32 +0100 Subject: [PATCH 02/15] Fix unit tests --- tests/test_postgresql.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 4c701c60..a63a90cc 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -38,7 +38,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, From 2f84e9f4ec38bec637b86b5b229d52c163658346 Mon Sep 17 00:00:00 2001 From: Oleksandr Shulgin Date: Tue, 23 Feb 2016 16:44:09 +0100 Subject: [PATCH 03/15] Add support for ZooKeeper/Exhibitor DCS URI in patronictl ... -d --- patroni/ctl.py | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/patroni/ctl.py b/patroni/ctl.py index fb4c241b..48a0c413 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -16,6 +16,7 @@ from six.moves.urllib_parse import urlparse import logging from .etcd import Etcd +from .zookeeper import ZooKeeper from .exceptions import PatroniCtlException from .postgresql import parseurl @@ -44,12 +45,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)} @@ -103,6 +104,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') From a1417875a2e2e223f3b8a309c4e19047db644a93 Mon Sep 17 00:00:00 2001 From: Oleksandr Shulgin Date: Tue, 23 Feb 2016 17:04:16 +0100 Subject: [PATCH 04/15] Add dummy patronictl tests with ZooKeeper --- tests/test_ctl.py | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index fb63aebd..e0e5fc99 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -16,10 +16,12 @@ from patroni.ctl import ctl, members, store_config, load_config, output_members, 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 patroni.zookeeper import ZooKeeper 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 test_etcd import etcd_read, etcd_write, requests_get, socket_getaddrinfo, MockResponse +from test_zookeeper import MockKazooClient from test_postgresql import MockConnect, psycopg2_connect CONFIG_FILE_PATH = './test-ctl.yaml' @@ -52,6 +54,7 @@ def test_rw_config(): class TestCtl(unittest.TestCase): @patch('socket.getaddrinfo', socket_getaddrinfo) + @patch('patroni.zookeeper.KazooClient', MockKazooClient) def setUp(self): self.runner = CliRunner() with patch.object(Client, 'machines') as mock_machines: @@ -61,6 +64,7 @@ class TestCtl(unittest.TestCase): self.e.client.read = etcd_read self.e.client.write = etcd_write self.e.client.delete = Mock(side_effect=etcd.EtcdException()) + self.zk = ZooKeeper('foo', {'ttl': 30, 'hosts': ['ok:2181'], 'scope': 'test'}) self.ha = Ha(MockPatroni(self.p, self.e)) self.ha._async_executor.run_async = run_async self.ha.old_cluster = self.e.get_cluster() @@ -377,3 +381,13 @@ leader''') ]) assert result.exit_code == 0 + + @patch('patroni.ctl.load_config', Mock(return_value={'dcs': {'scheme': 'zookeeper', 'hostname': 'localhost', 'port': 2181}})) + def test_zookeeper(self): + with patch('patroni.ctl.get_dcs', Mock(return_value=self.zk)): + self.runner.invoke(ctl, ['list']) + + @patch('patroni.ctl.load_config', Mock(return_value={'dcs': {'scheme': 'exhibitor', 'hostname': 'localhost', 'port': 8181}})) + def test_exhibitor(self): + with patch('patroni.ctl.get_dcs', Mock(return_value=self.zk)): + self.runner.invoke(ctl, ['list']) From 524cfafbbe07a1fe6a471245e5255edf803be7b5 Mon Sep 17 00:00:00 2001 From: Oleksandr Shulgin Date: Tue, 23 Feb 2016 17:18:16 +0100 Subject: [PATCH 05/15] Don't mock get_dcs() for ZK, we are trying to test it --- tests/test_ctl.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index 55612a94..f8f2783a 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -390,10 +390,8 @@ leader''') @patch('patroni.ctl.load_config', Mock(return_value={'dcs': {'scheme': 'zookeeper', 'hostname': 'localhost', 'port': 2181}})) def test_zookeeper(self): - with patch('patroni.ctl.get_dcs', Mock(return_value=self.zk)): - self.runner.invoke(ctl, ['list']) + self.runner.invoke(ctl, ['list']) @patch('patroni.ctl.load_config', Mock(return_value={'dcs': {'scheme': 'exhibitor', 'hostname': 'localhost', 'port': 8181}})) def test_exhibitor(self): - with patch('patroni.ctl.get_dcs', Mock(return_value=self.zk)): - self.runner.invoke(ctl, ['list']) + self.runner.invoke(ctl, ['list']) From 16b321e0a53494d70e5c606f023d78aef02c9d94 Mon Sep 17 00:00:00 2001 From: Oleksandr Shulgin Date: Tue, 23 Feb 2016 17:33:15 +0100 Subject: [PATCH 06/15] Add dummy cluster name to test_ctl / zookeeper --- tests/test_ctl.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index f8f2783a..4092a669 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -390,8 +390,8 @@ leader''') @patch('patroni.ctl.load_config', Mock(return_value={'dcs': {'scheme': 'zookeeper', 'hostname': 'localhost', 'port': 2181}})) def test_zookeeper(self): - self.runner.invoke(ctl, ['list']) + self.runner.invoke(ctl, ['list', 'foo']) @patch('patroni.ctl.load_config', Mock(return_value={'dcs': {'scheme': 'exhibitor', 'hostname': 'localhost', 'port': 8181}})) def test_exhibitor(self): - self.runner.invoke(ctl, ['list']) + self.runner.invoke(ctl, ['list', 'foo']) From 9c12eb671da4c83e7f3f2dae35bb2eed2b644741 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 24 Feb 2016 12:10:58 +0100 Subject: [PATCH 07/15] Fix unit tests --- tests/test_ctl.py | 17 +++++------------ 1 file changed, 5 insertions(+), 12 deletions(-) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index 4092a669..ed3aa1b9 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -9,7 +9,6 @@ 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.etcd import Etcd, Client -from patroni.zookeeper import ZooKeeper from patroni.exceptions import PatroniCtlException from psycopg2 import OperationalError from test_etcd import etcd_read, etcd_write, requests_get, socket_getaddrinfo, MockResponse @@ -48,7 +47,6 @@ def test_rw_config(): class TestCtl(unittest.TestCase): @patch('socket.getaddrinfo', socket_getaddrinfo) - @patch('patroni.zookeeper.KazooClient', MockKazooClient) def setUp(self): self.runner = CliRunner() with patch.object(Client, 'machines') as mock_machines: @@ -57,7 +55,6 @@ class TestCtl(unittest.TestCase): self.e.client.read = etcd_read self.e.client.write = etcd_write self.e.client.delete = Mock(side_effect=EtcdException) - self.zk = ZooKeeper('foo', {'ttl': 30, 'hosts': ['ok:2181'], 'scope': 'test'}) @patch('psycopg2.connect', psycopg2_connect) def test_get_cursor(self): @@ -177,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) @@ -387,11 +388,3 @@ leader''') ]) assert result.exit_code == 0 - - @patch('patroni.ctl.load_config', Mock(return_value={'dcs': {'scheme': 'zookeeper', 'hostname': 'localhost', 'port': 2181}})) - def test_zookeeper(self): - self.runner.invoke(ctl, ['list', 'foo']) - - @patch('patroni.ctl.load_config', Mock(return_value={'dcs': {'scheme': 'exhibitor', 'hostname': 'localhost', 'port': 8181}})) - def test_exhibitor(self): - self.runner.invoke(ctl, ['list', 'foo']) From cb38e50ac1d90b860767e6e914fe6a61270f0f31 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Fri, 26 Feb 2016 08:50:53 +0100 Subject: [PATCH 08/15] Remove unused code --- patroni/dcs.py | 8 -------- tests/test_etcd.py | 8 +++----- 2 files changed, 3 insertions(+), 13 deletions(-) diff --git a/patroni/dcs.py b/patroni/dcs.py index 4c2dddad..8cda0620 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 @@ -269,13 +268,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/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('', '')) From 9057ddeb7c80fb2fd54af83a09d62018beabc137 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Thu, 10 Mar 2016 16:06:31 +0100 Subject: [PATCH 09/15] First implementation of cloning from the replica. At the moment we just replace the master with the node at the 'clonefrom' tag if it's present. Master should be available anyway, otherwise, it will not even try to do cloning. Acceptance tests: https://github.com/zalando/patroni/pull/144/commits --- patroni/__init__.py | 4 ++++ patroni/dcs.py | 7 +++++++ patroni/ha.py | 4 +++- patroni/postgresql.py | 13 +++++++------ tests/test_ha.py | 1 + 5 files changed, 22 insertions(+), 7 deletions(-) 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/dcs.py b/patroni/dcs.py index 4c2dddad..b705105d 100644 --- a/patroni/dcs.py +++ b/patroni/dcs.py @@ -72,6 +72,10 @@ class Member(namedtuple('Member', 'index,name,session,data')): def replicatefrom(self): return self.data.get('tags', {}).get('replicatefrom') + @property + def clonefrom(self): + return self.data.get('tags', {}).get('clonefrom') + class Leader(namedtuple('Leader', 'index,session,member')): @@ -147,6 +151,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): diff --git a/patroni/ha.py b/patroni/ha.py index 98ba4398..61e27814 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -76,7 +76,9 @@ class Ha(object): 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, )) + clonefrom = self.patroni.clonefrom + source = self.cluster.get_member(clonefrom) if self.cluster.has_member(clonefrom) else self.cluster.leader + self._async_executor.run_async(self.clone, args=(source,)) return 'trying to bootstrap from leader' elif not self.cluster.initialize and not self.patroni.nofailover: # no initialize key if self.dcs.initialize(create_new=True): # race for initialization diff --git a/patroni/postgresql.py b/patroni/postgresql.py index a1cedb58..da561f8a 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -234,6 +234,7 @@ class Postgresql(object): return env def sync_replica(self, leader): + # add either the leader's or replica's credentials to pgpass env = self.write_pgpass(parseurl(leader.conn_url)) if leader else os.environ.copy() if self.create_replica(leader, env) == 0: self.delete_trigger_file() @@ -258,23 +259,23 @@ class Postgresql(object): 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) - def create_replica(self, leader, env): + def create_replica(self, source, 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 = source.conn_url if source 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 \ + # 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_leader(r)] if not source \ 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(source, env) if ret == 0: logger.info("replica has been created using basebackup") # if basebackup succeeds, exit with success @@ -702,7 +703,7 @@ $$""".format(name, options), name, password, password) """ 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 master or replica - initialize the replica using the replica creation method that works without the master (i.e. restore from on-disk base backup) diff --git a/tests/test_ha.py b/tests/test_ha.py index 98b1ad1c..2b5e0ada 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -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=()): From 805716ed689268f6be771cddf59d576434ad0994 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Fri, 11 Mar 2016 10:19:00 +0100 Subject: [PATCH 10/15] Variables and parameters renaming. Previously, "without_leader" suffix was used in the name of methods and functions that initialize a replica without an active replication connection, and leader was part of the name for parameters and messages that require an active replication conneciton. Since we support init from the members other than the leader, those conventions have to be changed. --- patroni/ha.py | 22 ++++++++++-------- patroni/postgresql.py | 49 +++++++++++++++++++++------------------- tests/test_ha.py | 4 ++-- tests/test_postgresql.py | 16 ++++++------- 4 files changed, 49 insertions(+), 42 deletions(-) diff --git a/patroni/ha.py b/patroni/ha.py index 61e27814..f9995d31 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -65,21 +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') clonefrom = self.patroni.clonefrom - source = self.cluster.get_member(clonefrom) if self.cluster.has_member(clonefrom) else self.cluster.leader - self._async_executor.run_async(self.clone, args=(source,)) - return 'trying to bootstrap from leader' + 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: @@ -98,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' diff --git a/patroni/postgresql.py b/patroni/postgresql.py index da561f8a..d8f649c7 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -233,10 +233,10 @@ class Postgresql(object): env['PGPASSFILE'] = self.pgpass return env - def sync_replica(self, leader): - # add either the leader's or replica's credentials to pgpass - 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 @@ -249,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, source, 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 = source.conn_url if source 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 source, 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 source \ - else replica_methods + 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(source, env) + ret = self.basebackup(clone_member, env) if ret == 0: logger.info("replica has been created using basebackup") # if basebackup succeeds, exit with success @@ -699,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 or replica + - 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). @@ -720,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() @@ -728,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 @@ -758,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/tests/test_ha.py b/tests/test_ha.py index 2b5e0ada..7eaf3808 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -88,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 @@ -212,7 +212,7 @@ class TestHa(unittest.TestCase): 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..1f5380d9 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -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')) From d965d21ada4b4b6f21896a8bfe4553e1bf2765d7 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Fri, 11 Mar 2016 10:48:58 +0100 Subject: [PATCH 11/15] Unit-tests for clone from the replica. Remove clonefrom function from dcs, since it's not used. --- patroni/dcs.py | 4 ---- patroni/ha.py | 2 +- tests/test_ha.py | 17 +++++++++++------ 3 files changed, 12 insertions(+), 11 deletions(-) diff --git a/patroni/dcs.py b/patroni/dcs.py index b705105d..66e07aae 100644 --- a/patroni/dcs.py +++ b/patroni/dcs.py @@ -72,10 +72,6 @@ class Member(namedtuple('Member', 'index,name,session,data')): def replicatefrom(self): return self.data.get('tags', {}).get('replicatefrom') - @property - def clonefrom(self): - return self.data.get('tags', {}).get('clonefrom') - class Leader(namedtuple('Leader', 'index,session,member')): diff --git a/patroni/ha.py b/patroni/ha.py index f9995d31..f6669e9d 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -80,7 +80,7 @@ class Ha(object): 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) + 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) diff --git a/tests/test_ha.py b/tests/test_ha.py index 7eaf3808..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): @@ -206,6 +206,11 @@ 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') From d3c2b8b2aa28a09f4c2c30169ba061fae70268c7 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Fri, 11 Mar 2016 15:18:58 +0100 Subject: [PATCH 12/15] replace unused variables with _ --- patroni/api.py | 4 ++-- patroni/ctl.py | 2 +- patroni/ha.py | 4 ++-- patroni/scripts/wale_restore.py | 2 +- 4 files changed, 6 insertions(+), 6 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index 1186321f..f2877b05 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' diff --git a/patroni/ctl.py b/patroni/ctl.py index c4cc9fe0..9d07267b 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -240,7 +240,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/ha.py b/patroni/ha.py index 98ba4398..d4697fe0 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -203,7 +203,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, _, xlog_location, 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 +223,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, _, xlog_location, 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/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, From 3319c3eeea7e3996b56d8e2d7eafc7ba2455674b Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Fri, 11 Mar 2016 15:36:24 +0100 Subject: [PATCH 13/15] replace unused variables with _ --- patroni/ha.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/patroni/ha.py b/patroni/ha.py index d4697fe0..5a7e50a9 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -203,7 +203,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, _, 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 +223,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, _, 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 From 9ca1b754a5d846e0772bec7ed4fe5c0280f5d05c Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Fri, 11 Mar 2016 16:33:26 +0100 Subject: [PATCH 14/15] Add zappr configuration. --- .zappr.yaml | 12 ++++++++++++ 1 file changed, 12 insertions(+) create mode 100644 .zappr.yaml diff --git a/.zappr.yaml b/.zappr.yaml new file mode 100644 index 00000000..4d628636 --- /dev/null +++ b/.zappr.yaml @@ -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 From 6c63d32253e81d7ef26b801f9f69610abe8e18b5 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Fri, 11 Mar 2016 16:39:57 +0100 Subject: [PATCH 15/15] zappr config must be .yml. --- .zappr.yaml => .zappr.yml | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename .zappr.yaml => .zappr.yml (100%) diff --git a/.zappr.yaml b/.zappr.yml similarity index 100% rename from .zappr.yaml rename to .zappr.yml