mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Remove patronictl scaffold (#2544)
The only reason for having it was a hacky way of running standby clusters.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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`
|
||||
"""
|
||||
|
||||
|
||||
+11
-14
@@ -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')
|
||||
|
||||
|
||||
+6
-10
@@ -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):
|
||||
|
||||
+15
-17
@@ -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
|
||||
|
||||
@@ -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'
|
||||
|
||||
+4
-5
@@ -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)
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
+1
-1
@@ -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())
|
||||
|
||||
Reference in New Issue
Block a user