diff --git a/patroni/ctl.py b/patroni/ctl.py index cb6978b9..b272d5ca 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -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) diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index 7b291d55..d7a69997 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -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` """ diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index ffa0ab12..36a93536 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -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 diff --git a/patroni/dcs/etcd.py b/patroni/dcs/etcd.py index 6f96beb4..d81cf968 100644 --- a/patroni/dcs/etcd.py +++ b/patroni/dcs/etcd.py @@ -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): diff --git a/patroni/dcs/zookeeper.py b/patroni/dcs/zookeeper.py index 9f8c9e12..68f106a3 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -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: diff --git a/tests/test_ctl.py b/tests/test_ctl.py index 0e4ea51b..0bf6eaf8 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -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