Merge branch 'master' of https://github.com/zalando/patroni into feature/acceptance_tests_behave

This commit is contained in:
Oleksii Kliukin
2016-03-11 16:58:56 +01:00
12 changed files with 110 additions and 71 deletions
+12
View File
@@ -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
+4
View File
@@ -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:
+5 -3
View File
@@ -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]
})
}
+10 -3
View File
@@ -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')
+3 -8
View File
@@ -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.
+16 -10
View File
@@ -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
+27 -23
View File
@@ -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):
+1 -1
View File
@@ -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,
+6 -1
View File
@@ -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)
+3 -5
View File
@@ -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('', ''))
+14 -8
View File
@@ -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:[email protected]: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:[email protected]: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:[email protected]: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:[email protected]: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):
+9 -9
View File
@@ -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'))