Compare commits

...
15 Commits
Author SHA1 Message Date
Alexander KukushkinandGitHub c6e70a9910 Release 1.5.5 (#979)
* Bump version
* Update release notes
2019-02-15 16:14:39 +01:00
Alexander KukushkinandGitHub 0ec1760397 Don't write primary_conninfo into recovery.conf for wal only standby cluster (#971)
It is useless and makes postgres to generate a lot of errors
2019-02-15 13:35:34 +01:00
Alexander KukushkinandGitHub 317721991b Fix handling of PATRONI_*_PASSWORD environment variables (#970)
Bug was introduced in https://github.com/zalando/patroni/pull/947 and https://github.com/zalando/patroni/pull/500
2019-02-15 13:35:22 +01:00
Alexander KukushkinandGitHub 0c516de147 Create headless service associated with $SCOPE-config endpoint (#958)
if there is no service defined k8s assumes that endpoint is orphaned and removes it.
Patroni tries to create the service only in case if use_endpoints is enabled if the following cases:
1. Upon start
2. When it tries to (re-)create the config endpoint

If for some reason creation of the service has failed, Patroni will retry it on every cycle of HA loop. Usually it fails due to lack of permissions and if you don't want to give such permissions to the service account used by Patroni, you can create the service explicitly in the deployment manifest.
2019-02-15 13:35:04 +01:00
Michael BanckandAlexander Kukushkin 073074f83e Run coverage as python -m coverage (#968)
Depending on the platform the coverage binary might not always be available under the standard name.
2019-02-13 16:02:12 +01:00
Michael BanckandAlexander Kukushkin 345e6d3131 Copy away output directories of failed acceptance tests. (#967)
And dump logs on travis from only failed features
2019-02-13 16:00:15 +01:00
Michael BanckandAlexander Kukushkin d01a9bdcd5 Change base port for acceptance tests from 5440 to 5360 (#966) 2019-02-13 15:59:13 +01:00
Alexander KukushkinandGitHub 10bbf0c3c5 Always use replication=1, otherwise it is considered a logical walsender (#952)
Fixes: https://github.com/zalando/patroni/issues/894
Fixes: https://github.com/zalando/patroni/issues/951
2019-01-30 12:38:51 +01:00
Alexander KukushkinandGitHub 4304560ce2 Adjust read timeout for leader watch blocking query (#950)
According to the Consul documentation the actual response timeout is increased by a small random amount of additional wait time added to the supplied maximum wait time to spread out the wake up time of any concurrent requests. It adds up to wait / 16 additional time to the maximum duration.
In our case we will add wait/15 or 1 second depending on what is bigger.

Fixes: https://github.com/zalando/patroni/issues/945
2019-01-30 12:38:24 +01:00
Alexander KukushkinandGitHub 254ee2acfc Show information about timelines in patronictl list (#949)
This information will help to detect stale replicas.

In addition to that Host will include ':{port}' if the port value isn't default or more than one member running on the same host.

Fixes: https://github.com/zalando/patroni/issues/942
2019-01-30 12:37:51 +01:00
Alexander KukushkinandGitHub 739329b590 Make it possible to automatically reinit the former master (#948)
If the pg_rewind is disabled or can't be used, the former master could fail to start as a new replica due to diverged timelines. In this case, the only way to fix it is wiping the data directory and reinitializing.

So far Patroni was able to remove the data directory only after failed attempt to run pg_rewind. This commit fixes it.
If the `postgresql.remove_data_directory_on_diverged_timelines` is set, Patroni will wipe the data directory and reinitialize the former master automatically.

Fixes: https://github.com/zalando/patroni/issues/941
2019-01-30 12:37:21 +01:00
Étienne MandAlexander Kukushkin bd2c54581a Add ETCD_(PROTOCOL|USERNAME|PASSWORD) env variables (#947)
Fix #944
2019-01-30 12:36:50 +01:00
Maxim IvanovandAlexander Kukushkin f0b12b7e2e Document create_replicas_methods in standby_cluster section (#939)
Fixes https://github.com/zalando/patroni/issues/935
2019-01-30 12:36:24 +01:00
Étienne MandAlexander Kukushkin 93d157dea3 Document how to start Patroni with an existing data directory (#918) 2019-01-30 12:35:57 +01:00
Alexander KukushkinandGitHub 2c128520cf Python34 compatibility (#933)
and some other minor fixes.

Closes https://github.com/zalando/patroni/issues/932
2019-01-16 14:40:05 +01:00
24 changed files with 314 additions and 123 deletions
+1 -1
View File
@@ -137,7 +137,7 @@ script:
echo Running acceptance tests using python${pv}
if ! PATH=.:/usr/lib/postgresql/9.6/bin:$PATH $TEST_SUITE; then
# output all log files when tests are failing
grep . features/output/*/*postgres?.*
grep . features/output/*_failed/*postgres?.*
exit 1
fi
fi
+8 -4
View File
@@ -46,13 +46,17 @@ Consul
Etcd
----
- **PATRONI\_ETCD\_HOST**: the host:port for the etcd endpoint.
- **PATRONI\_ETCD\_HOSTS**: list of etcd endpoints in format 'host1:port1','host2:port2',etc...
- **PATRONI\_ETCD\_URL**: url for the etcd, in format: http(s)://(username:password@)host:port
- **PATRONI\_ETCD\_PROXY**: proxy url for the etcd. If you are connecting to the etcd using proxy, use this parameter instead of **PATRONI\_ETCD\_URL**
- **PATRONI\_ETCD\_URL**: url for the etcd, in format: http(s)://(username:password@)host:port
- **PATRONI\_ETCD\_HOSTS**: list of etcd endpoints in format 'host1:port1','host2:port2',etc...
- **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\_USERNAME**: username for etcd authentication.
- **PATRONI\_ETCD\_PASSWORD**: password for etcd authentication.
- **PATRONI\_ETCD\_CACERT**: The ca certificate. If present it will enable validation.
- **PATRONI\_ETCD\_CERT**: File with the client certificate
- **PATRONI\_ETCD\_CERT**: File with the client certificate.
- **PATRONI\_ETCD\_KEY**: File with the client key. Can be empty if the key is part of certificate.
Exhibitor
+2
View File
@@ -66,6 +66,8 @@ Note that external tools to call in the replica creation or custom bootstrap scr
independently of Patroni.
.. _running_configuring:
Running and Configuring
-----------------------
+9 -6
View File
@@ -22,6 +22,8 @@ Log
- **patroni.postmaster: WARNING**
- **urllib3: DEBUG**
.. _bootstrap_settings:
Bootstrap configuration
-----------------------
- **dcs**: This section will be written into `/<namespace>/<scope>/config` of a given configuration store after initializing of new cluster. This is the global configuration for the cluster. If you want to change some parameters for all cluster nodes - just do it in DCS (or via Patroni API) and all nodes will apply this configuration.
@@ -35,7 +37,7 @@ Bootstrap configuration
- **postgresql**:
- **use\_pg\_rewind**: whether or not to use pg_rewind
- **use\_slots**: whether or not to use replication_slots. Must be False for PostgreSQL 9.3. You should comment out max_replication_slots before it becomes ineligible for leader status.
- **recovery\_conf**: additional configuration settings written to recovery.conf when configuring follower.
- **recovery\_conf**: additional configuration settings written to recovery.conf when configuring follower.
- **parameters**: list of configuration settings for Postgres. Many of these are required for replication to work.
- **standby\_cluster**: if this section is defined, we want to bootstrap a standby cluster.
- **host**: an address of remote master
@@ -99,10 +101,10 @@ Most of the parameters are optional, but you have to specify one of the **host**
- **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.
- **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
- **username**: (optional) username for etcd authentication.
- **password**: (optional) password for etcd authentication.
- **cacert**: (optional) The ca certificate. If present it will enable validation.
- **cert**: (optional) file with the client certificate
- **cert**: (optional) file with the client certificate.
- **key**: (optional) file with the client key. Can be empty if the key is part of **cert**.
Exhibitor
@@ -144,7 +146,7 @@ PostgreSQL
- **create\_replica\_methods**: an ordered list of the create methods for turning a Patroni node into a new replica.
"basebackup" is the default method; other methods are assumed to refer to scripts, each of which is configured as its
own config item. See :ref:`custom replica creation methods documentation <custom_replica_creation>` for further explanation.
- **data\_dir**: The location of the Postgres data directory, either existing or to be initialized by Patroni.
- **data\_dir**: The location of the Postgres data directory, either :ref:`existing <existing_data>` or to be initialized by Patroni.
- **config\_dir**: The location of the Postgres configuration directory, defaults to the data directory. Must be writable by Patroni.
- **bin\_dir**: Path to PostgreSQL binaries (pg_ctl, pg_rewind, pg_basebackup, postgres). The default value is an empty string meaning that PATH environment variable will be used to find the executables.
- **listen**: IP address + port that Postgres listens to; must be accessible from other nodes in the cluster, if you're using streaming replication. Multiple comma-separated addresses are permitted, as long as the port component is appended after to the last one with a colon, i.e. ``listen: 127.0.0.1,127.0.0.2:5432``. Patroni will use the first address from this list to establish local connections to the PostgreSQL node.
@@ -159,10 +161,11 @@ PostgreSQL
- **pg\_ctl\_timeout**: How long should pg_ctl wait when doing ``start``, ``stop`` or ``restart``. Default value is 60 seconds.
- **use\_pg\_rewind**: try to use pg\_rewind on the former leader when it joins cluster as a replica.
- **remove\_data\_directory\_on\_rewind\_failure**: If this option is enabled, Patroni will remove postgres data directory and recreate replica. Otherwise it will try to follow the new leader. Default value is **false**.
- **remove\_data\_directory\_on\_diverged\_timelines**: Patroni will remove postgres data directory and recreate replica if it notices that timelines are diverging and the former master can not start streaming from the new master. This option is useful when ``pg_rewind`` can not be used. Default value is **false**.
- **replica\_method**: for each create_replica_methods other than basebackup, you would add a configuration section of the same name. At a minimum, this should include "command" with a full path to the actual script to be executed. Other configuration parameters will be passed along to the script in the form "parameter=value".
REST API
--------
--------
- **connect\_address**: IP address (or hostname) and port, to access the Patroni's REST API. All the members of the cluster must be able to connect to this address, so unless the Patroni setup is intended for a demo inside the localhost, this address must be a non "localhost" or loopback addres (ie: "localhost" or "127.0.0.1"). It can serve as a endpoint for HTTP health checks (read below about the "listen" REST API parameter), and also for user queries (either directly or via the REST API), as well as for the health checks done by the cluster members during leader elections (for example, to determine whether the master is still running, or if there is a node which has a WAL position that is ahead of the one doing the query; etc.) The connect_address is put in the member key in DCS, making it possible to translate the member name into the address to connect to its REST API.
- **listen**: IP address (or hostname) and port that Patroni will listen to for the REST API - to provide also the same health checks and cluster messaging between the participating nodes, as described above. to provide health-check information for HAProxy (or any other load balancer capable of doing a HTTP "OPTION" or "GET" checks).
@@ -180,7 +183,7 @@ REST API
CTL
---
- **Optional**:
- **insecure**: Allow connections to REST API without verifying SSL certs.
- **insecure**: Allow connections to REST API without verifying SSL certs.
- **cacert**: Specifices the file with the CA_BUNDLE file or directory with certificates of trusted CAs to use while verifying REST API SSL certs.
- **certfile**: Specifies the file with the certificate in the PEM format to use while verifying REST API SSL certs. If not provided patronictl will use the value provided for REST API "certfile" parameter.
+35
View File
@@ -0,0 +1,35 @@
.. _existing_data:
Convert a Standalone to a Patroni Cluster
=========================================
This section describes the process for converting a standalone PostgreSQL instance into a Patroni cluster.
To deploy a Patroni cluster without using a pre-existing PostgreSQL instance, see :ref:`Running and Configuring <running_configuring>` instead.
Procedure
---------
A Patroni cluster can be started with a data directory from a single-node PostgreSQL database. This is achieved by following closely these steps:
#. Manually start PostgreSQL daemon
#. Create Patroni superuser and replication users as defined in the :ref:`authentication <postgresql_settings>` section of the Patroni configuration. If this user is created in SQL, the following queries achieve this:
.. code-block:: sql
CREATE USER $PATRONI_SUPERUSER_USERNAME WITH SUPERUSER ENCRYPTED PASSWORD '$PATRONI_SUPERUSER_PASSWORD';
CREATE USER $PATRONI_REPLICATION_USERNAME WITH REPLICATION ENCRYPTED PASSWORD '$PATRONI_REPLICATION_PASSWORD';
#. Start Patroni (e.g. ``patroni /etc/patroni/patroni.yml``). It automatically detects that PostgreSQL daemon is already running but its configuration might be out-of-date.
#. Ask Patroni to restart the node with ``patronictl restart cluster-name node-name``.
FAQ
---
#. During Patroni startup, Patroni complains that it cannot bind to the PostgreSQL port.
You need to verify ``listen_addresses`` and ``port`` in ``postgresql.conf`` and ``postgresql.listen`` in ``patroni.yml``. Don't forget that ``pg_hba.conf`` should allow such access.
#. After asking Patroni to restart the node, PostgreSQL displays the error message ``could not open configuration file "/etc/postgresql/10/main/pg_hba.conf": No such file or directory``
It can mean various things depending on how you manage PostgreSQL configuration. If you specified `postgresql.config_dir`, Patroni generates the ``pg_hba.conf`` based on the settings in the :ref:`bootstrap <bootstrap_settings>` section only when it bootstraps a new cluster. In this scenario the ``PGDATA`` was not empty, therefore no bootstrap happened. This file must exist beforehand.
+38
View File
@@ -3,6 +3,44 @@
Release notes
=============
Version 1.5.5
-------------
This version introduces the possibility of automatic reinit of the former master, improves patronictl list output and fixes a number of bugs.
**New features**
- Add support of `PATRONI_ETCD_PROTOCOL`, `PATRONI_ETCD_USERNAME` and `PATRONI_ETCD_PASSWORD` environment variables (Étienne M)
Before it was possible to configure them only in the config file or as a part of `PATRONI_ETCD_URL`, which is not always convenient.
- Make it possible to automatically reinit the former master (Alexander Kukushkin)
If the pg_rewind is disabled or can't be used, the former master could fail to start as a new replica due to diverged timelines. In this case, the only way to fix it is wiping the data directory and reinitializing. This behavior could be changed by setting `postgresql.remove_data_directory_on_diverged_timelines`. When it is set, Patroni will wipe the data directory and reinitialize the former master automatically.
- Show information about timelines in patronictl list (Alexander)
It helps to detect stale replicas. In addition to that, `Host` will include ':{port}' if the port value isn't default or there is more than one member running on the same host.
- Create a headless service associated with the $SCOPE-config endpoint (Alexander)
The "config" endpoint keeps information about the cluster-wide Patroni and Postgres configuration, history file, and last but the most important, it holds the `initialize` key. When the Kubernetes master node is restarted or upgraded, it removes endpoints without services. The headless service will prevent it from being removed.
**Bug fixes**
- Adjust the read timeout for the leader watch blocking query (Alexander)
According to the Consul documentation, the actual response timeout is increased by a small random amount of additional wait time added to the supplied maximum wait time to spread out the wake up time of any concurrent requests. It adds up to `wait / 16` additional time to the maximum duration. In our case we are adding `wait / 15` or 1 second depending on what is bigger.
- Always use replication=1 when connecting via replication protocol to the postgres (Alexander)
Starting from Postgres 10 the line in the pg_hba.conf with database=replication doesn't accept connections with the parameter replication=database.
- Don't write primary_conninfo into recovery.conf for wal-only standby cluster (Alexander)
Despite not having neither `host` nor `port` defined in the `standby_cluster` config, Patroni was putting the `primary_conninfo` into the `recovery.conf`, which is useless and generating a lot of errors.
Version 1.5.4
-------------
+11 -3
View File
@@ -177,9 +177,15 @@ standby nodes replicating from some remote master. This type of clusters has:
Standby leader holds and updates a leader lock in DCS. If the leader lock
expires, cascade replicas will perform an election to choose another leader
from the standbys. For the sake of flexibility, you can specify different
methods of creating a replica and recovery WAL records when a cluster is in the
"standby mode", and after it was detached to function as a normal cluster.
from the standbys.
For the sake of flexibility, you can specify methods of creating a replica and
recovery WAL records when a cluster is in the "standby mode" by providing
`create_replica_methods` key in `standby_cluster` section. It is distinct from
creating replicas, when cluster is detached and functions as a normal cluster,
which is controlled by `create_replica_methods` in `postgresql` section. Both
"standby" and "normal" `create_replica_methods` reference keys in `postgresql`
section.
To configure such cluster you need to specify the section ``standby_cluster``
in a patroni configuration:
@@ -192,6 +198,8 @@ in a patroni configuration:
host: 1.2.3.4
port: 5432
primary_slot_name: patroni
create_replica_methods:
- basebackup
Note, that these options will be applied only once during cluster bootstrap,
and the only way to change them afterwards is through DCS.
+8 -4
View File
@@ -12,6 +12,7 @@ import shutil
import signal
import six
import subprocess
import sys
import tempfile
import threading
import time
@@ -84,7 +85,7 @@ class AbstractController(object):
class PatroniController(AbstractController):
__PORT = 5440
__PORT = 5360
PATRONI_CONFIG = '{}.yml'
""" starts and stops individual patronis"""
@@ -142,7 +143,8 @@ class PatroniController(AbstractController):
if isinstance(self._context.dcs_ctl, KubernetesController):
self._context.dcs_ctl.create_pod(self._name[8:], self._scope)
os.environ['PATRONI_KUBERNETES_POD_IP'] = '10.0.0.' + self._name[-1]
return subprocess.Popen(['coverage', 'run', '--source=patroni', '-p', 'patroni.py', self._config],
return subprocess.Popen([sys.executable, '-m', 'coverage', 'run',
'--source=patroni', '-p', 'patroni.py', self._config],
stdout=self._log, stderr=subprocess.STDOUT, cwd=self._work_directory)
def stop(self, kill=False, timeout=15, postgres=False):
@@ -809,8 +811,8 @@ def before_all(context):
def after_all(context):
context.dcs_ctl.stop()
subprocess.call(['coverage', 'combine'])
subprocess.call(['coverage', 'report'])
subprocess.call([sys.executable, '-m', 'coverage', 'combine'])
subprocess.call([sys.executable, '-m', 'coverage', 'report'])
def before_feature(context, feature):
@@ -823,3 +825,5 @@ def after_feature(context, feature):
context.pctl.stop_all()
shutil.rmtree(os.path.join(context.pctl.patroni_path, 'data'))
context.dcs_ctl.cleanup_service_tree()
if feature.status == 'failed':
shutil.copytree(context.pctl.output_dir, context.pctl.output_dir + '_failed')
+2 -1
View File
@@ -5,6 +5,7 @@ import parse
import requests
import shlex
import subprocess
import sys
import time
import yaml
@@ -95,7 +96,7 @@ def do_request(context, request_method, url, data):
@step('I run {cmd}')
def do_run(context, cmd):
cmd = ['coverage', 'run', '--source=patroni', '-p'] + shlex.split(cmd)
cmd = [sys.executable, '-m', 'coverage', 'run', '--source=patroni', '-p'] + shlex.split(cmd)
try:
# XXX: Dirty hack! We need to take name/passwd from the config!
env = os.environ.copy()
+22
View File
@@ -1,3 +1,15 @@
# headless service to avoid deletion of patronidemo-config endpoint
apiVersion: v1
kind: Service
metadata:
name: patronidemo-config
labels:
application: patroni
cluster-name: patronidemo
spec:
clusterIP: None
---
apiVersion: apps/v1beta1
kind: StatefulSet
metadata:
@@ -173,6 +185,16 @@ rules:
- patch
- update
- watch
# The following privilege is only necessary for creation of headless service
# for patronidemo-config endpoint, in order to prevent cleaning it up by the
# k8s master. You can avoid giving this privilege by explicitly creating the
# service like it is done in this manifest (lines 2..10)
- apiGroups:
- ""
resources:
- services
verbs:
- create
---
apiVersion: rbac.authorization.k8s.io/v1
+33 -28
View File
@@ -272,8 +272,6 @@ class Config(object):
if authentication:
ret['postgresql']['authentication'] = authentication
users = {}
def _parse_list(value):
if not (value.strip().startswith('-') or '[' in value):
value = '[{0}]'.format(value)
@@ -285,33 +283,40 @@ class Config(object):
for param in list(os.environ.keys()):
if param.startswith(Config.PATRONI_ENV_PREFIX):
# PATRONI_(ETCD|CONSUL|ZOOKEEPER|EXHIBITOR|...)_(HOSTS?|PORT|..)
name, suffix = (param[8:].split('_', 1) + [''])[:2]
if name and suffix:
# PATRONI_(ETCD|CONSUL|ZOOKEEPER|EXHIBITOR|...)_(HOSTS?|PORT|..)
if suffix in ('HOST', 'HOSTS', 'PORT', 'SRV', 'URL', 'PROXY', 'CACERT', 'CERT', 'KEY', 'VERIFY',
'TOKEN', 'CHECKS', 'DC', 'REGISTER_SERVICE', 'SERVICE_CHECK_INTERVAL', 'NAMESPACE',
'CONTEXT', 'USE_ENDPOINTS', 'SCOPE_LABEL', 'ROLE_LABEL', 'POD_IP', 'PORTS', 'LABELS'):
value = os.environ.pop(param)
if suffix == 'PORT':
value = value and parse_int(value)
elif suffix in ('HOSTS', 'PORTS', 'CHECKS'):
value = value and _parse_list(value)
elif suffix == 'LABELS':
value = _parse_dict(value)
elif suffix == 'REGISTER_SERVICE':
value = parse_bool(value)
if value:
ret[name.lower()][suffix.lower()] = value
# PATRONI_<username>_PASSWORD=<password>, PATRONI_<username>_OPTIONS=<option1,option2,...>
# CREATE USER "<username>" WITH <OPTIONS> PASSWORD '<password>'
elif suffix == 'PASSWORD':
password = os.environ.pop(param)
if password:
users[name] = {'password': password}
options = os.environ.pop(param[:-9] + '_OPTIONS', None)
options = options and _parse_list(options)
if options:
users[name]['options'] = options
if suffix in ('HOST', 'HOSTS', 'PORT', 'PROTOCOL', 'SRV', 'URL', 'PROXY', 'CACERT', 'CERT', 'KEY',
'VERIFY', 'TOKEN', 'CHECKS', 'DC', 'REGISTER_SERVICE', 'SERVICE_CHECK_INTERVAL',
'NAMESPACE', 'CONTEXT', 'USE_ENDPOINTS', 'SCOPE_LABEL', 'ROLE_LABEL', 'POD_IP',
'PORTS', 'LABELS') and name:
value = os.environ.pop(param)
if suffix == 'PORT':
value = value and parse_int(value)
elif suffix in ('HOSTS', 'PORTS', 'CHECKS'):
value = value and _parse_list(value)
elif suffix == 'LABELS':
value = _parse_dict(value)
elif suffix == 'REGISTER_SERVICE':
value = parse_bool(value)
if value:
ret[name.lower()][suffix.lower()] = value
if 'etcd' in ret:
ret['etcd'].update(_get_auth('etcd'))
users = {}
for param in list(os.environ.keys()):
if param.startswith(Config.PATRONI_ENV_PREFIX):
name, suffix = (param[8:].rsplit('_', 1) + [''])[:2]
# PATRONI_<username>_PASSWORD=<password>, PATRONI_<username>_OPTIONS=<option1,option2,...>
# CREATE USER "<username>" WITH <OPTIONS> PASSWORD '<password>'
if name and suffix == 'PASSWORD':
password = os.environ.pop(param)
if password:
users[name] = {'password': password}
options = os.environ.pop(param[:-9] + '_OPTIONS', None)
options = options and _parse_list(options)
if options:
users[name]['options'] = options
if users:
ret['bootstrap']['users'] = users
+12 -4
View File
@@ -180,7 +180,7 @@ def print_output(columns, rows=None, alignment=None, fmt='pretty', header=True,
if fmt == 'tsv':
if columns is not None and header:
click.echo(delimiter.join(columns) + '\n')
click.echo(delimiter.join(columns))
for r in rows:
c = [str(c) for c in r]
@@ -716,6 +716,10 @@ def output_members(cluster, name, extended=False, fmt='pretty'):
has_scheduled_restarts = any(m.data.get('scheduled_restart') for m in cluster.members)
has_pending_restarts = any(m.data.get('pending_restart') for m in cluster.members)
# Show Host as 'host:port' if somebody is running on non-standard port or two nodes are running on the same host
append_port = any(str(m.conn_kwargs()['port']) != '5432' for m in cluster.members) or\
len(set(m.conn_kwargs()['host'] for m in cluster.members)) < len(cluster.members)
for m in cluster.members:
logging.debug(m)
@@ -732,7 +736,11 @@ def output_members(cluster, name, extended=False, fmt='pretty'):
elif xlog_location_cluster >= xlog_location:
lag = round((xlog_location_cluster - xlog_location)/1024/1024)
row = [name, m.name, m.conn_kwargs()['host'], role, m.data.get('state', ''), lag]
host = m.conn_kwargs()['host']
if append_port:
host += ':{0}'.format(m.conn_kwargs()['port'])
row = [name, m.name, host, role, m.data.get('state', ''), m.data.get('timeline', ''), lag]
if extended or has_pending_restarts:
row.append('*' if m.data.get('pending_restart') else '')
@@ -749,8 +757,8 @@ def output_members(cluster, name, extended=False, fmt='pretty'):
rows.append(row)
columns = ['Cluster', 'Member', 'Host', 'Role', 'State', 'Lag in MB']
alignment = {'Lag in MB': 'r'}
columns = ['Cluster', 'Member', 'Host', 'Role', 'State', 'TL', 'Lag in MB']
alignment = {'Lag in MB': 'r', 'TL': 'r'}
if extended or has_pending_restarts:
columns.append('Pending restart')
+6 -1
View File
@@ -102,7 +102,12 @@ class HTTPClient(object):
params = {k: v for k, v in params}
kwargs = {'retries': 0, 'preload_content': False, 'body': data}
if method == 'get' and isinstance(params, dict) and 'index' in params:
kwargs['timeout'] = (float(params['wait'][:-1]) if 'wait' in params else 300) + 1
timeout = float(params['wait'][:-1]) if 'wait' in params else 300
# According to the documentation a small random amount of additional wait time is added to the
# supplied maximum wait time to spread out the wake up time of any concurrent requests. This adds
# up to wait / 16 additional time to the maximum duration. Since our goal is actually getting a
# response rather read timeout we will add to the timeout a sligtly bigger value.
kwargs['timeout'] = timeout + max(timeout/15.0, 1)
else:
kwargs['timeout'] = self._read_timeout
token = params.pop('token', self.token) if isinstance(params, dict) else self.token
+28 -9
View File
@@ -92,13 +92,14 @@ class Kubernetes(AbstractDCS):
self.__subsets = None
use_endpoints = config.get('use_endpoints') and (config.get('patronictl') or 'pod_ip' in config)
if use_endpoints:
addresses = [k8s_client.V1EndpointAddress(ip=config['pod_ip'])]
addresses = [k8s_client.V1EndpointAddress(ip='127.0.0.1' if config.get('patronictl') else config['pod_ip'])]
ports = []
for p in config.get('ports', [{}]):
port = {'port': int(p.get('port', '5432'))}
port.update({n: p[n] for n in ('name', 'protocol') if p.get(n)})
ports.append(k8s_client.V1EndpointPort(**port))
self.__subsets = [k8s_client.V1EndpointSubset(addresses=addresses, ports=ports)]
self._should_create_config_service = True
self._api = CoreV1ApiProxy(use_endpoints)
self.set_retry_timeout(config['retry_timeout'])
self.set_ttl(config.get('ttl') or 30)
@@ -274,6 +275,24 @@ class Kubernetes(AbstractDCS):
body = k8s_client.V1ConfigMap(metadata=metadata)
return self.retry(func, self._namespace, body) if retry else func(self._namespace, body)
def patch_or_create_config(self, annotations, resource_version=None, patch=False, retry=True):
# SCOPE-config endpoint requires corresponding service otherwise it might be "cleaned" by k8s master
if self.__subsets and not patch and not resource_version:
self._should_create_config_service = True
self._create_config_service()
return self.patch_or_create(self.config_path, annotations, resource_version, patch, retry)
def _create_config_service(self):
metadata = k8s_client.V1ObjectMeta(namespace=self._namespace, name=self.config_path, labels=self._labels)
body = k8s_client.V1Service(metadata=metadata, spec=k8s_client.V1ServiceSpec(cluster_ip='None'))
try:
if not self._api.create_namespaced_service(self._namespace, body):
return
except Exception as e:
if not isinstance(e, k8s_client.rest.ApiException) or e.status != 409: # Service already exists
return logger.exception('create_config_service failed')
self._should_create_config_service = False
def _write_leader_optime(self, last_operation):
"""Unused"""
@@ -288,9 +307,7 @@ class Kubernetes(AbstractDCS):
if last_operation:
annotations[self._OPTIME] = last_operation
subsets = self.__subsets
if subsets is not None and access_is_restricted:
subsets = []
subsets = [] if access_is_restricted else self.__subsets
ret = self.patch_or_create(self.leader_path, annotations, self._leader_resource_version, subsets=subsets)
if ret:
@@ -333,13 +350,13 @@ class Kubernetes(AbstractDCS):
def set_config_value(self, value, index=None):
patch = bool(index or self.cluster and self.cluster.config and self.cluster.config.index)
return self.patch_or_create(self.config_path, {self._CONFIG: value}, index, patch, False)
return self.patch_or_create_config({self._CONFIG: value}, index, patch, False)
@catch_kubernetes_errors
def touch_member(self, data, ttl=None, permanent=False):
cluster = self.cluster
if cluster and cluster.leader and cluster.leader.name == self._name:
role = 'master'
role = 'promoted' if data['role'] in ('replica', 'promoted') else 'master'
elif data['state'] == 'running' and data['role'] != 'master':
role = data['role']
else:
@@ -354,12 +371,14 @@ class Kubernetes(AbstractDCS):
'annotations': {'status': json.dumps(data, separators=(',', ':'))}}
body = k8s_client.V1Pod(metadata=k8s_client.V1ObjectMeta(**metadata))
ret = self._api.patch_namespaced_pod(self._name, self._namespace, body)
if self.__subsets and self._should_create_config_service:
self._create_config_service()
return ret
def initialize(self, create_new=True, sysid=""):
cluster = self.cluster
resource_version = cluster.config.index if cluster and cluster.config and cluster.config.index else None
return self.patch_or_create(self.config_path, {self._INITIALIZE: sysid}, resource_version)
return self.patch_or_create_config({self._INITIALIZE: sysid}, resource_version)
def delete_leader(self):
if self.cluster and isinstance(self.cluster.leader, Leader) and self.cluster.leader.name == self._name:
@@ -367,7 +386,7 @@ class Kubernetes(AbstractDCS):
self.reset_cluster()
def cancel_initialization(self):
self.patch_or_create(self.config_path, {self._INITIALIZE: None}, self.cluster.config.index, True)
self.patch_or_create_config({self._INITIALIZE: None}, self.cluster.config.index, True)
@catch_kubernetes_errors
def delete_cluster(self):
@@ -375,7 +394,7 @@ class Kubernetes(AbstractDCS):
def set_history_value(self, value):
patch = bool(self.cluster and self.cluster.config and self.cluster.config.index)
return self.patch_or_create(self.config_path, {self._HISTORY: value}, None, patch, False)
return self.patch_or_create_config({self._HISTORY: value}, None, patch, False)
def set_sync_state_value(self, value, index=None):
"""Unused"""
+25 -22
View File
@@ -265,13 +265,22 @@ class Ha(object):
return result
def _handle_rewind(self):
def _handle_rewind_or_reinitialize(self):
leader = self.get_remote_master() if self.is_standby_cluster() else self.cluster.leader
if self.state_handler.rewind_needed_and_possible(leader):
if not self.state_handler.rewind_or_reinitialize_needed_and_possible(leader):
return None
if self.state_handler.can_rewind:
self._async_executor.schedule('running pg_rewind from ' + leader.name)
self._async_executor.run_async(self.state_handler.rewind, (leader,))
return True
# remove_data_directory_on_diverged_timelines is set
if not self.is_standby_cluster():
self._async_executor.schedule('reinitializing due to diverged timelines')
self._async_executor.run_async(self._do_reinitialize, args=(self.cluster, ))
return True
def recover(self):
# Postgres is not running and we will restart in standby mode. Watchdog is not needed until we promote.
self.watchdog.disable()
@@ -304,7 +313,7 @@ class Ha(object):
if self.is_standby_cluster() or not self.has_lock():
if not self.state_handler.rewind_executed:
self.state_handler.trigger_check_diverged_lsn()
if self._handle_rewind():
if self._handle_rewind_or_reinitialize():
return self._async_executor.scheduled_action
if self.has_lock(): # in standby cluster
@@ -338,9 +347,7 @@ class Ha(object):
else:
node_to_follow = cluster.leader
return (node_to_follow if
node_to_follow and
node_to_follow.name != self.state_handler.name else None)
return node_to_follow if node_to_follow and node_to_follow.name != self.state_handler.name else None
def follow(self, demote_reason, follow_reason, refresh=True):
if refresh:
@@ -351,7 +358,8 @@ class Ha(object):
node_to_follow = self._get_node_to_follow(self.cluster)
if self.is_paused():
if not (self.state_handler.need_rewind and self.state_handler.can_rewind) or self.cluster.is_unlocked():
if not (self.state_handler.need_rewind and self.state_handler.can_rewind_or_reinitialize_allowed)\
or self.cluster.is_unlocked():
self.state_handler.set_role('master' if is_leader else 'replica')
if is_leader:
return 'continue to run as master without lock'
@@ -361,7 +369,7 @@ class Ha(object):
self.demote('immediate-nolock')
return demote_reason
if self._handle_rewind():
if self._handle_rewind_or_reinitialize():
return self._async_executor.scheduled_action
if not self.state_handler.check_recovery_conf(node_to_follow):
@@ -733,7 +741,7 @@ class Ha(object):
else:
if self.is_synchronous_mode():
self.state_handler.set_synchronous_standby(None)
if self.state_handler.rewind_needed_and_possible(leader):
if self.state_handler.rewind_or_reinitialize_needed_and_possible(leader):
return False # do not start postgres, but run pg_rewind on the next iteration
self.state_handler.follow(node_to_follow)
@@ -1253,8 +1261,8 @@ class Ha(object):
# the demote code follows through to starting Postgres right away, however, in the rewind case
# it returns from demote and reaches this point to start PostgreSQL again after rewind. In that
# case it makes no sense to continue to recover() unless rewind has finished successfully.
elif (self.state_handler.rewind_failed or
not (self.state_handler.need_rewind and self.state_handler.can_rewind)):
elif self.state_handler.rewind_failed or not self.state_handler.need_rewind \
or not self.state_handler.can_rewind_or_reinitialize_allowed:
return 'postgres is not running'
# try to start dead postgres
@@ -1342,16 +1350,11 @@ class Ha(object):
if cluster_params:
unique_name = 'remote_master:{}'.format(uuid.uuid1())
data = {
'conn_kwargs': {
"host": cluster_params.get('host'),
"port": cluster_params.get('port'),
},
'no_replication_slot': 'primary_slot_name' not in cluster_params,
}
data.update({
k: v for k, v in cluster_params.items()
if k in RemoteMember.allowed_keys()
})
data = {k: v for k, v in cluster_params.items() if k in RemoteMember.allowed_keys()}
data['no_replication_slot'] = 'primary_slot_name' not in cluster_params
conn_kwargs = {k: cluster_params[k] for k in ('host', 'port') if k in cluster_params}
if conn_kwargs:
data['conn_kwargs'] = conn_kwargs
return RemoteMember(unique_name, data)
+20 -11
View File
@@ -16,7 +16,7 @@ from patroni.callback_executor import CallbackExecutor
from patroni.exceptions import PostgresConnectionException, PostgresException
from patroni.utils import compare_values, parse_bool, parse_int, Retry, RetryFailedError, polling_loop, split_host_port
from patroni.postmaster import PostmasterProcess
from patroni.dcs import slot_name_from_member_name, RemoteMember
from patroni.dcs import slot_name_from_member_name, RemoteMember, Leader
from requests.structures import CaseInsensitiveDict
from six import string_types
from six.moves.urllib.parse import quote_plus
@@ -404,6 +404,10 @@ class Postgresql(object):
return False
return self.configuration_allows_rewind(self.controldata())
@property
def can_rewind_or_reinitialize_allowed(self):
return self.config.get('remove_data_directory_on_diverged_timelines') or self.can_rewind
@property
def sysid(self):
if not self._sysid and not self.bootstrapping:
@@ -1250,9 +1254,7 @@ class Postgresql(object):
@contextmanager
def _get_replication_connection_cursor(self, host='localhost', port=5432, database=None, **kwargs):
database = database or self._database
replication = 'database' if self._major_version >= 90400 else 1
with self._get_connection_cursor(host=host, port=int(port), database=database, replication=replication,
with self._get_connection_cursor(host=host, port=int(port), database=database or self._database, replication=1,
user=self._replication['username'], password=self._replication['password'],
connect_timeout=3, options='-c statement_timeout=2000') as cur:
yield cur
@@ -1323,7 +1325,11 @@ class Postgresql(object):
if local_timeline is None or local_lsn is None:
return
if not self.check_leader_is_not_in_recovery(**leader.conn_kwargs(self._superuser)):
if isinstance(leader, Leader):
if leader.member.data.get('role') != 'master':
return
# standby cluster
elif not self.check_leader_is_not_in_recovery(**leader.conn_kwargs(self._superuser)):
return
history = need_rewind = None
@@ -1403,19 +1409,21 @@ class Postgresql(object):
else:
logger.error('Failed to rewind from healty master: %s', leader.name)
if self.config.get('remove_data_directory_on_rewind_failure', False):
logger.warning('remove_data_directory_on_rewind_failure is set. removing...')
self.remove_data_directory()
self._rewind_state = REWIND_STATUS.INITIAL
for name in ('remove_data_directory_on_rewind_failure', 'remove_data_directory_on_diverged_timelines'):
if self.config.get(name):
logger.warning('%s is set. removing...', name)
self.remove_data_directory()
self._rewind_state = REWIND_STATUS.INITIAL
break
else:
self._rewind_state = REWIND_STATUS.FAILED
return False
def trigger_check_diverged_lsn(self):
if self.can_rewind and self._rewind_state != REWIND_STATUS.NEED:
if self.can_rewind_or_reinitialize_allowed and self._rewind_state != REWIND_STATUS.NEED:
self._rewind_state = REWIND_STATUS.CHECK
def rewind_needed_and_possible(self, leader):
def rewind_or_reinitialize_needed_and_possible(self, leader):
if leader and leader.name != self.name and leader.conn_url and self._rewind_state == REWIND_STATUS.CHECK:
self._check_timeline_and_lsn(leader)
return leader and leader.conn_url and self._rewind_state == REWIND_STATUS.NEED
@@ -1652,6 +1660,7 @@ $$""".format(name, ' '.join(options)), name, password, password)
base backup)
"""
self._rewind_state = REWIND_STATUS.INITIAL
ret = self.create_replica(clone_member) == 0
if ret:
self._post_restore()
+1 -1
View File
@@ -1 +1 @@
__version__ = '1.5.4'
__version__ = '1.5.5'
+1 -1
View File
@@ -11,6 +11,6 @@ click>=4.1
prettytable>=0.7
tzlocal
python-dateutil
psutil
psutil>=2.0.0
cdiff
kubernetes>=2.0.0,<=7.0.0,!=4.0.*,!=5.0.*
+4 -1
View File
@@ -403,4 +403,7 @@ class TestRestApiServer(unittest.TestCase):
srv.reload_config({'listen': '127.0.0.2:8008'})
def test_handle_error(self):
self.assertIsNone(MockRestApiServer.handle_error(None, ('127.0.0.1', 55555)))
try:
raise Exception()
except Exception:
self.assertIsNone(MockRestApiServer.handle_error(None, ('127.0.0.1', 55555)))
+14 -4
View File
@@ -241,12 +241,20 @@ class TestHa(unittest.TestCase):
self.p.controldata = lambda: {'Database cluster state': 'in production', 'Database system identifier': SYSID}
self.assertEqual(self.ha.run_cycle(), 'doing crash recovery in a single user mode')
@patch.object(Postgresql, 'rewind_needed_and_possible', Mock(return_value=True))
@patch.object(Postgresql, 'rewind_or_reinitialize_needed_and_possible', Mock(return_value=True))
@patch.object(Postgresql, 'can_rewind', PropertyMock(return_value=True))
def test_recover_with_rewind(self):
self.p.is_running = false
self.ha.cluster = get_cluster_initialized_with_leader()
self.assertEqual(self.ha.run_cycle(), 'running pg_rewind from leader')
@patch.object(Postgresql, 'rewind_or_reinitialize_needed_and_possible', Mock(return_value=True))
@patch.object(Postgresql, 'create_replica', Mock(return_value=1))
def test_recover_with_reinitialize(self):
self.p.is_running = false
self.ha.cluster = get_cluster_initialized_with_leader()
self.assertEqual(self.ha.run_cycle(), 'reinitializing due to diverged timelines')
@patch('sys.exit', return_value=1)
@patch('patroni.ha.Ha.sysid_valid', MagicMock(return_value=True))
def test_sysid_no_match(self, exit_mock):
@@ -345,7 +353,8 @@ class TestHa(unittest.TestCase):
self.p.is_leader = false
self.assertEqual(self.ha.run_cycle(), 'PAUSE: no action')
@patch.object(Postgresql, 'rewind_needed_and_possible', Mock(return_value=True))
@patch.object(Postgresql, 'rewind_or_reinitialize_needed_and_possible', Mock(return_value=True))
@patch.object(Postgresql, 'can_rewind', PropertyMock(return_value=True))
def test_follow_triggers_rewind(self):
self.p.is_leader = false
self.p.trigger_check_diverged_lsn()
@@ -459,7 +468,7 @@ class TestHa(unittest.TestCase):
f = Failover(0, self.p.name, '', None)
self.ha.cluster = get_cluster_initialized_with_leader(f)
self.assertEqual(self.ha.run_cycle(), 'manual failover: demoting myself')
self.p.rewind_needed_and_possible = true
self.p.rewind_or_reinitialize_needed_and_possible = true
self.assertEqual(self.ha.run_cycle(), 'manual failover: demoting myself')
self.ha.fetch_node_status = get_node_status(nofailover=True)
self.assertEqual(self.ha.run_cycle(), 'no action. i am the leader with the lock')
@@ -681,7 +690,8 @@ class TestHa(unittest.TestCase):
msg = 'promoted self to a standby leader because i had the session lock'
self.assertEqual(self.ha.run_cycle(), msg)
@patch.object(Postgresql, 'rewind_needed_and_possible', Mock(return_value=True))
@patch.object(Postgresql, 'rewind_or_reinitialize_needed_and_possible', Mock(return_value=True))
@patch.object(Postgresql, 'can_rewind', PropertyMock(return_value=True))
def test_process_unhealthy_standby_cluster_as_cascade_replica(self):
self.p.is_leader = false
self.p.name = 'replica'
+13 -1
View File
@@ -72,7 +72,7 @@ class TestKubernetes(unittest.TestCase):
@patch.object(k8s_client.CoreV1Api, 'patch_namespaced_pod', Mock(return_value=True))
def test_touch_member(self):
self.k.touch_member({})
self.k.touch_member({'role': 'replica'})
self.k._name = 'p-1'
self.k.touch_member({'state': 'running', 'role': 'replica'})
self.k.touch_member({'state': 'stopped', 'role': 'master'})
@@ -112,3 +112,15 @@ class TestKubernetes(unittest.TestCase):
def test_set_history_value(self):
self.k.set_history_value('{}')
@patch('kubernetes.config.load_kube_config', Mock())
@patch.object(k8s_client.CoreV1Api, 'patch_namespaced_pod', Mock(return_value=True))
@patch.object(k8s_client.CoreV1Api, 'create_namespaced_endpoints', Mock())
@patch.object(k8s_client.CoreV1Api, 'create_namespaced_service',
Mock(side_effect=[True, False, k8s_client.rest.ApiException(500, '')]))
def test__create_config_service(self):
k = Kubernetes({'ttl': 30, 'scope': 'test', 'name': 'p-0', 'retry_timeout': 10,
'labels': {'f': 'b'}, 'use_endpoints': True, 'pod_ip': '10.0.0.0'})
self.assertIsNotNone(k.patch_or_create_config({'foo': 'bar'}))
self.assertIsNotNone(k.patch_or_create_config({'foo': 'bar'}))
k.touch_member({'state': 'running', 'role': 'replica'})
+10 -12
View File
@@ -11,8 +11,8 @@ from patroni.log import PatroniLogger
class TestPatroniLogger(unittest.TestCase):
@patch('logging.FileHandler._open', Mock())
def setUp(self):
self.config = {
def test_patroni_logger(self):
config = {
'log': {
'dir': 'foo',
'file_size': 4096,
@@ -24,15 +24,13 @@ class TestPatroniLogger(unittest.TestCase):
'restapi': {}, 'postgresql': {'data_dir': 'foo'}
}
sys.argv = ['patroni.py']
os.environ[Config.PATRONI_CONFIG_VARIABLE] = yaml.dump(self.config, default_flow_style=False)
self.logger = PatroniLogger()
config = Config()
self.logger.reload_config(config['log'])
os.environ[Config.PATRONI_CONFIG_VARIABLE] = yaml.dump(config, default_flow_style=False)
logger = PatroniLogger()
patroni_config = Config()
logger.reload_config(patroni_config['log'])
def test_rotating_handler(self):
self.assertEqual(self.logger.handler.maxBytes, self.config['log']['file_size'])
self.assertEqual(self.logger.handler.backupCount, self.config['log']['file_num'])
self.assertEqual(logger.handler.maxBytes, config['log']['file_size'])
self.assertEqual(logger.handler.backupCount, config['log']['file_num'])
def test_reload_config(self):
self.config['log'].pop('dir')
self.logger.reload_config(self.config['log'])
config['log'].pop('dir')
logger.reload_config(config['log'])
+10 -8
View File
@@ -345,10 +345,10 @@ class TestPostgresql(unittest.TestCase):
Mock(return_value={'Database cluster state': 'shut down in recovery',
'Minimum recovery ending location': '0/0',
"Min recovery ending loc's timeline": '0'})):
self.p.rewind_needed_and_possible(self.leader)
self.p.rewind_or_reinitialize_needed_and_possible(self.leader)
with patch.object(Postgresql, 'is_running', Mock(return_value=True)):
with patch.object(MockCursor, 'fetchone', Mock(side_effect=[(False, ), Exception])):
self.p.rewind_needed_and_possible(self.leader)
self.p.rewind_or_reinitialize_needed_and_possible(self.leader)
@patch.object(Postgresql, 'start', Mock())
@patch.object(Postgresql, 'can_rewind', PropertyMock(return_value=True))
@@ -357,21 +357,23 @@ class TestPostgresql(unittest.TestCase):
def test__check_timeline_and_lsn(self, mock_check_leader_is_not_in_recovery):
mock_check_leader_is_not_in_recovery.return_value = False
self.p.trigger_check_diverged_lsn()
self.assertFalse(self.p.rewind_needed_and_possible(self.leader))
self.assertFalse(self.p.rewind_or_reinitialize_needed_and_possible(self.leader))
self.leader = self.leader.member
self.assertFalse(self.p.rewind_or_reinitialize_needed_and_possible(self.leader))
mock_check_leader_is_not_in_recovery.return_value = True
self.assertFalse(self.p.rewind_needed_and_possible(self.leader))
self.assertFalse(self.p.rewind_or_reinitialize_needed_and_possible(self.leader))
self.p.trigger_check_diverged_lsn()
with patch('psycopg2.connect', Mock(side_effect=Exception)):
self.assertFalse(self.p.rewind_needed_and_possible(self.leader))
self.assertFalse(self.p.rewind_or_reinitialize_needed_and_possible(self.leader))
self.p.trigger_check_diverged_lsn()
with patch.object(MockCursor, 'fetchone', Mock(side_effect=[('', 2, '0/0'), ('', b'3\t0/40159C0\tn\n')])):
self.assertFalse(self.p.rewind_needed_and_possible(self.leader))
self.assertFalse(self.p.rewind_or_reinitialize_needed_and_possible(self.leader))
self.p.trigger_check_diverged_lsn()
with patch.object(MockCursor, 'fetchone', Mock(return_value=('', 1, '0/0'))):
with patch.object(Postgresql, '_get_local_timeline_lsn', Mock(return_value=(1, '0/0'))):
self.assertFalse(self.p.rewind_needed_and_possible(self.leader))
self.assertFalse(self.p.rewind_or_reinitialize_needed_and_possible(self.leader))
self.p.trigger_check_diverged_lsn()
self.assertTrue(self.p.rewind_needed_and_possible(self.leader))
self.assertTrue(self.p.rewind_or_reinitialize_needed_and_possible(self.leader))
@patch.object(MockCursor, 'fetchone', Mock(side_effect=[(True,), Exception]))
def test_check_leader_is_not_in_recovery(self):
+1 -1
View File
@@ -48,7 +48,7 @@ class TestPostmasterProcess(unittest.TestCase):
mock_init.side_effect = psutil.NoSuchProcess(123)
self.assertEqual(PostmasterProcess.from_pid(123), None)
mock_init.side_effect = None
self.assertNotEquals(PostmasterProcess.from_pid(123), None)
self.assertNotEqual(PostmasterProcess.from_pid(123), None)
@patch('psutil.Process.__init__', Mock())
@patch('psutil.Process.send_signal')