diff --git a/patroni/ctl.py b/patroni/ctl.py index a0f0c565..fae7cba0 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -20,8 +20,9 @@ import yaml 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 get_conn_kwargs +from patroni.exceptions import PatroniException, DCSError +from patroni.postgresql import Postgresql, get_conn_kwargs +from patroni.api import RestApiServer 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 @@ -646,3 +646,54 @@ def configure(config_file, dcs, namespace): config['dcs_api'] = str(dcs) config['namespace'] = str(namespace) store_config(config, config_file) + + +def touch_member(config, dcs): + ''' Rip of the ha.touch_member without inter-class dependencies ''' + p = Postgresql(config['postgresql']) + p.set_state('running') + p.set_role('master') + + api = RestApiServer(None, config['restapi']) + data = { + 'conn_url': p.connection_string, + 'api_url': api.connection_string, + 'state': p.state, + 'role': p.role + } + try: + dcs.touch_member(json.dumps(data, separators=(',', ':')), permanent=True) + except DCSError: + return False + return True + + +def set_defaults(config, cluster_name): + ''' fill-in some basic configuration parameters if config file is not set ''' + config['postgresql']['name'] = config['postgresql'].get('name') or cluster_name + config['postgresql']['scope'] = config['postgresql'].get('scope') or cluster_name + config['postgresql']['listen'] = config['postgresql'].get('listen') or "127.0.0.1" + config['postgresql']['authentication'] = {'replication': None} + config['restapi']['listen'] = ':' in config['restapi'].get('listen', ".") or '127.0.0.1:5432' + + +@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): + logging.debug("config_file = %s, cluster_name = %s, dcs = %s, sysid = %s", config_file, cluster_name, dcs, sysid) + config, dcs, cluster = ctl_load_config(cluster_name, config_file, dcs) + if cluster and cluster.initialize: + raise PatroniCtlException("This cluster is already initialized") + dcs.initialize(create_new=True, sysid=sysid) + + set_defaults(config, cluster_name) + + # make sure the leader key will never expire + if not (touch_member(config, dcs) and dcs.attempt_to_acquire_leader(permanent=True)): + dcs.delete_leader() + dcs.cancel_initialization() + return 1 + return 0 diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index 39049792..a2264a03 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -343,13 +343,14 @@ 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` + :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..53ab1b1b 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 = {'cas': None} if kwargs.get('permanent') else {} + self._client.kv.put(self.member_path, data, acquire=self._session, **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 = {'cas': None} if permanent else {} + ret = self._client.kv.put(self.leader_path, self._name, acquire=self._session, **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 c376ce8b..3f9f2c3c 100644 --- a/patroni/dcs/etcd.py +++ b/patroni/dcs/etcd.py @@ -281,16 +281,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, (ttl or self._ttl) if not permanent else None) @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=self._ttl if permanent else None, + 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 be6fe45a..3da9d1f1 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -162,8 +162,11 @@ 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 @@ -192,7 +195,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] path = self.member_path @@ -213,7 +216,7 @@ class ZooKeeper(AbstractDCS): if member: self._client.retry(self._client.set, path, data) else: - self._client.retry(self._client.create, path, data, makepath=True, ephemeral=True) + self._client.retry(self._client.create, path, data, makepath=True, ephemeral=not permanent) self._my_member_data = data return True except NodeExistsError: