Merge pull request #252 from zalando/feature/ctl_scaffolding

Add patronictl scaffold command.
This commit is contained in:
Oleksii Kliukin
2016-08-10 12:21:12 +02:00
committed by GitHub
6 changed files with 107 additions and 17 deletions
+57 -2
View File
@@ -21,6 +21,7 @@ from click import ClickException
from patroni.config import Config
from patroni.dcs import get_dcs as _get_dcs
from patroni.exceptions import PatroniException
from patroni.postgresql import Postgresql
from patroni.utils import is_valid_pg_version
from prettytable import PrettyTable
from six.moves.urllib_parse import urlparse
@@ -76,7 +77,6 @@ def load_config(path, dcs):
for d in DCS_DEFAULTS:
config.pop(d, None)
config.update(dcs)
return config
@@ -106,7 +106,7 @@ def ctl(ctx):
def get_dcs(config, scope):
config.setdefault('scope', scope)
config['scope'] = scope
config.setdefault('name', scope)
try:
return _get_dcs(config)
@@ -725,6 +725,61 @@ def configure(config_file, dcs, namespace):
store_config(config, config_file)
def touch_member(config, dcs):
''' Rip-off of the ha.touch_member without inter-class dependencies '''
p = Postgresql(config['postgresql'])
p.set_state('running')
p.set_role('master')
def restapi_connection_string(config):
protocol = 'https' if config.get('certfile') else 'http'
connect_address = config.get('connect_address')
listen = config['listen']
return '{0}://{1}/patroni'.format(protocol, connect_address or listen)
data = {
'conn_url': p.connection_string,
'api_url': restapi_connection_string(config['restapi']),
'state': p.state,
'role': p.role
}
return dcs.touch_member(json.dumps(data, separators=(',', ':')), permanent=True)
def set_defaults(config, cluster_name):
''' fill-in some basic configuration parameters if config file is not set '''
config['postgresql'].setdefault('name', cluster_name)
config['postgresql'].setdefault('scope', cluster_name)
config['postgresql'].setdefault('listen', '127.0.0.1')
config['postgresql']['authentication'] = {'replication': None}
config['restapi']['listen'] = ':' in config['restapi']['listen'] and config['restapi']['listen'] or '127.0.0.1:8008'
@ctl.command('scaffold', help='Create a structure for the cluster in DCS')
@click.argument('cluster_name')
@click.option('--sysid', '-s', help='System ID of the cluster to put into the initialize key', default="")
@option_config_file
@option_dcs
def scaffold(cluster_name, config_file, dcs, sysid):
config, dcs, cluster = ctl_load_config(cluster_name, config_file, dcs)
if cluster and cluster.initialize is not None:
raise PatroniCtlException("This cluster is already initialized")
if not dcs.initialize(create_new=True, sysid=sysid):
# initialize key already exists, don't touch this cluster
raise PatroniCtlException("Initialize key for cluster {0} already exists".format(cluster_name))
set_defaults(config, cluster_name)
# make sure the leader keys will never expire
if not (touch_member(config, dcs) and dcs.attempt_to_acquire_leader(permanent=True)):
# we did initialize this cluster, but failed to write the leader or member keys, wipe it down completely.
dcs.delete_cluster()
raise PatroniCtlException("Unable to install permanent leader for cluster {0}".format(cluster_name))
click.echo("Cluster {0} has been created successfully".format(cluster_name))
@ctl.command('flush', help='Flush scheduled events')
@click.argument('cluster_name')
@click.argument('member_names', nargs=-1)
+6 -2
View File
@@ -349,9 +349,11 @@ class AbstractDCS(object):
for example for etcd `prevValue` parameter must be used."""
@abc.abstractmethod
def attempt_to_acquire_leader(self):
def attempt_to_acquire_leader(self, permanent=False):
"""Attempt to acquire leader lock
This method should create `/leader` key with value=`~self._name`
:param permanent: if set to `!True`, the leader key will never expire.
Used in patronictl for the external master
:returns: `!True` if key has been created successfully.
Key must be created atomically. In case if key already exists it should not be
@@ -379,13 +381,15 @@ class AbstractDCS(object):
"""Create or update `/config` key"""
@abc.abstractmethod
def touch_member(self, data, ttl=None):
def touch_member(self, data, ttl=None, permanent=False):
"""Update member key in DCS.
This method should create or update key with the name = '/members/' + `~self._name`
and value = data in a given DCS.
:param data: json serialized information about instance (including connection strings)
:param ttl: ttl for member key, optional parameter. If it is None `~self.member_ttl will be used`
:param permanent: if set to `!True`, the member key will never expire.
Used in patronictl for the external master.
:returns: `!True` on success otherwise `!False`
"""
+5 -3
View File
@@ -191,7 +191,8 @@ class Consul(AbstractDCS):
return True
try:
self._client.kv.put(self.member_path, data, acquire=self._session)
args = {} if kwargs.get('permanent', False) else {'acquire': self._session}
self._client.kv.put(self.member_path, data, **args)
self._my_member_data = data
return True
except Exception:
@@ -199,8 +200,9 @@ class Consul(AbstractDCS):
return False
@catch_consul_errors
def attempt_to_acquire_leader(self):
ret = self._client.kv.put(self.leader_path, self._name, acquire=self._session)
def attempt_to_acquire_leader(self, permanent=False):
args = {} if permanent else {'acquire': self._session}
ret = self._client.kv.put(self.leader_path, self._name, **args)
if not ret:
logger.info('Could not take out TTL lock')
return ret
+8 -4
View File
@@ -307,16 +307,20 @@ class Etcd(AbstractDCS):
raise EtcdError('Etcd is not responding properly')
@catch_etcd_errors
def touch_member(self, data, ttl=None):
return self.retry(self._client.set, self.member_path, data, ttl or self._ttl)
def touch_member(self, data, ttl=None, permanent=False):
return self.retry(self._client.set, self.member_path, data, None if permanent else ttl or self._ttl)
@catch_etcd_errors
def take_leader(self):
return self.retry(self._client.set, self.leader_path, self._name, self._ttl)
def attempt_to_acquire_leader(self):
def attempt_to_acquire_leader(self, permanent=False):
try:
return bool(self.retry(self._client.write, self.leader_path, self._name, ttl=self._ttl, prevExist=False))
return bool(self.retry(self._client.write,
self.leader_path,
self._name,
ttl=None if permanent else self._ttl,
prevExist=False))
except etcd.EtcdAlreadyExist:
logger.info('Could not take out TTL lock')
except (RetryFailedError, etcd.EtcdException):
+4 -4
View File
@@ -200,8 +200,8 @@ class ZooKeeper(AbstractDCS):
except:
return False
def attempt_to_acquire_leader(self):
ret = self._create(self.leader_path, self._name, makepath=True, ephemeral=True)
def attempt_to_acquire_leader(self, permanent=False):
ret = self._create(self.leader_path, self._name, makepath=True, ephemeral=not permanent)
if not ret:
logger.info('Could not take out TTL lock')
return ret
@@ -230,7 +230,7 @@ class ZooKeeper(AbstractDCS):
return self._create(self.initialize_path, sysid, makepath=True) if create_new \
else self._client.retry(self._client.set, self.initialize_path, sysid.encode("utf-8"))
def touch_member(self, data, ttl=None):
def touch_member(self, data, ttl=None, permanent=False):
cluster = self.cluster
member = cluster and ([m for m in cluster.members if m.name == self._name] or [None])[0]
data = data.encode('utf-8')
@@ -248,7 +248,7 @@ class ZooKeeper(AbstractDCS):
return True
else:
try:
self._client.create_async(self.member_path, data, makepath=True, ephemeral=True).get(timeout=1)
self._client.create_async(self.member_path, data, makepath=True, ephemeral=not permanent).get(timeout=1)
self._my_member_data = data
return True
except Exception as e:
+27 -2
View File
@@ -12,7 +12,7 @@ from patroni.dcs.etcd import Client
from psycopg2 import OperationalError
from test_etcd import etcd_read, 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
get_cluster_initialized_with_only_leader, get_cluster_not_initialized_without_leader
from test_postgresql import MockConnect, psycopg2_connect
CONFIG_FILE_PATH = './test-ctl.yaml'
@@ -29,7 +29,9 @@ def test_rw_config():
os.rmdir(CONFIG_FILE_PATH)
@patch('patroni.ctl.load_config', Mock(return_value={'restapi': {'auth': 'u:p'}, 'etcd': {'host': 'localhost:4001'}}))
@patch('patroni.ctl.load_config', Mock(return_value={'postgresql': {'data_dir': '.', 'parameters': {}, 'retry_timeout': 5},
'restapi': {'auth': 'u:p', 'listen': ''},
'etcd': {'host': 'localhost:4001'}}))
class TestCtl(unittest.TestCase):
@patch('socket.getaddrinfo', socket_getaddrinfo)
@@ -345,6 +347,29 @@ class TestCtl(unittest.TestCase):
result = self.runner.invoke(configure, ['--dcs', 'abc', '-c', 'dummy', '-n', 'bla'])
assert result.exit_code == 0
@patch('patroni.ctl.get_dcs')
def test_scaffold(self, mock_get_dcs):
mock_get_dcs.return_value = self.e
mock_get_dcs.return_value.get_cluster = get_cluster_not_initialized_without_leader
mock_get_dcs.return_value.initialize = Mock(return_value=True)
mock_get_dcs.return_value.touch_member = Mock(return_value=True)
mock_get_dcs.return_value.attempt_to_acquire_leader = Mock(return_value=True)
with patch.object(self.e, 'initialize', return_value=False):
result = self.runner.invoke(ctl, ['scaffold', 'alpha'])
assert result.exception
with patch.object(mock_get_dcs.return_value, 'touch_member', Mock(return_value=False)):
result = self.runner.invoke(ctl, ['scaffold', 'alpha'])
assert result.exception
result = self.runner.invoke(ctl, ['scaffold', 'alpha'])
assert result.exit_code == 0
mock_get_dcs.return_value.get_cluster = get_cluster_initialized_with_leader
result = self.runner.invoke(ctl, ['scaffold', 'alpha'])
assert result.exception
@patch('patroni.ctl.get_dcs')
def test_list_extended(self, mock_get_dcs):
mock_get_dcs.return_value = self.e