mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-26 23:50:23 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
93efa91bbd | ||
|
|
c9f420fd4a | ||
|
|
ed0e308b9b | ||
|
|
313adb61ec | ||
|
|
db12051a5b | ||
|
|
c51234557d | ||
|
|
195b8bf049 | ||
|
|
ccfe30729e | ||
|
|
c81391e314 | ||
|
|
415180048a | ||
|
|
9288ce066b | ||
|
|
cb80f7ee06 |
@@ -59,7 +59,8 @@ Etcd
|
||||
- **PATRONI\_ETCD\_USE\_PROXIES**: If this parameter is set to true, Patroni will consider **hosts** as a list of proxies and will not perform a topology discovery of etcd cluster but stick to a fixed list of **hosts**.
|
||||
- **PATRONI\_ETCD\_PROTOCOL**: http or https, if not specified http is used. If the **url** or **proxy** is specified - will take protocol from them.
|
||||
- **PATRONI\_ETCD\_HOST**: the host:port for the etcd endpoint.
|
||||
- **PATRONI\_ETCD\_SRV**: Domain to search the SRV record(s) for cluster autodiscovery.
|
||||
- **PATRONI\_ETCD\_SRV**: Domain to search the SRV record(s) for cluster autodiscovery. Patroni will try to query these SRV service names for specified domain (in that order until first success): ``_etcd-client-ssl``, ``_etcd-client``, ``_etcd-ssl``, ``_etcd``, ``_etcd-server-ssl``, ``_etcd-server``. If SRV records for ``_etcd-server-ssl`` or ``_etcd-server`` are retrieved then ETCD peer protocol is used do query ETCD for available members. Otherwise hosts from SRV records will be used.
|
||||
- **PATRONI\_ETCD\_SRV\_SUFFIX**: Configures a suffix to the SRV name that is queried during discovery. Use this flag to differentiate between multiple etcd clusters under the same domain. Works only with conjunction with **PATRONI\_ETCD\_SRV**. For example, if ``PATRONI_ETCD_SRV_SUFFIX=foo`` and ``PATRONI_ETCD_SRV=example.org`` are set, the following DNS SRV query is made:``_etcd-client-ssl-foo._tcp.example.com`` (and so on for every possible ETCD SRV service name).
|
||||
- **PATRONI\_ETCD\_USERNAME**: username for etcd authentication.
|
||||
- **PATRONI\_ETCD\_PASSWORD**: password for etcd authentication.
|
||||
- **PATRONI\_ETCD\_CACERT**: The ca certificate. If present it will enable validation.
|
||||
@@ -105,6 +106,7 @@ Kubernetes
|
||||
- **PATRONI\_KUBERNETES\_USE\_ENDPOINTS**: (optional) if set to true, Patroni will use Endpoints instead of ConfigMaps to run leader elections and keep cluster state.
|
||||
- **PATRONI\_KUBERNETES\_POD\_IP**: (optional) IP address of the pod Patroni is running in. This value is required when `PATRONI_KUBERNETES_USE_ENDPOINTS` is enabled and is used to populate the leader endpoint subsets when the pod's PostgreSQL is promoted.
|
||||
- **PATRONI\_KUBERNETES\_PORTS**: (optional) if the Service object has the name for the port, the same name must appear in the Endpoint object, otherwise service won't work. For example, if your service is defined as ``{Kind: Service, spec: {ports: [{name: postgresql, port: 5432, targetPort: 5432}]}}``, then you have to set ``PATRONI_KUBERNETES_PORTS='[{"name": "postgresql", "port": 5432}]'`` and Patroni will use it for updating subsets of the leader Endpoint. This parameter is used only if `PATRONI_KUBERNETES_USE_ENDPOINTS` is set.
|
||||
- **PATRONI\_KUBERNETES\_CACERT**: (optional) Specifies the file with the CA_BUNDLE file with certificates of trusted CAs to use while verifying Kubernetes API SSL certs. If not provided, patroni will use the value provided by the ServiceAccount secret.
|
||||
|
||||
Raft
|
||||
----
|
||||
|
||||
+3
-1
@@ -156,7 +156,8 @@ Most of the parameters are optional, but you have to specify one of the **host**
|
||||
- **use\_proxies**: If this parameter is set to true, Patroni will consider **hosts** as a list of proxies and will not perform a topology discovery of etcd cluster.
|
||||
- **url**: url for the etcd.
|
||||
- **proxy**: proxy url for the etcd. If you are connecting to the etcd using proxy, use this parameter instead of **url**.
|
||||
- **srv**: Domain to search the SRV record(s) for cluster autodiscovery.
|
||||
- **srv**: Domain to search the SRV record(s) for cluster autodiscovery. Patroni will try to query these SRV service names for specified domain (in that order until first success): ``_etcd-client-ssl``, ``_etcd-client``, ``_etcd-ssl``, ``_etcd``, ``_etcd-server-ssl``, ``_etcd-server``. If SRV records for ``_etcd-server-ssl`` or ``_etcd-server`` are retrieved then ETCD peer protocol is used do query ETCD for available members. Otherwise hosts from SRV records will be used.
|
||||
- **srv\_suffix**: Configures a suffix to the SRV name that is queried during discovery. Use this flag to differentiate between multiple etcd clusters under the same domain. Works only with conjunction with **srv**. For example, if ``srv_suffix: foo`` and ``srv: example.org`` are set, the following DNS SRV query is made:``_etcd-client-ssl-foo._tcp.example.com`` (and so on for every possible ETCD SRV service name).
|
||||
- **protocol**: (optional) http or https, if not specified http is used. If the **url** or **proxy** is specified - will take protocol from them.
|
||||
- **username**: (optional) username for etcd authentication.
|
||||
- **password**: (optional) password for etcd authentication.
|
||||
@@ -204,6 +205,7 @@ Kubernetes
|
||||
- **use\_endpoints**: (optional) if set to true, Patroni will use Endpoints instead of ConfigMaps to run leader elections and keep cluster state.
|
||||
- **pod\_ip**: (optional) IP address of the pod Patroni is running in. This value is required when `use_endpoints` is enabled and is used to populate the leader endpoint subsets when the pod's PostgreSQL is promoted.
|
||||
- **ports**: (optional) if the Service object has the name for the port, the same name must appear in the Endpoint object, otherwise service won't work. For example, if your service is defined as ``{Kind: Service, spec: {ports: [{name: postgresql, port: 5432, targetPort: 5432}]}}``, then you have to set ``kubernetes.ports: [{"name": "postgresql", "port": 5432}]`` and Patroni will use it for updating subsets of the leader Endpoint. This parameter is used only if `kubernetes.use_endpoints` is set.
|
||||
- **cacert**: (optional) Specifies the file with the CA_BUNDLE file with certificates of trusted CAs to use while verifying Kubernetes API SSL certs. If not provided, patroni will use the value provided by the ServiceAccount secret.
|
||||
|
||||
|
||||
.. _raft_settings:
|
||||
|
||||
+38
-2
@@ -3,6 +3,42 @@
|
||||
Release notes
|
||||
=============
|
||||
|
||||
Version 2.1.1
|
||||
-------------
|
||||
|
||||
**New features**
|
||||
|
||||
- Support for ETCD SRV name suffix (David Pavlicek)
|
||||
|
||||
Etcd allows to differentiate between multiple Etcd clusters under the same domain and from now on Patroni also supports it.
|
||||
|
||||
- Enrich history with the new leader (huiyalin525)
|
||||
|
||||
It adds the new column to the ``patronictl history`` output.
|
||||
|
||||
- Make the CA bundle configurable for in-cluster Kubernetes config (Aron Parsons)
|
||||
|
||||
By default Patroni is using ``/var/run/secrets/kubernetes.io/serviceaccount/ca.crt`` and this new feature allows specifying the custom ``kubernetes.cacert``.
|
||||
|
||||
- Support dynamically registering/deregistering as a Consul service and changing tags (Tommy Li)
|
||||
|
||||
Previously it required Patroni restart.
|
||||
|
||||
**Bugfixes**
|
||||
|
||||
- Avoid unnecessary reload of REST API (Alexander Kukushkin)
|
||||
|
||||
The previous release added a feature of reloading REST API certificates if changed on disk. Unfortunately, the reload was happening unconditionally right after the start.
|
||||
|
||||
- Don't resolve cluster members when ``etcd.use_proxies`` is set (Alexander)
|
||||
|
||||
When starting up Patroni checks the healthiness of Etcd cluster by querying the list of members. In addition to that, it also tried to resolve their hostnames, which is not necessary when working with Etcd via proxy and was causing unnecessary warnings.
|
||||
|
||||
- Skip rows with NULL values in the ``pg_stat_replication`` (Alexander)
|
||||
|
||||
It seems that the ``pg_stat_replication`` view could contain NULL values in the ``replay_lsn``, ``flush_lsn``, or ``write_lsn`` fields even when ``state = 'streaming'``.
|
||||
|
||||
|
||||
Version 2.1.0
|
||||
-------------
|
||||
|
||||
@@ -39,13 +75,13 @@ This version adds compatibility with PostgreSQL v14, makes logical replication s
|
||||
When everything goes normal, only one line will be written for every run of HA loop.
|
||||
|
||||
|
||||
**Breaking chances**
|
||||
**Breaking changes**
|
||||
|
||||
- The old ``permanent logical replication slots`` feature will no longer work with PostgreSQL v10 and older (Alexander)
|
||||
|
||||
The strategy of creating the logical slots after performing a promotion can't guaranty that no logical events are lost and therefore disabled.
|
||||
|
||||
- The ``/leader`` always endpoint returns 200 if the node holds the lock (Alexander)
|
||||
- The ``/leader`` endpoint always returns 200 if the node holds the lock (Alexander)
|
||||
|
||||
Promoting the standby cluster requires updating load-balancer health checks, which is not very convenient and easy to forget. To solve it, we change the behavior of the ``/leader`` health check endpoint. It will return 200 without taking into account whether the cluster is normal or the ``standby_cluster``.
|
||||
|
||||
|
||||
@@ -20,6 +20,10 @@ metadata:
|
||||
spec:
|
||||
replicas: 3
|
||||
serviceName: *cluster_name
|
||||
selector:
|
||||
matchLabels:
|
||||
application: patroni
|
||||
cluster-name: *cluster_name
|
||||
template:
|
||||
metadata:
|
||||
labels:
|
||||
|
||||
+1
-1
@@ -74,7 +74,7 @@ class Patroni(AbstractPatroniDaemon):
|
||||
if local:
|
||||
self.tags = self.get_tags()
|
||||
self.request.reload_config(self.config)
|
||||
if local or self.api.reload_local_certificate():
|
||||
if local or sighup and self.api.reload_local_certificate():
|
||||
self.api.reload_config(self.config['restapi'])
|
||||
self.watchdog.reload_config(self.config)
|
||||
self.postgresql.reload_config(self.config['postgresql'], sighup)
|
||||
|
||||
+2
-2
@@ -646,10 +646,10 @@ class RestApiServer(ThreadingMixIn, HTTPServer, Thread):
|
||||
self.patroni = patroni
|
||||
self.__listen = None
|
||||
self.__ssl_options = None
|
||||
self.reload_config(config)
|
||||
self.daemon = True
|
||||
self.__ssl_serial_number = None
|
||||
self._received_new_cert = False
|
||||
self.reload_config(config)
|
||||
self.daemon = True
|
||||
|
||||
def query(self, sql, *params):
|
||||
cursor = None
|
||||
|
||||
+1
-1
@@ -347,7 +347,7 @@ class Config(object):
|
||||
if param.startswith(PATRONI_ENV_PREFIX):
|
||||
# PATRONI_(ETCD|CONSUL|ZOOKEEPER|EXHIBITOR|...)_(HOSTS?|PORT|..)
|
||||
name, suffix = (param[8:].split('_', 1) + [''])[:2]
|
||||
if suffix in ('HOST', 'HOSTS', 'PORT', 'USE_PROXIES', 'PROTOCOL', 'SRV', 'URL', 'PROXY',
|
||||
if suffix in ('HOST', 'HOSTS', 'PORT', 'USE_PROXIES', 'PROTOCOL', 'SRV', 'SRV_SUFFIX', 'URL', 'PROXY',
|
||||
'CACERT', 'CERT', 'KEY', 'VERIFY', 'TOKEN', 'CHECKS', 'DC', 'CONSISTENCY',
|
||||
'REGISTER_SERVICE', 'SERVICE_CHECK_INTERVAL', 'NAMESPACE', 'CONTEXT',
|
||||
'USE_ENDPOINTS', 'SCOPE_LABEL', 'ROLE_LABEL', 'POD_IP', 'PORTS', 'LABELS',
|
||||
|
||||
+6
-3
@@ -1299,10 +1299,13 @@ def version(obj, cluster_name, member_names):
|
||||
def history(obj, cluster_name, fmt):
|
||||
cluster = get_dcs(obj, cluster_name).get_cluster()
|
||||
history = cluster.history and cluster.history.lines or []
|
||||
table_header_row = ['TL', 'LSN', 'Reason', 'Timestamp', 'New Leader']
|
||||
for line in history:
|
||||
if len(line) < 4:
|
||||
line.append('')
|
||||
print_output(['TL', 'LSN', 'Reason', 'Timestamp'], history, {'TL': 'r', 'LSN': 'r'}, fmt)
|
||||
if len(line) < len(table_header_row):
|
||||
add_coloumn_num = len(table_header_row) - len(line)
|
||||
for _ in range(add_coloumn_num):
|
||||
line.append('')
|
||||
print_output(table_header_row, history, {'TL': 'r', 'LSN': 'r'}, fmt)
|
||||
|
||||
|
||||
def format_pg_version(version):
|
||||
|
||||
+38
-11
@@ -227,12 +227,11 @@ class Consul(AbstractDCS):
|
||||
self._last_session_refresh = 0
|
||||
self.__session_checks = config.get('checks', [])
|
||||
self._register_service = config.get('register_service', False)
|
||||
self._previous_loop_register_service = self._register_service
|
||||
self._service_tags = sorted(config.get('service_tags', []))
|
||||
self._previous_loop_service_tags = self._service_tags
|
||||
if self._register_service:
|
||||
self._service_tags = config.get('service_tags', [])
|
||||
self._service_name = service_name_from_scope_name(self._scope)
|
||||
if self._scope != self._service_name:
|
||||
logger.warning('Using %s as consul service name instead of scope name %s', self._service_name,
|
||||
self._scope)
|
||||
self._set_service_name()
|
||||
self._service_check_interval = config.get('service_check_interval', '5s')
|
||||
if not self._ctl:
|
||||
self.create_session()
|
||||
@@ -250,7 +249,18 @@ class Consul(AbstractDCS):
|
||||
|
||||
def reload_config(self, config):
|
||||
super(Consul, self).reload_config(config)
|
||||
self._client.reload_config(config.get('consul', {}))
|
||||
|
||||
consul_config = config.get('consul', {})
|
||||
self._client.reload_config(consul_config)
|
||||
self._previous_loop_service_tags = self._service_tags
|
||||
self._service_tags = sorted(consul_config.get('service_tags', []))
|
||||
|
||||
should_register_service = consul_config.get('register_service', False)
|
||||
if should_register_service and not self._register_service:
|
||||
self._set_service_name()
|
||||
|
||||
self._previous_loop_register_service = self._register_service
|
||||
self._register_service = should_register_service
|
||||
|
||||
def set_ttl(self, ttl):
|
||||
if self._client.http.set_ttl(ttl/2.0): # Consul multiplies the TTL by 2x
|
||||
@@ -402,14 +412,18 @@ class Consul(AbstractDCS):
|
||||
self._client.kv.delete(self.member_path)
|
||||
create_member = True
|
||||
|
||||
if self._register_service or self._previous_loop_register_service:
|
||||
try:
|
||||
self.update_service(not create_member and member and member.data or {}, data)
|
||||
except Exception:
|
||||
logger.exception('update_service')
|
||||
|
||||
if not create_member and member and deep_compare(data, member.data):
|
||||
return True
|
||||
|
||||
try:
|
||||
args = {} if permanent else {'acquire': self._session}
|
||||
self._client.kv.put(self.member_path, json.dumps(data, separators=(',', ':')), **args)
|
||||
if self._register_service:
|
||||
self.update_service(not create_member and member and member.data or {}, data)
|
||||
return True
|
||||
except InvalidSession:
|
||||
self._session = None
|
||||
@@ -418,6 +432,11 @@ class Consul(AbstractDCS):
|
||||
logger.exception('touch_member')
|
||||
return False
|
||||
|
||||
def _set_service_name(self):
|
||||
self._service_name = service_name_from_scope_name(self._scope)
|
||||
if self._scope != self._service_name:
|
||||
logger.warning('Using %s as consul service name instead of scope name %s', self._service_name, self._scope)
|
||||
|
||||
@catch_consul_errors
|
||||
def register_service(self, service_name, **kwargs):
|
||||
logger.info('Register service %s, params %s', service_name, kwargs)
|
||||
@@ -441,17 +460,22 @@ class Consul(AbstractDCS):
|
||||
deregister='{0}s'.format(self._client.http.ttl * 10))
|
||||
tags = self._service_tags[:]
|
||||
tags.append(role)
|
||||
self._previous_loop_service_tags = self._service_tags
|
||||
|
||||
params = {
|
||||
'service_id': '{0}/{1}'.format(self._scope, self._name),
|
||||
'address': conn_parts.hostname,
|
||||
'port': conn_parts.port,
|
||||
'check': check,
|
||||
'tags': tags
|
||||
'tags': tags,
|
||||
'enable_tag_override': True,
|
||||
}
|
||||
|
||||
if state == 'stopped':
|
||||
if state == 'stopped' or (not self._register_service and self._previous_loop_register_service):
|
||||
self._previous_loop_register_service = self._register_service
|
||||
return self.deregister_service(params['service_id'])
|
||||
|
||||
self._previous_loop_register_service = self._register_service
|
||||
if role in ['master', 'replica', 'standby-leader']:
|
||||
if state != 'running':
|
||||
return
|
||||
@@ -470,7 +494,10 @@ class Consul(AbstractDCS):
|
||||
if old_data.get(key) != new_data[key]:
|
||||
update = True
|
||||
|
||||
if force or update:
|
||||
if (
|
||||
force or update or self._register_service != self._previous_loop_register_service
|
||||
or self._service_tags != self._previous_loop_service_tags
|
||||
):
|
||||
return self._update_service(new_data)
|
||||
|
||||
@catch_consul_errors
|
||||
|
||||
+5
-3
@@ -188,7 +188,8 @@ class AbstractEtcdClientWithFailover(etcd.Client):
|
||||
logger.debug("Retrieved list of machines: %s", machines)
|
||||
if machines:
|
||||
random.shuffle(machines)
|
||||
self._update_dns_cache(self._dns_resolver.resolve_async, machines)
|
||||
if not self._use_proxies:
|
||||
self._update_dns_cache(self._dns_resolver.resolve_async, machines)
|
||||
return machines
|
||||
except Exception as e:
|
||||
self.http.clear()
|
||||
@@ -282,13 +283,14 @@ class AbstractEtcdClientWithFailover(etcd.Client):
|
||||
except DNSException:
|
||||
return []
|
||||
|
||||
def _get_machines_cache_from_srv(self, srv):
|
||||
def _get_machines_cache_from_srv(self, srv, srv_suffix=None):
|
||||
"""Fetch list of etcd-cluster member by resolving _etcd-server._tcp. SRV record.
|
||||
This record should contain list of host and peer ports which could be used to run
|
||||
'GET http://{host}:{port}/members' request (peer protocol)"""
|
||||
|
||||
ret = []
|
||||
for r in ['-client-ssl', '-client', '-ssl', '', '-server-ssl', '-server']:
|
||||
r = '{0}-{1}'.format(r, srv_suffix) if srv_suffix else r
|
||||
protocol = 'https' if '-ssl' in r else 'http'
|
||||
endpoint = '/members' if '-server' in r else ''
|
||||
for host, port in self.get_srv_record('_etcd{0}._tcp.{1}'.format(r, srv)):
|
||||
@@ -325,7 +327,7 @@ class AbstractEtcdClientWithFailover(etcd.Client):
|
||||
|
||||
machines_cache = []
|
||||
if 'srv' in self._config:
|
||||
machines_cache = self._get_machines_cache_from_srv(self._config['srv'])
|
||||
machines_cache = self._get_machines_cache_from_srv(self._config['srv'], self._config.get('srv_suffix'))
|
||||
|
||||
if not machines_cache and 'hosts' in self._config:
|
||||
machines_cache = list(self._config['hosts'])
|
||||
|
||||
@@ -617,7 +617,7 @@ class Etcd3(AbstractEtcd):
|
||||
return self.retry(self._do_refresh_lease)
|
||||
except (Etcd3ClientError, RetryFailedError):
|
||||
logger.exception('refresh_lease')
|
||||
raise Etcd3Error('Failed ro keepalive/grant lease')
|
||||
raise Etcd3Error('Failed to keepalive/grant lease')
|
||||
|
||||
def create_lease(self):
|
||||
while not self._lease:
|
||||
|
||||
@@ -55,16 +55,19 @@ class K8sConfig(object):
|
||||
if token:
|
||||
self._headers['authorization'] = 'Bearer ' + token
|
||||
|
||||
def load_incluster_config(self):
|
||||
def load_incluster_config(self, ca_certs=SERVICE_CERT_FILENAME):
|
||||
if SERVICE_HOST_ENV_NAME not in os.environ or SERVICE_PORT_ENV_NAME not in os.environ:
|
||||
raise self.ConfigException('Service host/port is not set.')
|
||||
if not os.environ[SERVICE_HOST_ENV_NAME] or not os.environ[SERVICE_PORT_ENV_NAME]:
|
||||
raise self.ConfigException('Service host/port is set but empty.')
|
||||
if not os.path.isfile(SERVICE_CERT_FILENAME):
|
||||
|
||||
if not os.path.isfile(ca_certs):
|
||||
raise self.ConfigException('Service certificate file does not exists.')
|
||||
with open(SERVICE_CERT_FILENAME) as f:
|
||||
with open(ca_certs) as f:
|
||||
if not f.read():
|
||||
raise self.ConfigException('Cert file exists but empty.')
|
||||
self.pool_config['ca_certs'] = ca_certs
|
||||
|
||||
if not os.path.isfile(SERVICE_TOKEN_FILENAME):
|
||||
raise self.ConfigException('Service token file does not exists.')
|
||||
with open(SERVICE_TOKEN_FILENAME) as f:
|
||||
@@ -72,7 +75,6 @@ class K8sConfig(object):
|
||||
if not token:
|
||||
raise self.ConfigException('Token file exists but empty.')
|
||||
self._make_headers(token=token)
|
||||
self.pool_config['ca_certs'] = SERVICE_CERT_FILENAME
|
||||
self._server = uri('https', (os.environ[SERVICE_HOST_ENV_NAME], os.environ[SERVICE_PORT_ENV_NAME]))
|
||||
|
||||
@staticmethod
|
||||
@@ -613,13 +615,14 @@ class Kubernetes(AbstractDCS):
|
||||
self._label_selector = ','.join('{0}={1}'.format(k, v) for k, v in self._labels.items())
|
||||
self._namespace = config.get('namespace') or 'default'
|
||||
self._role_label = config.get('role_label', 'role')
|
||||
self._ca_certs = os.environ.get('PATRONI_KUBERNETES_CACERT', config.get('cacert')) or SERVICE_CERT_FILENAME
|
||||
config['namespace'] = ''
|
||||
super(Kubernetes, self).__init__(config)
|
||||
self._retry = Retry(deadline=config['retry_timeout'], max_delay=1, max_tries=-1,
|
||||
retry_exceptions=KubernetesRetriableException)
|
||||
self._ttl = None
|
||||
try:
|
||||
k8s_config.load_incluster_config()
|
||||
k8s_config.load_incluster_config(ca_certs=self._ca_certs)
|
||||
except k8s_config.ConfigException:
|
||||
k8s_config.load_kube_config(context=config.get('context', 'local'))
|
||||
|
||||
|
||||
+4
-2
@@ -551,7 +551,7 @@ class Ha(object):
|
||||
if master_timeline == 1:
|
||||
if cluster_history:
|
||||
self.dcs.set_history_value('[]')
|
||||
elif not cluster_history or cluster_history[-1][0] != master_timeline - 1 or len(cluster_history[-1]) != 4:
|
||||
elif not cluster_history or cluster_history[-1][0] != master_timeline - 1 or len(cluster_history[-1]) != 5:
|
||||
cluster_history = {line[0]: line for line in cluster_history or []}
|
||||
history = self.state_handler.get_history(master_timeline)
|
||||
if history and self.cluster.config:
|
||||
@@ -559,9 +559,11 @@ class Ha(object):
|
||||
for line in history:
|
||||
# enrich current history with promotion timestamps stored in DCS
|
||||
if len(line) == 3 and line[0] in cluster_history \
|
||||
and len(cluster_history[line[0]]) == 4 \
|
||||
and len(cluster_history[line[0]]) >= 4 \
|
||||
and cluster_history[line[0]][1] == line[1]:
|
||||
line.append(cluster_history[line[0]][3])
|
||||
if len(cluster_history[line[0]]) == 5:
|
||||
line.append(cluster_history[line[0]][4])
|
||||
self.dcs.set_history_value(json.dumps(history, separators=(',', ':')))
|
||||
|
||||
def enforce_follow_remote_master(self, message):
|
||||
|
||||
@@ -843,6 +843,7 @@ class Postgresql(object):
|
||||
if history[-1][0] == timeline - 1:
|
||||
history_mtime = datetime.fromtimestamp(history_mtime).replace(tzinfo=tz.tzlocal())
|
||||
history[-1].append(history_mtime.isoformat())
|
||||
history[-1].append(self.name)
|
||||
return history
|
||||
except Exception:
|
||||
logger.exception('Failed to read and parse %s', (history_path,))
|
||||
@@ -1077,7 +1078,7 @@ class Postgresql(object):
|
||||
for app_name, sync_state, replica_lsn in self.query(
|
||||
"SELECT pg_catalog.lower(application_name), sync_state, pg_{2}_{1}_diff({0}_{1}, '0/0')::bigint"
|
||||
" FROM pg_catalog.pg_stat_replication"
|
||||
" WHERE state = 'streaming'"
|
||||
" WHERE state = 'streaming' AND {0}_{1} IS NOT NULL"
|
||||
" ORDER BY sync_state DESC, {0}_{1} DESC".format(sort_col, self.lsn_name, self.wal_name)):
|
||||
member = members.get(app_name)
|
||||
if member and not member.tags.get('nosync', False):
|
||||
|
||||
@@ -1001,8 +1001,9 @@ class ConfigHandler(object):
|
||||
if self._postgresql.major_version >= 90500:
|
||||
time.sleep(1)
|
||||
try:
|
||||
pending_restart = self._postgresql.query('SELECT COUNT(*) FROM pg_catalog.pg_settings'
|
||||
' WHERE pending_restart').fetchone()[0] > 0
|
||||
pending_restart = self._postgresql.query(
|
||||
'SELECT COUNT(*) FROM pg_catalog.pg_settings WHERE pg_catalog.lower(name) != ALL(%s)'
|
||||
' AND pending_restart', [n.lower() for n in self._RECOVERY_PARAMETERS]).fetchone()[0] > 0
|
||||
self._postgresql.set_pending_restart(pending_restart)
|
||||
except Exception as e:
|
||||
logger.warning('Exception %r when running query', e)
|
||||
|
||||
@@ -170,7 +170,6 @@ parameters = CaseInsensitiveDict({
|
||||
'DateStyle': String(90300, None),
|
||||
'db_user_namespace': Bool(90300, None),
|
||||
'deadlock_timeout': Integer(90300, None, 1, 2147483647, 'ms'),
|
||||
'debug_invalidate_system_caches_always': Integer(140000, None, '0', '0', None),
|
||||
'debug_pretty_print': Bool(90300, None),
|
||||
'debug_print_parse': Bool(90300, None),
|
||||
'debug_print_plan': Bool(90300, None),
|
||||
@@ -208,7 +207,6 @@ parameters = CaseInsensitiveDict({
|
||||
'enable_partition_pruning': Bool(110000, None),
|
||||
'enable_partitionwise_aggregate': Bool(110000, None),
|
||||
'enable_partitionwise_join': Bool(110000, None),
|
||||
'enable_resultcache': Bool(140000, None),
|
||||
'enable_seqscan': Bool(90300, None),
|
||||
'enable_sort': Bool(90300, None),
|
||||
'enable_tidscan': Bool(90300, None),
|
||||
|
||||
@@ -302,10 +302,11 @@ validate_host_port_listen.expected_type = string_types
|
||||
validate_host_port_listen_multiple_hosts.expected_type = string_types
|
||||
validate_data_dir.expected_type = string_types
|
||||
validate_etcd = {
|
||||
Or("host", "hosts", "srv", "url", "proxy"): Case({
|
||||
Or("host", "hosts", "srv", "srv_suffix", "url", "proxy"): Case({
|
||||
"host": validate_host_port,
|
||||
"hosts": Or(comma_separated_host_port, [validate_host_port]),
|
||||
"srv": str,
|
||||
"srv_suffix": str,
|
||||
"url": str,
|
||||
"proxy": str})
|
||||
}
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
__version__ = '2.1.0'
|
||||
__version__ = '2.1.1'
|
||||
|
||||
+50
-3
@@ -130,8 +130,9 @@ class TestConsul(unittest.TestCase):
|
||||
@patch.object(consul.Consul.KV, 'put', Mock(side_effect=[True, ConsulException, InvalidSession]))
|
||||
def test_touch_member(self):
|
||||
self.c.refresh_session = Mock(return_value=False)
|
||||
self.c.touch_member({'conn_url': 'postgres://replicator:[email protected]:5433/postgres',
|
||||
'api_url': 'http://127.0.0.1:8009/patroni'})
|
||||
with patch.object(Consul, 'update_service', Mock(side_effect=Exception)):
|
||||
self.c.touch_member({'conn_url': 'postgres://replicator:[email protected]:5433/postgres',
|
||||
'api_url': 'http://127.0.0.1:8009/patroni'})
|
||||
self.c._register_service = True
|
||||
self.c.refresh_session = Mock(return_value=True)
|
||||
for _ in range(0, 4):
|
||||
@@ -215,4 +216,50 @@ class TestConsul(unittest.TestCase):
|
||||
self.assertIsNone(self.c.update_service({}, d))
|
||||
|
||||
def test_reload_config(self):
|
||||
self.c.reload_config({'consul': {'token': 'foo'}, 'loop_wait': 10, 'ttl': 30, 'retry_timeout': 10})
|
||||
self.assertEqual([], self.c._service_tags)
|
||||
self.c.reload_config({'consul': {'token': 'foo', 'register_service': True, 'service_tags': ['foo']},
|
||||
'loop_wait': 10, 'ttl': 30, 'retry_timeout': 10})
|
||||
self.assertEqual(["foo"], self.c._service_tags)
|
||||
|
||||
self.c.refresh_session = Mock(return_value=False)
|
||||
|
||||
d = {'role': 'replica', 'api_url': 'http://a/t', 'conn_url': 'pg://c:1', 'state': 'running'}
|
||||
|
||||
# Changing register_service from True to False calls deregister()
|
||||
self.c.reload_config({'consul': {'register_service': False}, 'loop_wait': 10, 'ttl': 30, 'retry_timeout': 10})
|
||||
with patch('consul.Consul.Agent.Service.deregister') as mock_deregister:
|
||||
self.c.touch_member(d)
|
||||
mock_deregister.assert_called_once()
|
||||
|
||||
self.assertEqual([], self.c._service_tags)
|
||||
|
||||
# register_service staying False between reloads does not call deregister()
|
||||
self.c.reload_config({'consul': {'register_service': False}, 'loop_wait': 10, 'ttl': 30, 'retry_timeout': 10})
|
||||
with patch('consul.Consul.Agent.Service.deregister') as mock_deregister:
|
||||
self.c.touch_member(d)
|
||||
self.assertFalse(mock_deregister.called)
|
||||
|
||||
# Changing register_service from False to True calls register()
|
||||
self.c.reload_config({'consul': {'register_service': True}, 'loop_wait': 10, 'ttl': 30, 'retry_timeout': 10})
|
||||
with patch('consul.Consul.Agent.Service.register') as mock_register:
|
||||
self.c.touch_member(d)
|
||||
mock_register.assert_called_once()
|
||||
|
||||
# register_service staying True between reloads does not call register()
|
||||
self.c.reload_config({'consul': {'register_service': True}, 'loop_wait': 10, 'ttl': 30, 'retry_timeout': 10})
|
||||
with patch('consul.Consul.Agent.Service.register') as mock_register:
|
||||
self.c.touch_member(d)
|
||||
self.assertFalse(mock_deregister.called)
|
||||
|
||||
# register_service staying True between reloads does calls register() if other service data has changed
|
||||
self.c.reload_config({'consul': {'register_service': True}, 'loop_wait': 10, 'ttl': 30, 'retry_timeout': 10})
|
||||
with patch('consul.Consul.Agent.Service.register') as mock_register:
|
||||
self.c.touch_member(d)
|
||||
mock_register.assert_called_once()
|
||||
|
||||
# register_service staying True between reloads does calls register() if service_tags have changed
|
||||
self.c.reload_config({'consul': {'register_service': True, 'service_tags': ['foo']}, 'loop_wait': 10,
|
||||
'ttl': 30, 'retry_timeout': 10})
|
||||
with patch('consul.Consul.Agent.Service.register') as mock_register:
|
||||
self.c.touch_member(d)
|
||||
mock_register.assert_called_once()
|
||||
|
||||
+3
-1
@@ -87,7 +87,8 @@ def dns_query(name, _):
|
||||
raise DNSException()
|
||||
srv = Mock()
|
||||
srv.port = 2380
|
||||
srv.target.to_text.return_value = 'localhost' if name == '_etcd-server._tcp.foobar' else '127.0.0.1'
|
||||
srv.target.to_text.return_value = \
|
||||
'localhost' if name in ['_etcd-server._tcp.foobar', '_etcd-server-baz._tcp.foobar'] else '127.0.0.1'
|
||||
return [srv]
|
||||
|
||||
|
||||
@@ -183,6 +184,7 @@ class TestClient(unittest.TestCase):
|
||||
|
||||
def test__get_machines_cache_from_srv(self):
|
||||
self.client._get_machines_cache_from_srv('foobar')
|
||||
self.client._get_machines_cache_from_srv('foobar', 'baz')
|
||||
self.client.get_srv_record = Mock(return_value=[('localhost', 2380)])
|
||||
self.client._get_machines_cache_from_srv('blabla')
|
||||
|
||||
|
||||
+2
-2
@@ -35,8 +35,8 @@ def false(*args, **kwargs):
|
||||
|
||||
def get_cluster(initialize, leader, members, failover, sync, cluster_config=None):
|
||||
t = datetime.datetime.now().isoformat()
|
||||
history = TimelineHistory(1, '[[1,67197376,"no recovery target specified","' + t + '"]]',
|
||||
[(1, 67197376, 'no recovery target specified', t)])
|
||||
history = TimelineHistory(1, '[[1,67197376,"no recovery target specified","' + t + '","foo"]]',
|
||||
[(1, 67197376, 'no recovery target specified', t, 'foo')])
|
||||
cluster_config = cluster_config or ClusterConfig(1, {'check_timeline': True}, 1)
|
||||
return Cluster(initialize, cluster_config, leader, 10, members, failover, sync, history, None)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user