diff --git a/patroni/ctl.py b/patroni/ctl.py index 1e1824f3..047fc4bb 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -34,7 +34,6 @@ except ImportError: # pragma: no cover from .dcs import get_dcs as _get_dcs from .exceptions import PatroniException -from .postgresql import Postgresql from .postgresql.misc import postgres_version_to_int from .utils import cluster_as_json, find_executable, patch_config, polling_loop, is_standby_cluster from .request import PatroniRequest @@ -949,62 +948,6 @@ def timestamp(precision=6): return datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S.%f')[:precision - 7] -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(data, 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') -@option_citus_group -@click.option('--sysid', '-s', help='System ID of the cluster to put into the initialize key', default="") -@click.pass_obj -def scaffold(obj, cluster_name, group, sysid): - dcs = get_dcs(obj, cluster_name, group) - cluster = dcs.get_cluster() - 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(obj, cluster_name) - - # make sure the leader keys will never expire - if not (touch_member(obj, 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='Discard scheduled events') @click.argument('cluster_name') @option_citus_group diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index d14a73c1..ca5e05c2 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -921,11 +921,9 @@ class AbstractDCS(object): return ret @abc.abstractmethod - def attempt_to_acquire_leader(self, permanent=False): + def attempt_to_acquire_leader(self): """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 primary :returns: `!True` if key has been created successfully. Key must be created atomically. In case if key already exists it should not be @@ -955,15 +953,13 @@ class AbstractDCS(object): """Create or update `/config` key""" @abc.abstractmethod - def touch_member(self, data, permanent=False): + def touch_member(self, data): """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: 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 primary :returns: `!True` on success otherwise `!False` """ diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index 310553e3..0f657467 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -415,12 +415,12 @@ class Consul(AbstractDCS): raise ConsulError('Consul is not responding properly') @catch_consul_errors - def touch_member(self, data, permanent=False): + def touch_member(self, data): cluster = self.cluster member = cluster and cluster.get_member(self._name, fallback_to_leader=False) try: - create_member = not permanent and self.refresh_session() + create_member = self.refresh_session() except DCSError: return False @@ -438,8 +438,7 @@ class Consul(AbstractDCS): return True try: - args = {} if permanent else {'acquire': self._session} - self._client.kv.put(self.member_path, json.dumps(data, separators=(',', ':')), **args) + self._client.kv.put(self.member_path, json.dumps(data, separators=(',', ':')), acquire=self._session) return True except InvalidSession: self._session = None @@ -524,10 +523,9 @@ class Consul(AbstractDCS): ): return self._update_service(new_data) - def _do_attempt_to_acquire_leader(self, permanent, retry): + def _do_attempt_to_acquire_leader(self, retry): try: - kwargs = {} if permanent else {'acquire': self._session} - return retry(self._client.kv.put, self.leader_path, self._name, **kwargs) + return retry(self._client.kv.put, self.leader_path, self._name, acquire=self._session) except InvalidSession: logger.error('Our session disappeared from Consul. Will try to get a new one and retry attempt') self._session = None @@ -542,16 +540,15 @@ class Consul(AbstractDCS): return retry(self._client.kv.put, self.leader_path, self._name, acquire=self._session) @catch_return_false_exception - def attempt_to_acquire_leader(self, permanent=False): + def attempt_to_acquire_leader(self): retry = self._retry.copy() - if not permanent: - self._run_and_handle_exceptions(self._do_refresh_session, retry=retry) + self._run_and_handle_exceptions(self._do_refresh_session, retry=retry) - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: - raise ConsulError('attempt_to_acquire_leader timeout') + retry.deadline = retry.stoptime - time.time() + if retry.deadline < 1: + raise ConsulError('attempt_to_acquire_leader timeout') - ret = self._run_and_handle_exceptions(self._do_attempt_to_acquire_leader, permanent, retry, retry=None) + ret = self._run_and_handle_exceptions(self._do_attempt_to_acquire_leader, retry, retry=None) if not ret: logger.info('Could not take out TTL lock') diff --git a/patroni/dcs/etcd.py b/patroni/dcs/etcd.py index 9908b863..604c73ba 100644 --- a/patroni/dcs/etcd.py +++ b/patroni/dcs/etcd.py @@ -694,28 +694,24 @@ class Etcd(AbstractEtcd): return cluster @catch_etcd_errors - def touch_member(self, data, permanent=False): + def touch_member(self, data): data = json.dumps(data, separators=(',', ':')) - return self._client.set(self.member_path, data, None if permanent else self._ttl) + return self._client.set(self.member_path, data, self._ttl) @catch_etcd_errors def take_leader(self): return self.retry(self._client.write, self.leader_path, self._name, ttl=self._ttl) - def _do_attempt_to_acquire_leader(self, permanent=False): + def _do_attempt_to_acquire_leader(self): try: - return bool(self.retry(self._client.write, - self.leader_path, - self._name, - ttl=None if permanent else self._ttl, - prevExist=False)) + return bool(self.retry(self._client.write, self.leader_path, self._name, ttl=self._ttl, prevExist=False)) except etcd.EtcdAlreadyExist: logger.info('Could not take out TTL lock') return False @catch_return_false_exception - def attempt_to_acquire_leader(self, permanent=False): - return self._run_and_handle_exceptions(self._do_attempt_to_acquire_leader, permanent=permanent, retry=None) + def attempt_to_acquire_leader(self): + return self._run_and_handle_exceptions(self._do_attempt_to_acquire_leader, retry=None) @catch_etcd_errors def set_failover_value(self, value, index=None): diff --git a/patroni/dcs/etcd3.py b/patroni/dcs/etcd3.py index 11e5763e..bf9ea156 100644 --- a/patroni/dcs/etcd3.py +++ b/patroni/dcs/etcd3.py @@ -743,12 +743,11 @@ class Etcd3(AbstractEtcd): return cluster @catch_etcd_errors - def touch_member(self, data, permanent=False): - if not permanent: - try: - self.refresh_lease() - except Etcd3Error: - return False + def touch_member(self, data): + try: + self.refresh_lease() + except Etcd3Error: + return False cluster = self.cluster member = cluster and cluster.get_member(self._name, fallback_to_leader=False) @@ -758,7 +757,7 @@ class Etcd3(AbstractEtcd): data = json.dumps(data, separators=(',', ':')) try: - return self._client.put(self.member_path, data, None if permanent else self._lease) + return self._client.put(self.member_path, data, self._lease) except LeaseNotFound: self._lease = None logger.error('Our lease disappeared from Etcd, can not "touch_member"') @@ -767,13 +766,13 @@ class Etcd3(AbstractEtcd): def take_leader(self): return self.retry(self._client.put, self.leader_path, self._name, self._lease) - def _do_attempt_to_acquire_leader(self, permanent, retry): + def _do_attempt_to_acquire_leader(self, retry): def _retry(*args, **kwargs): kwargs['retry'] = retry return retry(*args, **kwargs) try: - return _retry(self._client.put, self.leader_path, self._name, None if permanent else self._lease, 0) + return _retry(self._client.put, self.leader_path, self._name, self._lease, 0) except LeaseNotFound: logger.error('Our lease disappeared from Etcd. Will try to get a new one and retry attempt') self._lease = None @@ -785,24 +784,23 @@ class Etcd3(AbstractEtcd): if retry.deadline < 1: raise Etcd3Error('_do_attempt_to_acquire_leader timeout') - return _retry(self._client.put, self.leader_path, self._name, None if permanent else self._lease, 0) + return _retry(self._client.put, self.leader_path, self._name, self._lease, 0) @catch_return_false_exception - def attempt_to_acquire_leader(self, permanent=False): + def attempt_to_acquire_leader(self): retry = self._retry.copy() def _retry(*args, **kwargs): kwargs['retry'] = retry return retry(*args, **kwargs) - if not permanent: - self._run_and_handle_exceptions(self._do_refresh_lease, retry=_retry) + self._run_and_handle_exceptions(self._do_refresh_lease, retry=_retry) - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: - raise Etcd3Error('attempt_to_acquire_leader timeout') + retry.deadline = retry.stoptime - time.time() + if retry.deadline < 1: + raise Etcd3Error('attempt_to_acquire_leader timeout') - ret = self._run_and_handle_exceptions(self._do_attempt_to_acquire_leader, permanent, retry, retry=None) + ret = self._run_and_handle_exceptions(self._do_attempt_to_acquire_leader, retry, retry=None) if not ret: logger.info('Could not take out TTL lock') return ret diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index 03505cde..736b93aa 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -8,7 +8,6 @@ import os import random import socket import six -import sys import tempfile import time import urllib3 @@ -1132,9 +1131,9 @@ class Kubernetes(AbstractDCS): resource_version = kind and kind.metadata.resource_version return self._update_leader_with_retry(annotations, resource_version, self.__ips) - def attempt_to_acquire_leader(self, permanent=False): + def attempt_to_acquire_leader(self): now = datetime.datetime.now(tzutc).isoformat() - annotations = {self._LEADER: self._name, 'ttl': str(sys.maxsize if permanent else self._ttl), + annotations = {self._LEADER: self._name, 'ttl': str(self._ttl), 'renewTime': now, 'acquireTime': now, 'transitions': '0'} if self._leader_observed_record: try: @@ -1186,7 +1185,7 @@ class Kubernetes(AbstractDCS): return self.patch_or_create_config({self._CONFIG: value}, index, bool(self._config_resource_version), False) @catch_kubernetes_errors - def touch_member(self, data, permanent=False): + def touch_member(self, data): cluster = self.cluster if cluster and cluster.leader and cluster.leader.name == self._name: role = 'master' diff --git a/patroni/dcs/raft.py b/patroni/dcs/raft.py index cc60fb60..0d0e4568 100644 --- a/patroni/dcs/raft.py +++ b/patroni/dcs/raft.py @@ -415,9 +415,8 @@ class Raft(AbstractDCS): ret = self.attempt_to_acquire_leader() return ret - def attempt_to_acquire_leader(self, permanent=False): - return self._sync_obj.set(self.leader_path, self._name, ttl=None if permanent else self._ttl, - handle_raft_error=False, prevExist=False) + def attempt_to_acquire_leader(self): + return self._sync_obj.set(self.leader_path, self._name, ttl=self._ttl, handle_raft_error=False, prevExist=False) def set_failover_value(self, value, index=None): return self._sync_obj.set(self.failover_path, value, prevIndex=index) @@ -425,9 +424,9 @@ class Raft(AbstractDCS): def set_config_value(self, value, index=None): return self._sync_obj.set(self.config_path, value, prevIndex=index) - def touch_member(self, data, permanent=False): + def touch_member(self, data): data = json.dumps(data, separators=(',', ':')) - return self._sync_obj.set(self.member_path, data, None if permanent else self._ttl, timeout=2) + return self._sync_obj.set(self.member_path, data, self._ttl, timeout=2) def take_leader(self): return self._sync_obj.set(self.leader_path, self._name, ttl=self._ttl) diff --git a/patroni/dcs/zookeeper.py b/patroni/dcs/zookeeper.py index ab4ca5f5..476b1c51 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -340,10 +340,10 @@ class ZooKeeper(AbstractDCS): logger.exception('Failed to create %s', path) return False - def attempt_to_acquire_leader(self, permanent=False): + def attempt_to_acquire_leader(self): try: self._client.retry(self._client.create, self.leader_path, self._name.encode('utf-8'), - makepath=True, ephemeral=not permanent) + makepath=True, ephemeral=True) return True except (ConnectionClosedError, RetryFailedError) as e: raise ZooKeeperError(e) @@ -383,7 +383,7 @@ class ZooKeeper(AbstractDCS): return self._create(self.initialize_path, sysid, retry=True) if create_new \ else self._client.retry(self._client.set, self.initialize_path, sysid) - def touch_member(self, data, permanent=False): + def touch_member(self, data): cluster = self.cluster member = cluster and cluster.get_member(self._name, fallback_to_leader=False) member_data = self.__last_member_data or member and member.data @@ -408,8 +408,7 @@ class ZooKeeper(AbstractDCS): return True else: try: - self._client.create_async(self.member_path, encoded_data, makepath=True, - ephemeral=not permanent).get(timeout=1) + self._client.create_async(self.member_path, encoded_data, makepath=True, ephemeral=True).get(timeout=1) self.__last_member_data = data return True except Exception as e: diff --git a/tests/test_ctl.py b/tests/test_ctl.py index decd2779..55964c9f 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -403,30 +403,6 @@ class TestCtl(unittest.TestCase): with patch('patroni.ctl.load_config', Mock(return_value={})): self.runner.invoke(ctl, ['list']) - @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) - mock_get_dcs.return_value.delete_cluster = Mock() - - 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 diff --git a/tests/test_etcd.py b/tests/test_etcd.py index 898bdd96..7e91b9e9 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -277,7 +277,7 @@ class TestEtcd(unittest.TestCase): self.assertIsInstance(cluster.workers[1], Cluster) def test_touch_member(self): - self.assertFalse(self.etcd.touch_member('', '')) + self.assertFalse(self.etcd.touch_member('')) def test_take_leader(self): self.assertFalse(self.etcd.take_leader())