mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-26 07:30:14 +00:00
Compare commits
@@ -25,6 +25,15 @@ Example: defining ``PATRONI_admin_PASSWORD=strongpasswd`` and ``PATRONI_admin_OP
|
||||
Consul
|
||||
------
|
||||
- **PATRONI\_CONSUL\_HOST**: the host:port for the Consul endpoint.
|
||||
- **PATRONI\_CONSUL\_URL**: url for the Consul, in format: http(s)://host:port
|
||||
- **PATRONI\_CONSUL\_PORT**: (optional) Consul port
|
||||
- **PATRONI\_CONSUL\_SCHEME**: (optional) **http** or **https**, defaults to **http**
|
||||
- **PATRONI\_CONSUL\_TOKEN**: (optional) ACL token
|
||||
- **PATRONI\_CONSUL\_VERIFY**: (optional) whether to verify the SSL certificate for HTTPS requests
|
||||
- **PATRONI\_CONSUL\_CACERT**: (optional) The ca certificate. If pressent it will enable validation.
|
||||
- **PATRONI\_CONSUL\_CERT**: (optional) File with the client certificate
|
||||
- **PATRONI\_CONSUL\_KEY**: (optional) File with the client key. Can be empty if the key is part of certificate.
|
||||
- **PATRONI\_CONSUL\_CHECKS**: (optional) list of Consul health checks used for the session. If not specified Consul will use "serfHealth" in additional to the TTL based check created by Patroni. Additional checks, in particular the "serfHealth", may cause the leader lock to expire faster than in `ttl` seconds when the leader instance becomes unavailable.
|
||||
|
||||
Etcd
|
||||
----
|
||||
|
||||
+15
-1
@@ -43,13 +43,27 @@ Bootstrap configuration
|
||||
- **- createdb**
|
||||
- **post\_bootstrap** or **post\_init**: An additional script that will be executed after initializing the cluster. The script receives a connection string URL (with the cluster superuser as a user name). The PGPASSFILE variable is set to the location of pgpass file.
|
||||
|
||||
.. _consul_settings:
|
||||
|
||||
Consul
|
||||
------
|
||||
- **host**: the host:port for the Consul endpoint.
|
||||
Most of the parameters are optional, but you have to specify one of the **host** or **url**
|
||||
|
||||
- **host**: the host:port for the Consul endpoint, in format: http(s)://host:port
|
||||
- **url**: url for the Consul endpoint
|
||||
- **port**: (optional) Consul port
|
||||
- **scheme**: (optional) **http** or **https**, defaults to **http**
|
||||
- **token**: (optional) ACL token
|
||||
- **verify** (optional) whether to verify the SSL certificate for HTTPS requests
|
||||
- **cacert**: (optional) The ca certificate. If pressent it will enable validation.
|
||||
- **cert**: (optional) file with the client certificate
|
||||
- **key**: (optional) file with the client key. Can be empty if the key is part of **cert**.
|
||||
- **checks**: (optional) list of Consul health checks used for the session. If not specified Consul will use "serfHealth" in additional to the TTL based check created by Patroni. Additional checks, in particular the "serfHealth", may cause the leader lock to expire faster than in `ttl` seconds when the leader instance becomes unavailable
|
||||
|
||||
Etcd
|
||||
----
|
||||
Most of the parameters are optional, but you have to specify one of the **host**, **url**, **proxy** or **srv**
|
||||
|
||||
- **host**: the host:port for the etcd endpoint.
|
||||
- **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**
|
||||
|
||||
@@ -3,6 +3,101 @@
|
||||
Release notes
|
||||
=============
|
||||
|
||||
Version 1.3.5
|
||||
-------------
|
||||
|
||||
**Bugfix**
|
||||
|
||||
- Set role to 'uninitialized' if data directory was removed (Alexander Kukushkin)
|
||||
|
||||
If the node was running as a master it was preventing from failover.
|
||||
|
||||
**Stability improvement**
|
||||
|
||||
- Try to run postmaster in a single-user mode if we tried and failed to start postgres (Alexander)
|
||||
|
||||
Usually such problem happens when node running as a master was terminated and timelines were diverged.
|
||||
If ``recovery.conf`` has ``restore_command`` defined, there are really high chances that postgres will abort startup and leave controldata unchanged.
|
||||
It makes impossible to use ``pg_rewind``, which requires a clean shutdown.
|
||||
|
||||
**Consul improvements**
|
||||
|
||||
- Make it possible to specify health checks when creating session (Alexander)
|
||||
|
||||
If not specified, Consul will use "serfHealth". From one side it allows fast detection of isolated master, but from another side it makes it impossible for Patroni to tolerate short network lags.
|
||||
|
||||
**Bugfix**
|
||||
|
||||
- Fix watchdog on Python 3 (Ants Aasma)
|
||||
|
||||
A misunderstanding of the ioctl() call interface. If mutable=False then fcntl.ioctl() actually returns the arg buffer back.
|
||||
This accidentally worked on Python2 because int and str comparison did not return an error.
|
||||
Error reporting is actually done by raising IOError on Python2 and OSError on Python3.
|
||||
|
||||
Version 1.3.4
|
||||
-------------
|
||||
|
||||
**Different Consul improvements**
|
||||
|
||||
- Pass the consul token as a header (Andrew Colin Kissa)
|
||||
|
||||
Headers are now the prefered way to pass the token to the consul `API <https://www.consul.io/api/index.html#authentication>`__.
|
||||
|
||||
|
||||
- Advanced configuration for Consul (Alexander Kukushkin)
|
||||
|
||||
possibility to specify ``scheme``, ``token``, client and ca certificates :ref:`details <consul_settings>`.
|
||||
|
||||
- compatibility with python-consul-0.7.1 and above (Alexander)
|
||||
|
||||
new python-consul module has changed signature of some methods
|
||||
|
||||
- "Could not take out TTL lock" message was never logged (Alexander)
|
||||
|
||||
Not a critical bug, but lack of proper logging complicates investigation in case of problems.
|
||||
|
||||
|
||||
**Quote synchronous_standby_names using quote_ident**
|
||||
|
||||
- When writing ``synchronous_standby_names`` into the ``postgresql.conf`` its value must be quoted (Alexander)
|
||||
|
||||
If it is not quoted properly, PostgreSQL will effectively disable synchronous replication and continue to work.
|
||||
|
||||
|
||||
**Different bugfixes around pause state, mostly related to watchdog** (Alexander)
|
||||
|
||||
- Do not send keepalives if watchdog is not active
|
||||
- Avoid activating watchdog in a pause mode
|
||||
- Set correct postgres state in pause mode
|
||||
- Do not try to run queries from API if postgres is stopped
|
||||
|
||||
|
||||
Version 1.3.3
|
||||
-------------
|
||||
|
||||
**Bugfixes**
|
||||
|
||||
- synchronous replication was disabled shortly after promotion even when synchronous_mode_strict was turned on (Alexander Kukushkin)
|
||||
- create empty ``pg_ident.conf`` file if it is missing after restoring from the backup (Alexander)
|
||||
- open access in ``pg_hba.conf`` to all databases, not only postgres (Franco Bellagamba)
|
||||
|
||||
|
||||
Version 1.3.2
|
||||
-------------
|
||||
|
||||
**Bugfix**
|
||||
|
||||
- patronictl edit-config didn't work with ZooKeeper (Alexander Kukushkin)
|
||||
|
||||
|
||||
Version 1.3.1
|
||||
-------------
|
||||
|
||||
**Bugfix**
|
||||
|
||||
- failover via API was broken due to change in ``_MemberStatus`` (Alexander Kukushkin)
|
||||
|
||||
|
||||
Version 1.3
|
||||
-----------
|
||||
|
||||
|
||||
@@ -34,9 +34,9 @@ Scenario: check local configuration reload
|
||||
Then I receive a response code 202
|
||||
|
||||
Scenario: check dynamic configuration change via DCS
|
||||
Given I issue a PATCH request to http://127.0.0.1:8008/config with {"ttl": 10, "loop_wait": 2, "postgresql": {"parameters": {"max_connections": 101}}}
|
||||
Then I receive a response code 200
|
||||
And I receive a response loop_wait 2
|
||||
Given I run patronictl.py edit-config -s 'ttl=10' -s 'loop_wait=2' -p 'max_connections=101' --force batman
|
||||
Then I receive a response returncode 0
|
||||
And I receive a response output "+loop_wait: 2"
|
||||
And Response on GET http://127.0.0.1:8008/patroni contains pending_restart after 11 seconds
|
||||
When I issue a GET request to http://127.0.0.1:8008/config
|
||||
Then I receive a response code 200
|
||||
@@ -65,8 +65,8 @@ Scenario: check API requests for the primary-replica pair in the pause mode
|
||||
Then postgres1 role is the secondary after 15 seconds
|
||||
|
||||
Scenario: check the failover via the API in the pause mode
|
||||
Given I run patronictl.py failover batman --master postgres0 --candidate postgres1 --force
|
||||
Then I receive a response returncode 0
|
||||
Given I issue a POST request to http://127.0.0.1:8008/failover with {"leader": "postgres0", "candidate": "postgres1"}
|
||||
Then I receive a response code 200
|
||||
And postgres1 is a leader after 5 seconds
|
||||
And postgres1 role is the primary after 10 seconds
|
||||
And postgres0 role is the secondary after 10 seconds
|
||||
|
||||
@@ -31,7 +31,7 @@ def watchdog_was_closed(context, name):
|
||||
assert context.pctl.get_watchdog(name).was_closed
|
||||
|
||||
|
||||
@step('I wait for next {name:w} watchdog ping')
|
||||
@step('I reset {name:w} watchdog state')
|
||||
def watchdog_reset_pinged(context, name):
|
||||
context.pctl.get_watchdog(name).reset()
|
||||
|
||||
|
||||
@@ -1,19 +1,31 @@
|
||||
Feature: watchdog
|
||||
Verify that watchdog gets pinged and triggered under appropriate circumstances.
|
||||
|
||||
Scenario: watchdog is opened, pinged and closed
|
||||
Scenario: watchdog is opened and pinged
|
||||
Given I start postgres0 with watchdog
|
||||
Then postgres0 is a leader after 10 seconds
|
||||
And postgres0 role is the primary after 10 seconds
|
||||
And postgres0 watchdog has been pinged after 10 seconds
|
||||
When I shut down postgres0
|
||||
|
||||
Scenario: watchdog is disabled during pause
|
||||
Given I run patronictl.py pause batman
|
||||
Then I receive a response returncode 0
|
||||
When I sleep for 2 seconds
|
||||
Then postgres0 watchdog has been closed
|
||||
|
||||
#TODO: test watchdog is disabled during pause
|
||||
#TODO: test watchdog is disabled properly when shutting down
|
||||
Scenario: watchdog is opened and pinged after resume
|
||||
Given I reset postgres0 watchdog state
|
||||
And I run patronictl.py resume batman
|
||||
Then I receive a response returncode 0
|
||||
And postgres0 watchdog has been pinged after 10 seconds
|
||||
|
||||
Scenario: watchdog is disabled when shutting down
|
||||
Given I shut down postgres0
|
||||
Then postgres0 watchdog has been closed
|
||||
|
||||
Scenario: watchdog is triggered if patroni stops responding
|
||||
Given I start postgres0 with watchdog
|
||||
Given I reset postgres0 watchdog state
|
||||
And I start postgres0 with watchdog
|
||||
Then postgres0 role is the primary after 10 seconds
|
||||
When postgres0 hangs for 30 seconds
|
||||
Then postgres0 watchdog is triggered after 30 seconds
|
||||
|
||||
@@ -293,9 +293,15 @@ class RestApiHandler(BaseHTTPRequestHandler):
|
||||
if leader and (not cluster.leader or cluster.leader.name != leader):
|
||||
return 'leader name does not match'
|
||||
if candidate:
|
||||
if cluster.is_synchronous_mode() and cluster.sync.sync_standby != candidate:
|
||||
return 'candidate name does not match with sync_standby'
|
||||
members = [m for m in cluster.members if m.name == candidate]
|
||||
if not members:
|
||||
return 'candidate does not exists'
|
||||
elif cluster.is_synchronous_mode():
|
||||
members = [m for m in cluster.members if m.name == cluster.sync.sync_standby]
|
||||
if not members:
|
||||
return 'failover is not possible: can not find sync_standby'
|
||||
else:
|
||||
members = [m for m in cluster.members if m.name != cluster.leader.name and m.api_url]
|
||||
if not members:
|
||||
@@ -376,6 +382,8 @@ class RestApiHandler(BaseHTTPRequestHandler):
|
||||
|
||||
def get_postgresql_status(self, retry=False):
|
||||
try:
|
||||
if self.server.patroni.postgresql.state not in ('running', 'restarting', 'starting'):
|
||||
raise RetryFailedError('')
|
||||
row = self.query("""WITH replication_info AS (
|
||||
SELECT usename, application_name, client_addr, state, sync_state, sync_priority
|
||||
FROM pg_stat_replication
|
||||
|
||||
+3
-3
@@ -242,12 +242,12 @@ class Config(object):
|
||||
name, suffix = (param[8:].rsplit('_', 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') \
|
||||
and '_' not in name:
|
||||
if suffix in ('HOST', 'HOSTS', 'PORT', 'SRV', 'URL', 'PROXY', 'CACERT', 'CERT', 'KEY',
|
||||
'VERIFY', 'TOKEN', 'CHECKS') and '_' not in name:
|
||||
value = os.environ.pop(param)
|
||||
if suffix == 'PORT':
|
||||
value = value and parse_int(value)
|
||||
elif suffix == 'HOSTS':
|
||||
elif suffix in ('HOSTS', 'CHECKS'):
|
||||
value = value and _parse_list(value)
|
||||
if value:
|
||||
ret[name.lower()][suffix.lower()] = value
|
||||
|
||||
+5
-5
@@ -903,14 +903,14 @@ def apply_config_changes(before_editing, data, kvpairs):
|
||||
if prefix == ('postgresql', 'parameters'):
|
||||
path = ['.'.join(path)]
|
||||
|
||||
key = path[0]
|
||||
if len(path) == 1:
|
||||
if value is None:
|
||||
config.pop(path[0], None)
|
||||
config.pop(key, None)
|
||||
else:
|
||||
config[path[0]] = value
|
||||
config[key] = value
|
||||
else:
|
||||
key = path[0]
|
||||
if key not in config:
|
||||
if not isinstance(config.get(key), dict):
|
||||
config[key] = {}
|
||||
set_path_value(config[key], path[1:], value, prefix + (key,))
|
||||
if config[key] == {}:
|
||||
@@ -1017,7 +1017,7 @@ def edit_config(obj, cluster_name, force, quiet, kvpairs, pgkvpairs, apply_filen
|
||||
return
|
||||
|
||||
if force or click.confirm('Apply these changes?'):
|
||||
if not dcs.set_config_value(json.dumps(changed_data), cluster.config.modify_index):
|
||||
if not dcs.set_config_value(json.dumps(changed_data), cluster.config.index):
|
||||
raise PatroniCtlException("Config modification aborted due to concurrent changes")
|
||||
click.echo("Configuration changed")
|
||||
|
||||
|
||||
@@ -320,6 +320,12 @@ class Cluster(namedtuple('Cluster', 'initialize,config,leader,last_leader_operat
|
||||
def is_paused(self):
|
||||
return self.config and self.config.data.get('pause', False) or False
|
||||
|
||||
def is_synchronous_mode(self):
|
||||
return bool(self.config and self.config.data.get('synchronous_mode'))
|
||||
|
||||
def is_synchronous_mode_strict(self):
|
||||
return bool(self.config and self.config.data.get('synchronous_mode_strict'))
|
||||
|
||||
|
||||
@six.add_metaclass(abc.ABCMeta)
|
||||
class AbstractDCS(object):
|
||||
|
||||
+65
-17
@@ -2,15 +2,16 @@ from __future__ import absolute_import
|
||||
import logging
|
||||
import os
|
||||
import socket
|
||||
import ssl
|
||||
import time
|
||||
import urllib3
|
||||
|
||||
from consul import ConsulException, NotFound, base
|
||||
from patroni.dcs import AbstractDCS, ClusterConfig, Cluster, Failover, Leader, Member, SyncState
|
||||
from patroni.exceptions import DCSError
|
||||
from patroni.utils import Retry, RetryFailedError
|
||||
from patroni.utils import parse_bool, Retry, RetryFailedError
|
||||
from urllib3.exceptions import HTTPError
|
||||
from six.moves.urllib.parse import urlencode
|
||||
from six.moves.urllib.parse import urlencode, urlparse
|
||||
from six.moves.http_client import HTTPException
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -26,14 +27,23 @@ class ConsulInternalError(ConsulException):
|
||||
|
||||
class HTTPClient(object):
|
||||
|
||||
def __init__(self, host='127.0.0.1', port=8500, scheme='http', verify=True, timeout=10):
|
||||
self.host = host
|
||||
self.port = port
|
||||
self.scheme = scheme
|
||||
self.verify = verify
|
||||
self.set_read_timeout(timeout)
|
||||
self.base_uri = '{0}://{1}:{2}'.format(self.scheme, self.host, self.port)
|
||||
self.http = urllib3.PoolManager(num_pools=10)
|
||||
def __init__(self, host='127.0.0.1', port=8500, scheme='http', verify=True, cert=None, ca_cert=None):
|
||||
self._read_timeout = 10
|
||||
self.base_uri = '{0}://{1}:{2}'.format(scheme, host, port)
|
||||
kwargs = {}
|
||||
if cert:
|
||||
if isinstance(cert, tuple):
|
||||
# Key and cert are separate
|
||||
kwargs['cert_file'] = cert[0]
|
||||
kwargs['key_file'] = cert[1]
|
||||
else:
|
||||
# combined certificate
|
||||
kwargs['cert_file'] = cert
|
||||
if ca_cert:
|
||||
kwargs['ca_certs'] = ca_cert
|
||||
if verify or ca_cert:
|
||||
kwargs['cert_reqs'] = ssl.CERT_REQUIRED
|
||||
self.http = urllib3.PoolManager(num_pools=10, **kwargs)
|
||||
self._ttl = None
|
||||
|
||||
def set_read_timeout(self, timeout):
|
||||
@@ -72,15 +82,26 @@ class HTTPClient(object):
|
||||
kwargs['timeout'] = (float(params['wait'][:-1]) if 'wait' in params else 300) + 1
|
||||
else:
|
||||
kwargs['timeout'] = self._read_timeout
|
||||
if isinstance(params, dict) and 'token' in params and params['token']:
|
||||
kwargs['headers'] = {'X-Consul-Token': params.pop('token')}
|
||||
return callback(self.response(self.http.request(method.upper(), self.uri(path, params), **kwargs)))
|
||||
return wrapper
|
||||
|
||||
|
||||
class ConsulClient(base.Consul):
|
||||
|
||||
@staticmethod
|
||||
def connect(host, port, scheme, verify=True):
|
||||
return HTTPClient(host, port, scheme, verify)
|
||||
def __init__(self, *args, **kwargs):
|
||||
self._cert = kwargs.pop('cert', None)
|
||||
self._ca_cert = kwargs.pop('ca_cert', None)
|
||||
super(ConsulClient, self).__init__(*args, **kwargs)
|
||||
|
||||
def connect(self, *args, **kwargs):
|
||||
kwargs.update(dict(zip(['host', 'port', 'scheme', 'verify'], args)))
|
||||
if self._cert:
|
||||
kwargs['cert'] = self._cert
|
||||
if self._ca_cert:
|
||||
kwargs['ca_cert'] = self._ca_cert
|
||||
return HTTPClient(**kwargs)
|
||||
|
||||
|
||||
def catch_consul_errors(func):
|
||||
@@ -104,11 +125,35 @@ class Consul(AbstractDCS):
|
||||
HTTPError, socket.error, socket.timeout))
|
||||
|
||||
self._my_member_data = None
|
||||
host, port = config.get('host', '127.0.0.1:8500').split(':')
|
||||
self._client = ConsulClient(host=host, port=port)
|
||||
kwargs = {}
|
||||
if 'url' in config:
|
||||
r = urlparse(config['url'])
|
||||
config.update({'scheme': r.scheme, 'host': r.hostname, 'port': r.port or 8500})
|
||||
elif 'host' in config:
|
||||
host, port = (config.get('host', '127.0.0.1:8500') + ':8500').split(':')[:2]
|
||||
config['host'] = host
|
||||
if 'port' not in config:
|
||||
config['port'] = int(port)
|
||||
|
||||
if config.get('cacert'):
|
||||
config['ca_cert'] = config.pop('cacert')
|
||||
|
||||
if config.get('key') and config.get('cert'):
|
||||
config['cert'] = (config['cert'], config['key'])
|
||||
|
||||
kwargs = {p: config.get(p) for p in ('host', 'port', 'token', 'scheme', 'cert', 'ca_cert') if config.get(p)}
|
||||
|
||||
verify = config.get('verify')
|
||||
if not isinstance(verify, bool):
|
||||
verify = parse_bool(verify)
|
||||
if isinstance(verify, bool):
|
||||
kwargs['verify'] = verify
|
||||
|
||||
self._client = ConsulClient(**kwargs)
|
||||
self.set_retry_timeout(config['retry_timeout'])
|
||||
self.set_ttl(config.get('ttl') or 30)
|
||||
self._last_session_refresh = 0
|
||||
self.__session_checks = config.get('checks')
|
||||
if not self._ctl:
|
||||
self.create_session()
|
||||
|
||||
@@ -145,6 +190,7 @@ class Consul(AbstractDCS):
|
||||
ret = not self._session
|
||||
if ret:
|
||||
self._session = self._client.session.create(name=self._scope + '-' + self._name,
|
||||
checks=self.__session_checks,
|
||||
lock_delay=0.001, behavior='delete')
|
||||
self._last_session_refresh = time.time()
|
||||
return ret
|
||||
@@ -245,12 +291,14 @@ class Consul(AbstractDCS):
|
||||
return False
|
||||
|
||||
@catch_consul_errors
|
||||
def _do_attempt_to_acquire_leader(self, kwargs):
|
||||
return self.retry(self._client.kv.put, self.leader_path, self._name, **kwargs)
|
||||
|
||||
def attempt_to_acquire_leader(self, permanent=False):
|
||||
if not self._session and not permanent:
|
||||
self.refresh_session()
|
||||
|
||||
args = {} if permanent else {'acquire': self._session}
|
||||
ret = self.retry(self._client.kv.put, self.leader_path, self._name, **args)
|
||||
ret = self._do_attempt_to_acquire_leader({} if permanent else {'acquire': self._session})
|
||||
if not ret:
|
||||
logger.info('Could not take out TTL lock')
|
||||
return ret
|
||||
|
||||
+41
-16
@@ -184,7 +184,10 @@ class Ha(object):
|
||||
if timeout == 0:
|
||||
# We are requested to prefer failing over to restarting master. But see first if there
|
||||
# is anyone to fail over to.
|
||||
if self.is_failover_possible(self.cluster.members):
|
||||
members = self.cluster.members
|
||||
if self.is_synchronous_mode():
|
||||
members = [m for m in members if self.cluster.sync.matches(m.name)]
|
||||
if self.is_failover_possible(members):
|
||||
logger.info("Master crashed. Failing over.")
|
||||
self.demote('immediate')
|
||||
return 'stopped PostgreSQL to fail over after a crash'
|
||||
@@ -204,6 +207,16 @@ class Ha(object):
|
||||
msg = "starting as a secondary"
|
||||
node_to_follow = self._get_node_to_follow(self.cluster)
|
||||
|
||||
# once we already tried to start postgres but failed, single user mode is a rescue in this case
|
||||
if self.recovering and not self.state_handler.rewind_executed and self.state_handler.can_rewind:
|
||||
data = self.state_handler.controldata()
|
||||
if data.get('Database cluster state') not in ('shut down', 'shut down in recovery'):
|
||||
self.recovering = False
|
||||
msg = 'fixing cluster state in a single user mode'
|
||||
self._async_executor.schedule(msg)
|
||||
self._async_executor.run_async(self.state_handler.fix_cluster_state)
|
||||
return msg
|
||||
|
||||
self.recovering = True
|
||||
|
||||
self._async_executor.schedule('restarting after failure')
|
||||
@@ -249,10 +262,10 @@ class Ha(object):
|
||||
return follow_reason
|
||||
|
||||
def is_synchronous_mode(self):
|
||||
return bool(self.cluster and self.cluster.config and self.cluster.config.data.get('synchronous_mode'))
|
||||
return bool(self.cluster and self.cluster.is_synchronous_mode())
|
||||
|
||||
def is_synchronous_mode_strict(self):
|
||||
return bool(self.cluster and self.cluster.config and self.cluster.config.data.get('synchronous_mode_strict'))
|
||||
return bool(self.cluster and self.cluster.is_synchronous_mode_strict())
|
||||
|
||||
def process_sync_replication(self):
|
||||
"""Process synchronous standby beahvior.
|
||||
@@ -341,14 +354,13 @@ class Ha(object):
|
||||
self._disable_sync -= 1
|
||||
|
||||
def enforce_master_role(self, message, promote_message):
|
||||
if not self.watchdog.is_running:
|
||||
if not self.watchdog.activate():
|
||||
if self.state_handler.is_leader():
|
||||
self.demote('immediate')
|
||||
return 'Demoting self because watchdog could not be activated'
|
||||
else:
|
||||
self.release_leader_key_voluntarily()
|
||||
return 'Not promoting self because watchdog could not be actived'
|
||||
if not self.is_paused() and not self.watchdog.is_running and not self.watchdog.activate():
|
||||
if self.state_handler.is_leader():
|
||||
self.demote('immediate')
|
||||
return 'Demoting self because watchdog could not be activated'
|
||||
else:
|
||||
self.release_leader_key_voluntarily()
|
||||
return 'Not promoting self because watchdog could not be activated'
|
||||
|
||||
if self.state_handler.is_leader() or self.state_handler.role == 'master':
|
||||
# Inform the state handler about its master role.
|
||||
@@ -364,7 +376,7 @@ class Ha(object):
|
||||
# Somebody else updated sync state, it may be due to us losing the lock. To be safe, postpone
|
||||
# promotion until next cycle. TODO: trigger immediate retry of run_cycle
|
||||
return 'Postponing promotion because synchronous replication state was updated by somebody else'
|
||||
self.state_handler.set_synchronous_standby(None)
|
||||
self.state_handler.set_synchronous_standby('*' if self.is_synchronous_mode_strict() else None)
|
||||
self.state_handler.promote()
|
||||
return promote_message
|
||||
|
||||
@@ -549,8 +561,8 @@ class Ha(object):
|
||||
self.state_handler.set_role('demoted')
|
||||
|
||||
if mode_control['release']:
|
||||
self.release_leader_key_voluntarily()
|
||||
time.sleep(2) # Give a time to somebody to take the leader lock
|
||||
self.release_leader_key_voluntarily()
|
||||
time.sleep(2) # Give a time to somebody to take the leader lock
|
||||
if mode_control['offline']:
|
||||
node_to_follow, leader = None, None
|
||||
else:
|
||||
@@ -564,6 +576,8 @@ class Ha(object):
|
||||
self._async_executor.schedule('starting after demotion')
|
||||
self._async_executor.run_async(self.state_handler.follow, (node_to_follow,))
|
||||
else:
|
||||
if self.is_synchronous_mode():
|
||||
self.state_handler.set_synchronous_standby(None)
|
||||
if self.state_handler.rewind_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)
|
||||
@@ -622,8 +636,16 @@ class Ha(object):
|
||||
if not failover.candidate and self.is_paused():
|
||||
logger.warning('Failover is possible only to a specific candidate in a paused state')
|
||||
else:
|
||||
members = [m for m in self.cluster.members
|
||||
if not failover.candidate or m.name == failover.candidate]
|
||||
if self.is_synchronous_mode():
|
||||
if failover.candidate and not self.cluster.sync.matches(failover.candidate):
|
||||
logger.warning('Failover candidate=%s does not match with sync_standby=%s',
|
||||
failover.candidate, self.cluster.sync.sync_standby)
|
||||
members = []
|
||||
else:
|
||||
members = [m for m in self.cluster.members if self.cluster.sync.matches(m.name)]
|
||||
else:
|
||||
members = [m for m in self.cluster.members
|
||||
if not failover.candidate or m.name == failover.candidate]
|
||||
if self.is_failover_possible(members): # check that there are healthy members
|
||||
self._async_executor.schedule('manual failover: demote')
|
||||
self._async_executor.run_async(self.demote, ('graceful',))
|
||||
@@ -915,6 +937,8 @@ class Ha(object):
|
||||
# Check if we are in startup, when paused defer to main loop for manual failovers.
|
||||
if not self.state_handler.check_for_startup() or self.is_paused():
|
||||
self.set_start_timeout(None)
|
||||
if self.is_paused():
|
||||
self.state_handler.set_state(self.state_handler.is_running() and 'running' or 'stopped')
|
||||
return None
|
||||
|
||||
# state_handler.state == 'starting' here
|
||||
@@ -990,6 +1014,7 @@ class Ha(object):
|
||||
|
||||
# is data directory empty?
|
||||
if self.state_handler.data_directory_empty():
|
||||
self.state_handler.set_role('uninitialized')
|
||||
# In case datadir went away while we were master. TODO: check for this and try to stop postgresql.
|
||||
self.watchdog.disable()
|
||||
|
||||
|
||||
+86
-11
@@ -17,7 +17,7 @@ from contextlib import contextmanager
|
||||
from patroni import call_self
|
||||
from patroni.callback_executor import CallbackExecutor
|
||||
from patroni.exceptions import PostgresConnectionException
|
||||
from patroni.utils import compare_values, parse_bool, parse_int, Retry, RetryFailedError, polling_loop, null_context
|
||||
from patroni.utils import compare_values, parse_bool, parse_int, Retry, RetryFailedError, polling_loop
|
||||
from six import string_types
|
||||
from six.moves.urllib.parse import quote_plus
|
||||
from threading import current_thread, Lock
|
||||
@@ -42,6 +42,12 @@ STOP_SIGNALS = {
|
||||
}
|
||||
STOP_POLLING_INTERVAL = 1
|
||||
REWIND_STATUS = type('Enum', (), {'INITIAL': 0, 'CHECK': 1, 'NEED': 2, 'NOT_NEED': 3, 'SUCCESS': 4, 'FAILED': 5})
|
||||
sync_standby_name_re = re.compile('^[A-Za-z_][A-Za-z_0-9\$]*$')
|
||||
|
||||
|
||||
def quote_ident(value):
|
||||
"""Very simplified version of quote_ident"""
|
||||
return value if sync_standby_name_re.match(value) else '"' + value + '"'
|
||||
|
||||
|
||||
def slot_name_from_member_name(member_name):
|
||||
@@ -60,6 +66,11 @@ def slot_name_from_member_name(member_name):
|
||||
return slot_name[0:63]
|
||||
|
||||
|
||||
@contextmanager
|
||||
def null_context():
|
||||
yield
|
||||
|
||||
|
||||
class Postgresql(object):
|
||||
|
||||
# List of parameters which must be always passed to postmaster as command line options
|
||||
@@ -1137,7 +1148,7 @@ class Postgresql(object):
|
||||
with open(self._pg_hba_conf, 'w') as f:
|
||||
f.write(self._CONFIG_WARNING_HEADER)
|
||||
for address, t in addresses.items():
|
||||
f.write('{0}\t{1}\t{2}\t{3}\ttrust\n'.format(t, self._database,
|
||||
f.write('{0}\t{1}\t{2}\t{3}\ttrust\n'.format(t, 'all',
|
||||
self._superuser.get('username') or 'all', address))
|
||||
elif not self._server_parameters.get('hba_file') and self.config.get('pg_hba'):
|
||||
with open(self._pg_hba_conf, 'w') as f:
|
||||
@@ -1250,12 +1261,14 @@ class Postgresql(object):
|
||||
else: # otherwise analyze pg_controldata output
|
||||
data = self.controldata()
|
||||
try:
|
||||
if data.get('Database cluster state') == 'shut down in recovery':
|
||||
lsn = data.get('Minimum recovery ending location')
|
||||
timeline = int(data.get("Min recovery ending loc's timeline"))
|
||||
if lsn == '0/0' or timeline == 0: # it was a master when it crashed
|
||||
data['Database cluster state'] = 'shut down'
|
||||
if data.get('Database cluster state') == 'shut down':
|
||||
lsn = data.get('Latest checkpoint location')
|
||||
timeline = int(data.get("Latest checkpoint's TimeLineID"))
|
||||
elif data.get('Database cluster state') == 'shut down in recovery':
|
||||
lsn = data.get('Minimum recovery ending location')
|
||||
timeline = int(data.get("Min recovery ending loc's timeline"))
|
||||
except (TypeError, ValueError):
|
||||
logger.exception('Failed to get local timeline and lsn from pg_controldata output')
|
||||
logger.info('Local timeline=%s lsn=%s', timeline, lsn)
|
||||
@@ -1396,8 +1409,12 @@ class Postgresql(object):
|
||||
for f in self._configuration_to_save:
|
||||
config_file = os.path.join(self._config_dir, f)
|
||||
backup_file = os.path.join(self._data_dir, f + '.backup')
|
||||
if not os.path.isfile(config_file) and os.path.isfile(backup_file):
|
||||
shutil.copy(backup_file, config_file)
|
||||
if not os.path.isfile(config_file):
|
||||
if os.path.isfile(backup_file):
|
||||
shutil.copy(backup_file, config_file)
|
||||
# Previously we didn't backup pg_ident.conf, if file is missing just create empty
|
||||
elif f == 'pg_ident.conf':
|
||||
open(config_file, 'w').close()
|
||||
except IOError:
|
||||
logger.exception('unable to restore configuration files from backup')
|
||||
|
||||
@@ -1648,12 +1665,13 @@ $$""".format(name, ' '.join(options)), name, password, password)
|
||||
:returns tuple of candidate name or None, and bool showing if the member is the active synchronous standby.
|
||||
"""
|
||||
current = cluster.sync.sync_standby
|
||||
members = {m.name: m for m in cluster.members}
|
||||
current = current.lower() if current else current
|
||||
members = {m.name.lower(): m for m in cluster.members}
|
||||
candidates = []
|
||||
# Pick candidates based on who has flushed WAL farthest.
|
||||
# TODO: for synchronous_commit = remote_write we actually want to order on write_location
|
||||
for app_name, state, sync_state in self.query(
|
||||
"""SELECT application_name, state, sync_state
|
||||
"""SELECT LOWER(application_name), state, sync_state
|
||||
FROM pg_stat_replication
|
||||
ORDER BY flush_{0} DESC""".format(self.lsn_name)):
|
||||
member = members.get(app_name)
|
||||
@@ -1673,14 +1691,17 @@ $$""".format(name, ' '.join(options)), name, password, password)
|
||||
|
||||
def set_synchronous_standby(self, name):
|
||||
"""Sets a node to be synchronous standby and if changed does a reload for PostgreSQL."""
|
||||
if name and name != '*':
|
||||
name = quote_ident(name)
|
||||
if name != self._synchronous_standby_names:
|
||||
if name is None:
|
||||
self._server_parameters.pop('synchronous_standby_names', None)
|
||||
else:
|
||||
self._server_parameters['synchronous_standby_names'] = name
|
||||
self._synchronous_standby_names = name
|
||||
self._write_postgresql_conf()
|
||||
self.reload()
|
||||
if self.state == 'running':
|
||||
self._write_postgresql_conf()
|
||||
self.reload()
|
||||
|
||||
@staticmethod
|
||||
def postgres_version_to_int(pg_version):
|
||||
@@ -1725,3 +1746,57 @@ $$""".format(name, ' '.join(options)), name, password, password)
|
||||
90600
|
||||
"""
|
||||
return Postgresql.postgres_version_to_int(pg_version + '.0')
|
||||
|
||||
def read_postmaster_opts(self):
|
||||
"""returns the list of option names/values from postgres.opts, Empty dict if read failed or no file"""
|
||||
result = {}
|
||||
try:
|
||||
with open(os.path.join(self._data_dir, 'postmaster.opts')) as f:
|
||||
data = f.read()
|
||||
for opt in data.split('" "'):
|
||||
if '=' in opt and opt.startswith('--'):
|
||||
name, val = opt.split('=', 1)
|
||||
result[name.strip('-')] = val.rstrip('"\n')
|
||||
except IOError:
|
||||
logger.exception('Error when reading postmaster.opts')
|
||||
return result
|
||||
|
||||
def single_user_mode(self, command=None, options=None):
|
||||
"""run a given command in a single-user mode. If the command is empty - then just start and stop"""
|
||||
cmd = [self._pgcommand('postgres'), '--single', '-D', self._data_dir]
|
||||
for opt, val in sorted((options or {}).items()):
|
||||
cmd.extend(['-c', '{0}={1}'.format(opt, val)])
|
||||
# need a database name to connect
|
||||
cmd.append(self._database)
|
||||
p = subprocess.Popen(cmd, stdin=subprocess.PIPE, stdout=open(os.devnull, 'w'), stderr=subprocess.STDOUT)
|
||||
if p:
|
||||
if command:
|
||||
p.communicate('{0}\n'.format(command))
|
||||
p.stdin.close()
|
||||
return p.wait()
|
||||
return 1
|
||||
|
||||
def cleanup_archive_status(self):
|
||||
status_dir = os.path.join(self._data_dir, 'pg_' + self.wal_name, 'archive_status')
|
||||
try:
|
||||
for f in os.listdir(status_dir):
|
||||
path = os.path.join(status_dir, f)
|
||||
try:
|
||||
if os.path.islink(path):
|
||||
os.unlink(path)
|
||||
elif os.path.isfile(path):
|
||||
os.remove(path)
|
||||
except OSError:
|
||||
logger.exception('Unable to remove %s', path)
|
||||
except OSError:
|
||||
logger.exception('Unable to list %s', status_dir)
|
||||
|
||||
def fix_cluster_state(self):
|
||||
self.cleanup_archive_status()
|
||||
|
||||
# Start in a single user mode and stop to produce a clean shutdown
|
||||
opts = self.read_postmaster_opts()
|
||||
opts.update({'archive_mode': 'on', 'archive_command': 'false'})
|
||||
if os.path.isfile(self._recovery_conf) or os.path.islink(self._recovery_conf):
|
||||
os.unlink(self._recovery_conf)
|
||||
return self.single_user_mode(options=opts) == 0 or None
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
import contextlib
|
||||
import random
|
||||
import time
|
||||
import re
|
||||
@@ -281,8 +280,3 @@ def polling_loop(timeout, interval=1):
|
||||
yield iteration
|
||||
iteration += 1
|
||||
time.sleep(interval)
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def null_context():
|
||||
yield
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
__version__ = '1.3.1'
|
||||
__version__ = '1.3.5'
|
||||
|
||||
@@ -133,6 +133,7 @@ class Watchdog(object):
|
||||
|
||||
try:
|
||||
self.impl.open()
|
||||
actual_timeout = self._set_timeout()
|
||||
except WatchdogError as e:
|
||||
logger.warning("Could not activate %s: %s", self.impl.describe(), e)
|
||||
self.impl = NullWatchdog()
|
||||
@@ -141,8 +142,6 @@ class Watchdog(object):
|
||||
logger.warning("Watchdog implementation can't be disabled."
|
||||
" Watchdog will trigger after Patroni loses leader key.")
|
||||
|
||||
actual_timeout = self._set_timeout()
|
||||
|
||||
if not self.impl.is_running or actual_timeout > self.config.timeout:
|
||||
if self.config.mode == MODE_REQUIRED:
|
||||
if self.impl.is_null:
|
||||
@@ -202,7 +201,8 @@ class Watchdog(object):
|
||||
@synchronized
|
||||
def keepalive(self):
|
||||
try:
|
||||
self.impl.keepalive()
|
||||
if self.active:
|
||||
self.impl.keepalive()
|
||||
# In case there are any pending configuration changes apply them now.
|
||||
if self.active and self.config != self.active_config:
|
||||
if self.config.mode != MODE_OFF and self.active_config.mode == MODE_OFF:
|
||||
|
||||
+20
-10
@@ -155,21 +155,25 @@ class LinuxWatchdogDevice(WatchdogBase):
|
||||
def can_be_disabled(self):
|
||||
return self.get_support().has_MAGICCLOSE
|
||||
|
||||
def _ioctl(self, func, arg, mutate_arg=False):
|
||||
def _ioctl(self, func, arg):
|
||||
"""Runs the specified ioctl on the underlying fd.
|
||||
|
||||
Raises WatchdogError if the device is closed.
|
||||
Raises OSError or IOError (Python 2) when the ioctl fails."""
|
||||
if self._fd is None:
|
||||
raise WatchdogError("Watchdog device is closed")
|
||||
|
||||
result = fcntl.ioctl(self._fd, func, arg, mutate_arg)
|
||||
if result < 0:
|
||||
raise IOError(result)
|
||||
fcntl.ioctl(self._fd, func, arg, True)
|
||||
|
||||
def get_support(self):
|
||||
if self._support_cache is None:
|
||||
info = watchdog_info()
|
||||
self._ioctl(WDIOC_GETSUPPORT, info, True)
|
||||
try:
|
||||
self._ioctl(WDIOC_GETSUPPORT, info)
|
||||
except (WatchdogError, OSError, IOError) as e:
|
||||
raise WatchdogError("Could not get information about watchdog device: {}".format(e))
|
||||
self._support_cache = WatchdogInfo(info.options,
|
||||
info.firmware_version,
|
||||
str(bytearray(info.identity)).rstrip('\x00'))
|
||||
bytearray(info.identity).decode(errors='ignore').rstrip('\x00'))
|
||||
return self._support_cache
|
||||
|
||||
def describe(self):
|
||||
@@ -180,7 +184,7 @@ class LinuxWatchdogDevice(WatchdogBase):
|
||||
try:
|
||||
_, version, identity = self.get_support()
|
||||
ver_str = " (firmware {0})".format(version) if version else ""
|
||||
except WatchdogError: # XXX: Can it really be raise when self._fd is not None?
|
||||
except WatchdogError:
|
||||
pass
|
||||
|
||||
return identity + ver_str + dev_str
|
||||
@@ -199,11 +203,17 @@ class LinuxWatchdogDevice(WatchdogBase):
|
||||
timeout = int(timeout)
|
||||
if not 0 < timeout < 0xFFFF:
|
||||
raise WatchdogError("Invalid timeout {0}. Supported values are between 1 and 65535".format(timeout))
|
||||
self._ioctl(WDIOC_SETTIMEOUT, ctypes.c_int(timeout))
|
||||
try:
|
||||
self._ioctl(WDIOC_SETTIMEOUT, ctypes.c_int(timeout))
|
||||
except (WatchdogError, OSError, IOError) as e:
|
||||
raise WatchdogError("Could not set timeout on watchdog device: {}".format(e))
|
||||
|
||||
def get_timeout(self):
|
||||
timeout = ctypes.c_int()
|
||||
self._ioctl(WDIOC_GETTIMEOUT, timeout, True)
|
||||
try:
|
||||
self._ioctl(WDIOC_GETTIMEOUT, timeout)
|
||||
except (WatchdogError, OSError, IOError) as e:
|
||||
raise WatchdogError("Could not get timeout on watchdog device: {}".format(e))
|
||||
return timeout.value
|
||||
|
||||
|
||||
|
||||
+1
-1
@@ -6,7 +6,7 @@ requests
|
||||
six >= 1.7
|
||||
kazoo==2.2.1
|
||||
python-etcd>=0.4.3,<0.5
|
||||
python-consul==0.7.0
|
||||
python-consul>=0.7.0
|
||||
click>=4.1
|
||||
prettytable>=0.7
|
||||
tzlocal
|
||||
|
||||
+7
-3
@@ -3,7 +3,7 @@ import json
|
||||
import psycopg2
|
||||
import unittest
|
||||
|
||||
from mock import Mock, patch
|
||||
from mock import Mock, PropertyMock, patch
|
||||
from patroni.api import RestApiHandler, RestApiServer
|
||||
from patroni.dcs import ClusterConfig, Member
|
||||
from patroni.ha import _MemberStatus
|
||||
@@ -152,6 +152,7 @@ class TestRestApiHandler(unittest.TestCase):
|
||||
def test_do_OPTIONS(self):
|
||||
self.assertIsNotNone(MockRestApiServer(RestApiHandler, 'OPTIONS / HTTP/1.0'))
|
||||
|
||||
@patch.object(MockPostgresql, 'state', PropertyMock(return_value='stopped'))
|
||||
def test_do_GET_patroni(self):
|
||||
self.assertIsNotNone(MockRestApiServer(RestApiHandler, 'GET /patroni'))
|
||||
|
||||
@@ -280,6 +281,7 @@ class TestRestApiHandler(unittest.TestCase):
|
||||
def test_do_POST_failover(self, dcs):
|
||||
dcs.loop_wait = 10
|
||||
cluster = dcs.get_cluster.return_value
|
||||
cluster.is_synchronous_mode.return_value = False
|
||||
|
||||
post = 'POST /failover HTTP/1.0' + self._authorization + '\nContent-Length: '
|
||||
|
||||
@@ -291,14 +293,16 @@ class TestRestApiHandler(unittest.TestCase):
|
||||
cluster.leader.name = 'postgresql1'
|
||||
MockRestApiServer(RestApiHandler, request)
|
||||
|
||||
MockRestApiServer(RestApiHandler, post + '25\n\n{"leader": "postgresql1"}')
|
||||
for cluster.is_synchronous_mode.return_value in (True, False):
|
||||
MockRestApiServer(RestApiHandler, post + '25\n\n{"leader": "postgresql1"}')
|
||||
|
||||
cluster.leader.name = 'postgresql2'
|
||||
request = post + '53\n\n{"leader": "postgresql1", "candidate": "postgresql2"}'
|
||||
MockRestApiServer(RestApiHandler, request)
|
||||
|
||||
cluster.leader.name = 'postgresql1'
|
||||
MockRestApiServer(RestApiHandler, request)
|
||||
for cluster.is_synchronous_mode.return_value in (True, False):
|
||||
MockRestApiServer(RestApiHandler, request)
|
||||
|
||||
cluster.members = [Member(0, 'postgresql0', 30, {'api_url': 'http'}),
|
||||
Member(0, 'postgresql2', 30, {'api_url': 'http'})]
|
||||
|
||||
@@ -45,7 +45,7 @@ class TestHTTPClient(unittest.TestCase):
|
||||
|
||||
def test_get(self):
|
||||
self.client.get(Mock(), '')
|
||||
self.client.get(Mock(), '', {'wait': '1s', 'index': 1})
|
||||
self.client.get(Mock(), '', {'wait': '1s', 'index': 1, 'token': 'foo'})
|
||||
self.client.http.request.return_value.status = 500
|
||||
self.assertRaises(ConsulInternalError, self.client.get, Mock(), '')
|
||||
|
||||
@@ -69,6 +69,10 @@ class TestConsul(unittest.TestCase):
|
||||
@patch.object(consul.Consul.KV, 'get', kv_get)
|
||||
@patch.object(consul.Consul.KV, 'delete', Mock())
|
||||
def setUp(self):
|
||||
Consul({'ttl': 30, 'scope': 't', 'name': 'p', 'url': 'https://l:1', 'retry_timeout': 10,
|
||||
'verify': 'on', 'key': 'foo', 'cert': 'bar', 'cacert': 'buz'})
|
||||
Consul({'ttl': 30, 'scope': 't', 'name': 'p', 'url': 'https://l:1', 'retry_timeout': 10,
|
||||
'verify': 'on', 'cert': 'bar', 'cacert': 'buz'})
|
||||
self.c = Consul({'ttl': 30, 'scope': 'test', 'name': 'postgresql1', 'host': 'localhost:1', 'retry_timeout': 10})
|
||||
self.c._base_path = '/service/good'
|
||||
self.c._load_cluster()
|
||||
|
||||
+27
-3
@@ -63,6 +63,7 @@ def get_node_status(reachable=True, in_recovery=True, wal_position=10, nofailove
|
||||
return _MemberStatus(e, reachable, in_recovery, wal_position, tags, watchdog_failed)
|
||||
return fetch_node_status
|
||||
|
||||
|
||||
future_restart_time = datetime.datetime.now(tzutc) + datetime.timedelta(days=5)
|
||||
postmaster_start_time = datetime.datetime.now(tzutc)
|
||||
|
||||
@@ -157,7 +158,6 @@ class TestHa(unittest.TestCase):
|
||||
self.ha.old_cluster = self.e.get_cluster()
|
||||
self.ha.cluster = get_cluster_not_initialized_without_leader()
|
||||
self.ha.load_cluster_from_dcs = Mock()
|
||||
self.ha.is_synchronous_mode = false
|
||||
|
||||
def test_update_lock(self):
|
||||
self.p.last_operation = Mock(side_effect=PostgresConnectionException(''))
|
||||
@@ -193,6 +193,15 @@ class TestHa(unittest.TestCase):
|
||||
self.ha.cluster = get_cluster_initialized_with_leader()
|
||||
self.assertEquals(self.ha.run_cycle(), 'running pg_rewind from leader')
|
||||
|
||||
@patch.object(Postgresql, 'can_rewind', PropertyMock(return_value=True))
|
||||
@patch.object(Postgresql, 'fix_cluster_state', Mock())
|
||||
def test_single_user_after_recover_failed(self):
|
||||
self.p.controldata = lambda: {'Database cluster state': 'in production'}
|
||||
self.p.is_running = false
|
||||
self.p.follow = false
|
||||
self.assertEquals(self.ha.run_cycle(), 'starting as a secondary')
|
||||
self.assertEquals(self.ha.run_cycle(), 'fixing cluster state in a single user mode')
|
||||
|
||||
@patch('sys.exit', return_value=1)
|
||||
@patch('patroni.ha.Ha.sysid_valid', MagicMock(return_value=True))
|
||||
def test_sysid_no_match(self, exit_mock):
|
||||
@@ -246,7 +255,7 @@ class TestHa(unittest.TestCase):
|
||||
with patch.object(Watchdog, 'activate', Mock(return_value=False)):
|
||||
self.assertEquals(self.ha.run_cycle(), 'Demoting self because watchdog could not be activated')
|
||||
self.p.is_leader = false
|
||||
self.assertEquals(self.ha.run_cycle(), 'Not promoting self because watchdog could not be actived')
|
||||
self.assertEquals(self.ha.run_cycle(), 'Not promoting self because watchdog could not be activated')
|
||||
|
||||
def test_leader_with_lock(self):
|
||||
self.ha.cluster.is_unlocked = false
|
||||
@@ -437,6 +446,19 @@ class TestHa(unittest.TestCase):
|
||||
self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, self.p.name, '', None))
|
||||
self.assertEquals('PAUSE: no action. i am the leader with the lock', self.ha.run_cycle())
|
||||
|
||||
@patch('requests.get', requests_get)
|
||||
def test_manual_failover_from_leader_in_synchronous_mode(self):
|
||||
self.p.is_leader = true
|
||||
self.ha.has_lock = true
|
||||
self.ha.is_synchronous_mode = true
|
||||
self.ha.is_failover_possible = false
|
||||
self.ha.process_sync_replication = Mock()
|
||||
self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, self.p.name, 'a', None), (self.p.name, None))
|
||||
self.assertEquals('no action. i am the leader with the lock', self.ha.run_cycle())
|
||||
self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, self.p.name, 'a', None), (self.p.name, 'a'))
|
||||
self.ha.is_failover_possible = true
|
||||
self.assertEquals('manual failover: demoting myself', self.ha.run_cycle())
|
||||
|
||||
@patch('requests.get', requests_get)
|
||||
def test_manual_failover_process_no_leader(self):
|
||||
self.p.is_leader = false
|
||||
@@ -634,7 +656,8 @@ class TestHa(unittest.TestCase):
|
||||
@patch('patroni.ha.Ha.demote')
|
||||
def test_failover_immediately_on_zero_master_start_timeout(self, demote):
|
||||
self.p.is_running = false
|
||||
self.ha.cluster = get_cluster_initialized_with_leader()
|
||||
self.ha.cluster = get_cluster_initialized_with_leader(sync=(self.p.name, 'other'))
|
||||
self.ha.cluster.config.data['synchronous_mode'] = True
|
||||
self.ha.patroni.config.set_dynamic_configuration({'master_start_timeout': 0})
|
||||
self.ha.has_lock = true
|
||||
self.ha.update_lock = true
|
||||
@@ -839,6 +862,7 @@ class TestHa(unittest.TestCase):
|
||||
self.ha.has_lock = true
|
||||
self.p.data_directory_empty = true
|
||||
self.assertEquals(self.ha.run_cycle(), 'released leader key voluntarily as data dir empty and currently leader')
|
||||
self.assertEquals(self.p.role, 'uninitialized')
|
||||
|
||||
# as has_lock is mocked out, we need to fake the leader key release
|
||||
self.ha.has_lock = false
|
||||
|
||||
@@ -174,7 +174,7 @@ class TestPostgresql(unittest.TestCase):
|
||||
if not os.path.exists(self.data_dir):
|
||||
os.makedirs(self.data_dir)
|
||||
self.p = Postgresql({'name': 'test0', 'scope': 'batman', 'data_dir': self.data_dir,
|
||||
'config_dir': self.config_dir, 'retry_timeout': 10,
|
||||
'config_dir': self.config_dir, 'retry_timeout': 10, 'pgpass': '/tmp/pgpass0',
|
||||
'listen': '127.0.0.2, 127.0.0.3:5432', 'connect_address': '127.0.0.2:5432',
|
||||
'authentication': {'superuser': {'username': 'test', 'password': 'test'},
|
||||
'replication': {'username': 'replicator', 'password': 'rep-pass'}},
|
||||
@@ -327,10 +327,10 @@ class TestPostgresql(unittest.TestCase):
|
||||
@patch.object(Postgresql, 'can_rewind', PropertyMock(return_value=True))
|
||||
def test__get_local_timeline_lsn(self):
|
||||
self.p.trigger_check_diverged_lsn()
|
||||
with patch.object(Postgresql, 'controldata', Mock(return_value={'Database cluster state': 'shut down'})):
|
||||
self.p.rewind_needed_and_possible(self.leader)
|
||||
with patch.object(Postgresql, 'controldata',
|
||||
Mock(return_value={'Database cluster state': 'shut down in recovery'})):
|
||||
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)
|
||||
with patch.object(Postgresql, 'is_running', Mock(return_value=True)):
|
||||
with patch.object(MockCursor, 'fetchone', Mock(side_effect=[(False, ), Exception])):
|
||||
@@ -883,3 +883,41 @@ class TestPostgresql(unittest.TestCase):
|
||||
def test_terminate_starting_postmaster(self):
|
||||
self.p.terminate_starting_postmaster(123)
|
||||
self.p.terminate_starting_postmaster(123)
|
||||
|
||||
def test_read_postmaster_opts(self):
|
||||
m = mock_open(read_data='/usr/lib/postgres/9.6/bin/postgres "-D" "data/postgresql0" \
|
||||
"--listen_addresses=127.0.0.1" "--port=5432" "--hot_standby=on" "--wal_level=hot_standby" \
|
||||
"--wal_log_hints=on" "--max_wal_senders=5" "--max_replication_slots=5"\n')
|
||||
with patch.object(builtins, 'open', m):
|
||||
data = self.p.read_postmaster_opts()
|
||||
self.assertEquals(data['wal_level'], 'hot_standby')
|
||||
self.assertEquals(int(data['max_replication_slots']), 5)
|
||||
self.assertEqual(data.get('D'), None)
|
||||
|
||||
m.side_effect = IOError
|
||||
data = self.p.read_postmaster_opts()
|
||||
self.assertEqual(data, dict())
|
||||
|
||||
@patch('subprocess.Popen')
|
||||
@patch.object(builtins, 'open', Mock(return_value=42))
|
||||
def test_single_user_mode(self, subprocess_popen_mock):
|
||||
subprocess_popen_mock.return_value.wait.return_value = 0
|
||||
self.assertEquals(self.p.single_user_mode(command="CHECKPOINT"), 0)
|
||||
subprocess_popen_mock.return_value = None
|
||||
self.assertEquals(self.p.single_user_mode(), 1)
|
||||
self.assertEquals(self.p.single_user_mode(options={'archive_mode': 'on'}), 1)
|
||||
|
||||
@patch('os.listdir', Mock(side_effect=[OSError, ['a', 'b']]))
|
||||
@patch('os.unlink', Mock(side_effect=OSError))
|
||||
@patch('os.remove', Mock())
|
||||
@patch('os.path.islink', Mock(side_effect=[True, False]))
|
||||
@patch('os.path.isfile', Mock(return_value=True))
|
||||
def test_cleanup_archive_status(self):
|
||||
self.p.cleanup_archive_status()
|
||||
self.p.cleanup_archive_status()
|
||||
|
||||
@patch('os.unlink', Mock())
|
||||
@patch('os.path.isfile', Mock(return_value=True))
|
||||
@patch.object(Postgresql, 'single_user_mode', Mock(return_value=0))
|
||||
def test_fix_cluster_state(self):
|
||||
self.assertTrue(self.p.fix_cluster_state())
|
||||
|
||||
+13
-3
@@ -132,8 +132,9 @@ class TestWatchdog(unittest.TestCase):
|
||||
def test_exceptions(self):
|
||||
wd = Watchdog({'ttl': 30, 'loop_wait': 10, 'watchdog': {'mode': 'bad'}})
|
||||
wd.impl.close = wd.impl.keepalive = Mock(side_effect=WatchdogError(''))
|
||||
self.assertIsNone(wd.disable())
|
||||
self.assertTrue(wd.activate())
|
||||
self.assertIsNone(wd.keepalive())
|
||||
self.assertIsNone(wd.disable())
|
||||
|
||||
@patch('platform.system', Mock(return_value='Linux'))
|
||||
def test_config_reload(self):
|
||||
@@ -193,15 +194,24 @@ class TestLinuxWatchdogDevice(unittest.TestCase):
|
||||
self.assertRaises(WatchdogError, self.impl.set_timeout, -1)
|
||||
|
||||
@patch('os.open', Mock(return_value=3))
|
||||
@patch('fcntl.ioctl', Mock(return_value=-1))
|
||||
@patch('fcntl.ioctl', Mock(side_effect=OSError))
|
||||
def test__ioctl(self):
|
||||
self.assertRaises(WatchdogError, self.impl.get_support)
|
||||
self.impl.open()
|
||||
self.assertRaises(IOError, self.impl.get_support)
|
||||
self.assertRaises(WatchdogError, self.impl.get_support)
|
||||
|
||||
def test_is_healthy(self):
|
||||
self.assertFalse(self.impl.is_healthy)
|
||||
|
||||
@patch('os.open', Mock(return_value=3))
|
||||
@patch('fcntl.ioctl', Mock(side_effect=OSError))
|
||||
def test_error_handling(self):
|
||||
self.impl.open()
|
||||
self.assertRaises(WatchdogError, self.impl.get_timeout)
|
||||
self.assertRaises(WatchdogError, self.impl.set_timeout, 10)
|
||||
# We still try to output a reasonable string even if getting info errors
|
||||
self.assertEquals(self.impl.describe(), "Linux watchdog device")
|
||||
|
||||
@patch('os.open', Mock(side_effect=OSError))
|
||||
def test_open(self):
|
||||
self.assertRaises(WatchdogError, self.impl.open)
|
||||
|
||||
Reference in New Issue
Block a user