Compare commits

...
36 Commits
Author SHA1 Message Date
Alexander KukushkinandGitHub b341ab2e2f Release 2.0.2 (#1851)
* bump version
* update release notes
* implement missing unit-test
2021-02-22 12:28:19 +01:00
Alex BrasetvikandGitHub 82b918e10d Constant time comparison of auth key (#1847)
… to avoid [timing attacks](https://codahale.com/a-lesson-in-timing-attacks/).
2021-02-19 11:16:41 +01:00
Mark MercadoandGitHub 77f08af682 Create raft data_dir if necessary (#1846)
Close #1760
2021-02-18 11:35:31 +01:00
Mark MercadoandGitHub 32daef3939 Fixing grammar and typos (#1845) 2021-02-18 10:28:38 +01:00
Alexander KukushkinandGitHub b698df374f Fix build (#1843)
run apt-get update before installing packages
2021-02-16 09:34:36 +01:00
Alexander KukushkinandGitHub 9f252d246e Improve handling of concurrent update error (#1796)
The old strategy was waiting for 1 second and hoping that we will get an update event from the WATCH connection.
Unfortunately, it didn't work well in practice. Instead, we will get the current value from the API by performing an explicit read request.

Close https://github.com/zalando/patroni/issues/1767
2021-02-11 15:55:05 +01:00
Alexander KukushkinandGitHub 39332c93ed Treat PATRONI_KUBERNETES_USE_ENDPOINTS env as boolean (#1832)
Close https://github.com/zalando/patroni/issues/1814
2021-02-03 09:44:59 +01:00
krishnaandGitHub b3dc765e6d Choose synchronous nodes based on replication lag (#1786)
This commit makes it possible to configure the maximum lag (`maximum_lag_on_syncnode`) after which Patroni will "demote" the node from synchronous and replace it with another node.

The previous implementation always tried to stick to the same synchronous nodes (even if they are not optimal ones).
2021-02-02 15:45:02 +01:00
Kaarel MoppelandGitHub 9d7d4423e3 Docs: document the need for special configuration for symlinked pg_wal (#1818)
If an existing instance was configured with WAL residing outside of
PGDATA then currently a 'reinit' would lose such symlinks. So add some
bits of information on that to draw attention to this cornercase issue
and also add the --waldir option to the sample `postgresql.basebackup`
configuration sections to increase visibility.

Discussion: https://github.com/zalando/patroni/issues/1817
2021-02-02 11:50:35 +01:00
Doug WhitfieldandGitHub 9c5c0e71c2 Add a new section to end on testing HA solution (#1827)
Thanks to @ants for the suggestion and some tips on testing via slack.
2021-02-02 11:49:50 +01:00
Alexander KukushkinandGitHub cdfc4ea50f Handle case with psutil cmdline() returning empty list (#1829)
Close https://github.com/zalando/patroni/issues/1828
2021-02-02 11:48:44 +01:00
Alexander KukushkinandGitHub 6bf205b190 Don't use bypass_api_service when running patronictl (#1830)
It could happen that the cluster role wither not configured or doesn't provide enough permissions. In this case bypass_api_service is ignored, but the warning is logged, which is rather annoying when patronictl is used.
Since the bypass_api_service is most useful for Patroni, we will simply ignore it when patronictl is used.
2021-02-02 11:47:46 +01:00
Jonathan S. KatzandGitHub accba93cbe Add support for encrypted TLS keys for REST API (#1825)
The Python SSL library allows for the inclusion of a password in its "load_cert_chain" function when setting up a SSLContext[1].
This allows for loading an encrypted key file in PEM representation to be loaded into the certificate chain.

This commit adds the optional "keyfile_password" parameter to the REST API block of configuration so that Patroni can load in encrypted private keys when establishing its TLS socket.

This also adds the corollary "PATRONI_RESTAPI_KEYFILE_PASSWORD" environmental variable, which has the same effect.

[1] https://docs.python.org/3/library/ssl.html#ssl.SSLContext.load_cert_chain
2021-02-02 11:47:09 +01:00
Gunnar "Nick" BluthandGitHub ba4ab58d40 Support cipher suite limitation for REST API (#1824)
Many environments require a limitation of allowed TLS cipher suites / levels.
See e.g. the german BSI requirements: 
https://www.bsi.bund.de/SharedDocs/Downloads/EN/BSI/Publications/TechGuidelines/TG02102/BSI-TR-02102-2.pdf?__blob=publicationFile&v=10

This implements an optional "ciphers" setting that - if given - enforces the ciphers on the REST API socket.

See also #1730.
2021-01-27 13:53:28 +01:00
Alexander KukushkinandGitHub 8b5cb85536 Exit only if authentication explicitly failed (#1806)
It could happen that one of etcd servers is not accessible on Patroni start.
In this case Patroni was trying to perform authentication and exiting, while it should exit only if Etcd explicitly responded with the `AuthFailed` error.

Close https://github.com/zalando/patroni/issues/1805
2021-01-15 14:31:33 +01:00
Alexander KukushkinandGitHub a9f86aa195 Add compatibility with python-consul2 (#1812)
the good old python-consul is not maintained for a few years in a row, therefore someone forked under a different name, but package files are installed into the same location as for the old.

The API of both modules is mostly compatible therefore it wasn't hard to add the support of both modules in Patroni.

Taking into account that python-consul is not a direct requirement for Patroni, but extra, now the end-user has a choice what to install.

Close https://github.com/zalando/patroni/issues/1810
2021-01-15 14:30:48 +01:00
Alexander KukushkinandGitHub 4a8c4cfc53 Make tests more reliable (#1808)
1.  Fix flaky behave tests with zookeeper. First, install/start binaries (zookeeper/localkube) and only after that continue with installing requirements and running behave. Previously zookeeper didn't had enough time to start and tests sometimes were failing.
2.  Fix flaky raft tests. Despite observations of MacOS slowness, for some unknown reason the delete test with a very small timeout was not timing out, but succeeding, causing unit-tests to fail. The solution - do not rely on the actual timeout, but mock it.
2021-01-15 14:29:55 +01:00
Alexander KukushkinandGitHub 8446077fb3 Fixes around pg_rewind (#1794)
1. If the superuser name is different from postgres, the pg_rewind in the standby cluster was failing because the connection string didn't contain the database name.
2. Provide output if the single-user mode recovery failed.

Close https://github.com/zalando/patroni/pull/1736
2020-12-16 19:54:19 +01:00
Alexander KukushkinandGitHub 94b9f8fae6 Silence unhandled exceptions in Thread.run() during unit-tests (#1802)
Python 3.8 changed the way how exceptions raised from the Thread.run() method are handled.
It resulted in unit-tests showing a couple of warnings. They are not important and we just silence them.
2020-12-16 19:37:51 +01:00
Alexander KukushkinandGitHub df9cadcdc4 Handle AttributeError in etcd.py (#1801)
When the__del__() method is executed the python interpreter already unloaded some of the modules that are still used down the http.clear() method.
The only we could do in this case - silence some exceptions like ReferenceError, TypeError, AttributeError.

Close https://github.com/zalando/patroni/issues/1785
2020-12-16 19:23:56 +01:00
Alexander KukushkinandGitHub 9b263dc6c9 Move find_executable to utils (#1799)
We don't want to import the whole patroni.ctl into the patroni
2020-12-16 18:58:54 +01:00
Alexander KukushkinandGitHub 3a87d0e99b Implement missing validators for etcd3 and raft (#1798)
Close https://github.com/zalando/patroni/issues/1771
2020-12-16 18:44:58 +01:00
Alexander KukushkinandGitHub 89a15a2df4 Fix small issues with ignore-slots feature (#1797)
When there is no config key in DCS Patroni shouldn't try accessing ignore_slots, otherwise an exception is raised.

In addition to that implement missing unit-tests and fix linting issues in behave tests.
2020-12-16 18:10:12 +01:00
Alexander KukushkinandGitHub 61bf412ac6 Fix bug in the post_bootstrap method() (#1795)
The config.write_pgpass(r) should be called unconditionally, no matter whether the password defined in the config or not.

Close https://github.com/zalando/patroni/issues/1727
2020-12-16 17:33:44 +01:00
Nicolas ThauvinandGitHub 492f755bd5 Warn the user when the required watchdog is not healthy (#1779)
When the watchdog device is not writable or missing in required mode, the member cannot be promoted. It only keeps logging that it is not the healthiest member.
Add a warning to show the user where to search for this misconfiguration.
2020-12-16 16:46:14 +01:00
Igor YanchenkoandGitHub 583fccfb34 Start postgres with hot_standby=off when doing PITR (#1788)
When doing a custom bootstrap the configured value of max_connections could be smaller than the actual value stored in the pg_contol. As a workaround Patroni already did an adjustment of such parameters. Unfortunately, it doesn't cover the case when the parameter value was increased even higher during the WAL replay. That was causing postgres to stop.
The hot_standby=off disable connections to the instance, but allows postgres to run with parameter values smaller than required.
The hot_standby=off doesn't affect the running primary. When the recovery finishes and postgres promotes it starts accepting connections as usual. If the node will be demoted, Patroni will stop postgres and start it up as replica with hot_standby=on.
2020-12-14 15:59:46 +01:00
Pascal GOUHIERandGitHub 53b525202c Add existing_data to toc (#1790) 2020-12-14 15:53:02 +01:00
Kaarel MoppelandGitHub 464019eaf7 Mention currently supported PostgreSQL versions (#1777) 2020-12-14 15:51:06 +01:00
andrewlecuyerandGitHub 37bff48915 Fix Invalid os.symlink Calls when Moving Data Dir (#1781)
Fixes #1780

Specifically fixes the calls to `os.symlink` within the `move_data_directory` function, as used to update the symlink(s) for any WAL or tablespace directories following an init failure.  The new name for a WAL and/or tablespace directory is now passed in as in as the first argument when calling `os.symlink`, while the second argument is now the name of the symlink that is being recreated.

This aligns with the Python documentation for `os.symlink`, which states the following regarding the first two arguments for the `os.symlink` function (arguments `src` and `dst`): 

> Create a symbolic link pointing to src named dst (https://docs.python.org/3/library/os.html#os.symlink)

With this change, all WAL and/or tablespaces directories will be properly renamed following a init failure, and the `FileExistsError` that would otherwise occur when attempting to recreate the associated symlinks is no longer thrown, e.g.:

```bash
2020-12-05 20:35:42.219 GMT [262] LOG:  unrecognized configuration parameter "invalid" in file "/pgdata/mycluster1/postgresql.auto.conf" line 5
2020-12-05 20:35:42.219 GMT [262] FATAL:  configuration file "/pgdata/mycluster1/postgresql.auto.conf" contains errors
2020-12-05 20:35:43,220 ERROR: postmaster is not running
2020-12-05 20:35:43,223 INFO: removing initialize key after failed attempt to bootstrap the cluster
2020-12-05 20:35:43,292 INFO: renaming user defined tablespace directory and updating symlink: /tablespaces/ts1/ts1
2020-12-05 20:35:43,314 INFO: renaming data directory to /pgdata/mycluster1_2020-12-05-20-35-43
Traceback (most recent call last):
```

Additionally, any/all WAL and/or tablespace symlinks are properly recreated, e.g.:

```bash
$ ls -l pg_tblspc
total 0
lrwxrwxrwx. 1 postgres postgres 40 Dec  5 20:35 16414 -> /tablespaces/ts1/ts1_2020-12-05-20-35-43
```
2020-12-14 15:50:14 +01:00
Alexander KukushkinandGitHub a43baece68 Submit coverage to codacy only if secret is available (#1791)
If PR is open from the external GH repo secrets are not set due to security reasons. It makes codacy coverage report to fail.
2020-12-14 15:49:37 +01:00
Alexander KukushkinandGitHub e3ef9ac306 Fix issues with zookeeper (#1792)
1. The `ttl` was incorrectly returned 1000 times higher then it should
2. The `watch()` method must return True if the parent method returned True. Not doing so resulted in the incorrect calculation of sleep time.
3. Move mock of exhibitor api to the features/environment.py. It simplifies testing with behave.
2020-12-14 15:12:57 +01:00
Alexander KukushkinandGitHub 1530ed0b9c Switch to GH actions (#1778)
it allows up to 20 parallel builds
2020-12-04 21:52:34 +01:00
Nicolas LimageandGitHub e2b2daf0b3 add missing shutdown_request (#1770)
This patch fixes the error handling of cases where there are runtime errors in `socketserver`.
For example, when creating a new thread (to handle a request) fails.

`get_request` handles ssl connections by replacing the new client socket by a tuple containing `(server_socket, new_client_socket)` in order to later deal with handshakes in `process_request_thread`

During the processing of a request, the socketserver `BaseServer` calls `handle_request`, calling the `_handle_request_noblock`, which is calling the following functions (https://github.com/python/cpython/blob/3.8/Lib/socketserver.py#L303):

```
request, client_addr = get_request()
verify_request(request, client_address):
process_request(request, client_address)
handle_error(request, client_address)
shutdown_request(request)
```

- `get_request` is overloaded in patroni and returns `request` as a tuple in case of ssl calls
- `verify_request` defaults to `return True` and should be fixed if used but is fine in this case
- `process_request` just calls `process_request_thread` (which is overloaded in patroni and handles tuple-style requests)
- `handle_error` is overloaded in patroni and handles tuple-style requests)
- but `shutdown_request` is not overloaded and thus missing support for tuple-style requests

This patch adds support for tuple-style requests in patroni api
2020-11-27 19:25:06 +01:00
James ColemanandGitHub d7f579ee61 Feature: ability to ignore externally managed replication slots (#1742)
There are sometimes good reasons to manage replication slots externally
to Patroni. For example, a consumer may wish to manage its own slots (so
that it can more easily track when a failover has a occurred and whether
it is ahead of or behind the WAL position on the new primary).
Additionally tooling like pglogical actually replicates slots to all
replicas so that the current position can be maintained on failover
targets (this also aids consumers by supplying primitives so that they
can verify data hasn't been lost or a split brain occurred relative to
the physical cluster).

To support these use cases this new feature allows configuring Patroni
to entirely ignore sets of slots specified by any subset of name,
database, slot type, and plugin.
2020-11-24 11:45:14 +01:00
James ColemanandGitHub acf512712e Make it easier to develop locally (#1748)
Previously the only documentation for how to run tests was the
implementation in the Travis configuration file. Here we add
instructions as well as move development dependencies to an easily used
and shared (with Travis config) separate requirements.dev.txt file.
2020-11-24 08:29:36 +01:00
Alexander KukushkinandGitHub e8e87bf0a1 Don't interrupt restart or promote if lost leader lock in pause (#1726)
In pause it is allowed to run postgres as master without lock.
2020-10-08 08:56:53 +02:00
55 changed files with 1101 additions and 328 deletions
+182
View File
@@ -0,0 +1,182 @@
import inspect
import os
import shutil
import subprocess
import stat
import sys
import tarfile
import time
import zipfile
def install_requirements(what):
old_path = sys.path[:]
w = os.path.join(os.getcwd(), os.path.dirname(inspect.getfile(inspect.currentframe())))
sys.path.insert(0, os.path.dirname(os.path.dirname(w)))
try:
from setup import EXTRAS_REQUIRE, read
finally:
sys.path = old_path
requirements = ['mock>=2.0.0', 'flake8', 'pytest', 'pytest-cov'] if what == 'all' else ['behave']
requirements += ['psycopg2-binary', 'coverage']
for r in read('requirements.txt').split('\n'):
r = r.strip()
if r != '':
extras = {e for e, v in EXTRAS_REQUIRE.items() if v and r.startswith(v[0])}
if not extras or what == 'all' or what in extras:
requirements.append(r)
subprocess.call([sys.executable, '-m', 'pip', 'install', '--upgrade', 'pip'])
r = subprocess.call([sys.executable, '-m', 'pip', 'install'] + requirements)
s = subprocess.call([sys.executable, '-m', 'pip', 'install', '--upgrade', 'setuptools'])
return s | r
def install_packages(what):
packages = {
'zookeeper': ['zookeeper', 'zookeeper-bin', 'zookeeperd'],
'consul': ['consul'],
}
packages['exhibitor'] = packages['zookeeper']
packages = packages.get(what, [])
ver = str({'etcd': '9.6', 'etcd3': '9.6', 'consul': 10, 'exhibitor': 11, 'kubernetes': 12, 'raft': 13}.get(what))
subprocess.call(['sudo', 'apt-get', 'update', '-y'])
return subprocess.call(['sudo', 'apt-get', 'install', '-y', 'postgresql-' + ver, 'expect-dev', 'wget'] + packages)
def get_file(url, name):
try:
from urllib.request import urlretrieve
except ImportError:
from urllib import urlretrieve
print('Downloading ' + url)
urlretrieve(url, name)
def untar(archive, name):
with tarfile.open(archive) as tar:
f = tar.extractfile(name)
dest = os.path.basename(name)
with open(dest, 'wb') as d:
shutil.copyfileobj(f, d)
return dest
def unzip(archive, name):
with zipfile.ZipFile(archive, 'r') as z:
name = z.extract(name)
dest = os.path.basename(name)
shutil.move(name, dest)
return dest
def unzip_all(archive):
print('Extracting ' + archive)
with zipfile.ZipFile(archive, 'r') as z:
z.extractall()
def chmod_755(name):
os.chmod(name, stat.S_IRUSR | stat.S_IWUSR | stat.S_IXUSR |
stat.S_IRGRP | stat.S_IXGRP | stat.S_IROTH | stat.S_IXOTH)
def unpack(archive, name):
print('Extracting {0} from {1}'.format(name, archive))
func = unzip if archive.endswith('.zip') else untar
name = func(archive, name)
chmod_755(name)
return name
def install_etcd():
version = os.environ.get('ETCDVERSION', '3.3.13')
platform = {'linux2': 'linux', 'win32': 'windows', 'cygwin': 'windows'}.get(sys.platform, sys.platform)
dirname = 'etcd-v{0}-{1}-amd64'.format(version, platform)
ext = 'tar.gz' if platform == 'linux' else 'zip'
name = '{0}.{1}'.format(dirname, ext)
url = 'https://github.com/etcd-io/etcd/releases/download/v{0}/{1}'.format(version, name)
get_file(url, name)
ext = '.exe' if platform == 'windows' else ''
return int(unpack(name, '{0}/etcd{1}'.format(dirname, ext)) is None)
def install_postgres():
version = os.environ.get('PGVERSION', '12.1-1')
platform = {'darwin': 'osx', 'win32': 'windows-x64', 'cygwin': 'windows-x64'}[sys.platform]
name = 'postgresql-{0}-{1}-binaries.zip'.format(version, platform)
get_file('http://get.enterprisedb.com/postgresql/' + name, name)
unzip_all(name)
bin_dir = os.path.join('pgsql', 'bin')
for f in os.listdir(bin_dir):
chmod_755(os.path.join(bin_dir, f))
subprocess.call(['pgsql/bin/postgres', '-V'])
return 0
def setup_kubernetes():
get_file('https://storage.googleapis.com/minikube/k8sReleases/v1.7.0/localkube-linux-amd64', 'localkube')
chmod_755('localkube')
devnull = open(os.devnull, 'w')
subprocess.Popen(['sudo', 'nohup', './localkube', '--logtostderr=true', '--enable-dns=false'],
stdout=devnull, stderr=devnull)
for _ in range(0, 120):
if subprocess.call(['wget', '-qO', '-', 'http://127.0.0.1:8080/'], stdout=devnull, stderr=devnull) == 0:
break
time.sleep(1)
else:
print('localkube did not start')
return 1
subprocess.call('sudo chmod 644 /var/lib/localkube/certs/*', shell=True)
print('Set up .kube/config')
kube = os.path.join(os.path.expanduser('~'), '.kube')
os.makedirs(kube)
with open(os.path.join(kube, 'config'), 'w') as f:
f.write("""apiVersion: v1
clusters:
- cluster:
certificate-authority: /var/lib/localkube/certs/ca.crt
server: https://127.0.0.1:8443
name: local
contexts:
- context:
cluster: local
user: myself
name: local
current-context: local
kind: Config
preferences: {}
users:
- name: myself
user:
client-certificate: /var/lib/localkube/certs/apiserver.crt
client-key: /var/lib/localkube/certs/apiserver.key
""")
return 0
def main():
what = os.environ.get('DCS', sys.argv[1] if len(sys.argv) > 1 else 'all')
if what != 'all':
if sys.platform.startswith('linux'):
r = install_packages(what)
if r == 0 and what == 'kubernetes':
r = setup_kubernetes()
else:
r = install_postgres()
if r == 0 and what.startswith('etcd'):
r = install_etcd()
if r != 0:
return r
return install_requirements(what)
if __name__ == '__main__':
sys.exit(main())
+47
View File
@@ -0,0 +1,47 @@
import os
import shutil
import subprocess
import sys
import tempfile
def main():
what = os.environ.get('DCS', sys.argv[1] if len(sys.argv) > 1 else 'all')
if what == 'all':
flake8 = subprocess.call([sys.executable, 'setup.py', 'flake8'])
test = subprocess.call([sys.executable, 'setup.py', 'test'])
version = '.'.join(map(str, sys.version_info[:2]))
shutil.move('.coverage', os.path.join(tempfile.gettempdir(), '.coverage.' + version))
return flake8 | test
elif what == 'combine':
tmp = tempfile.gettempdir()
for name in os.listdir(tmp):
if name.startswith('.coverage.'):
shutil.move(os.path.join(tmp, name), name)
return subprocess.call([sys.executable, '-m', 'coverage', 'combine'])
env = os.environ.copy()
if sys.platform.startswith('linux'):
version = {'etcd': '9.6', 'etcd3': '9.6', 'consul': 10, 'exhibitor': 11, 'kubernetes': 12, 'raft': 13}.get(what)
path = '/usr/lib/postgresql/{0}/bin:.'.format(version)
unbuffer = ['timeout', '600', 'unbuffer']
else:
path = os.path.abspath(os.path.join('pgsql', 'bin'))
if sys.platform == 'darwin':
path += ':.'
unbuffer = []
env['PATH'] = path + os.pathsep + env['PATH']
env['DCS'] = what
ret = subprocess.call(unbuffer + [sys.executable, '-m', 'behave'], env=env)
if ret != 0:
if subprocess.call('grep . features/output/*_failed/*postgres?.*', shell=True) != 0:
subprocess.call('grep . features/output/*/*postgres?.*', shell=True)
return 1
return 0
if __name__ == '__main__':
sys.exit(main())
+175
View File
@@ -0,0 +1,175 @@
name: Tests
on:
pull_request:
push:
branches:
- master
tags:
- v.*
jobs:
unit:
runs-on: ${{ matrix.os }}-latest
strategy:
fail-fast: false
matrix:
os: [ubuntu, windows, macos]
steps:
- uses: actions/checkout@v1
- name: Set up Python 2.7
uses: actions/setup-python@v2
with:
python-version: 2.7
if: matrix.os != 'windows'
- name: Install dependencies
run: python .github/workflows/install_deps.py
if: matrix.os != 'windows'
- name: Run tests and flake8
run: python .github/workflows/run_tests.py
if: matrix.os != 'windows'
- name: Set up Python 3.5
uses: actions/setup-python@v2
with:
python-version: 3.5
- name: Install dependencies
run: python .github/workflows/install_deps.py
- name: Run tests and flake8
run: python .github/workflows/run_tests.py
- name: Set up Python 3.6
uses: actions/setup-python@v2
with:
python-version: 3.6
- name: Install dependencies
run: python .github/workflows/install_deps.py
- name: Run tests and flake8
run: python .github/workflows/run_tests.py
- name: Set up Python 3.7
uses: actions/setup-python@v2
with:
python-version: 3.7
- name: Install dependencies
run: python .github/workflows/install_deps.py
- name: Run tests and flake8
run: python .github/workflows/run_tests.py
- name: Set up Python 3.8
uses: actions/setup-python@v2
with:
python-version: 3.8
- name: Install dependencies
run: python .github/workflows/install_deps.py
- name: Run tests and flake8
run: python .github/workflows/run_tests.py
- name: Set up Python 3.9
uses: actions/setup-python@v2
with:
python-version: 3.9
- name: Install dependencies
run: python .github/workflows/install_deps.py
- name: Run tests and flake8
run: python .github/workflows/run_tests.py
- name: Combine coverage
run: python .github/workflows/run_tests.py combine
- name: Install coveralls
run: python -m pip install coveralls
- name: Upload Coverage
env:
COVERALLS_FLAG_NAME: unit-${{ matrix.os }}
COVERALLS_PARALLEL: 'true'
GITHUB_TOKEN: ${{ secrets.github_token }}
run: python -m coveralls --service=github
- name: Run codacy-coverage-reporter
uses: codacy/codacy-coverage-reporter-action@master
env:
SECRETS_AVAILABLE: ${{ secrets.CODACY_PROJECT_TOKEN != '' }}
with:
project-token: ${{ secrets.CODACY_PROJECT_TOKEN }}
coverage-reports: coverage.xml
if: ${{ matrix.os == 'ubuntu' && env.SECRETS_AVAILABLE == 'true' }}
behave:
runs-on: ${{ matrix.os }}-latest
env:
DCS: ${{ matrix.dcs }}
ETCDVERSION: 3.3.13
strategy:
fail-fast: false
matrix:
os: [ubuntu]
python-version: [2.7, 3.5, 3.8]
dcs: [etcd, etcd3, consul, exhibitor, kubernetes, raft]
steps:
- uses: actions/checkout@v1
- name: Set up Python
uses: actions/setup-python@v2
with:
python-version: ${{ matrix.python-version }}
- name: Install dependencies
run: python .github/workflows/install_deps.py
- name: Run behave tests
run: python .github/workflows/run_tests.py
- uses: actions/setup-python@v2
with:
python-version: 3.9
- name: Install coveralls
run: python -m pip install coveralls
- name: Upload Coverage
env:
COVERALLS_FLAG_NAME: behave-${{ matrix.os }}-${{ matrix.dcs }}-${{ matrix.python-version }}
COVERALLS_PARALLEL: 'true'
GITHUB_TOKEN: ${{ secrets.github_token }}
run: python -m coveralls --service=github
behavem:
runs-on: ${{ matrix.os }}-latest
env:
DCS: ${{ matrix.dcs }}
ETCDVERSION: 3.3.13
PGVERSION: 12.1-1 # for windows and macos
strategy:
fail-fast: false
matrix:
os: [macos] #, windows]
python-version: [3.7]
dcs: [etcd, etcd3, raft]
steps:
- uses: actions/checkout@v1
- name: Set up Python
uses: actions/setup-python@v2
with:
python-version: ${{ matrix.python-version }}
- name: Install dependencies
run: python .github/workflows/install_deps.py
- name: Run behave tests
run: python .github/workflows/run_tests.py
- name: Install coveralls
run: python -m pip install coveralls
- name: Upload Coverage
env:
COVERALLS_FLAG_NAME: behave-${{ matrix.os }}-${{ matrix.dcs }}-${{ matrix.python-version }}
COVERALLS_PARALLEL: 'true'
GITHUB_TOKEN: ${{ secrets.github_token }}
run: python -m coveralls --service=github
coveralls-finish:
name: Finalize coveralls.io
needs: [unit, behave, behavem]
runs-on: ubuntu-latest
steps:
- uses: actions/setup-python@v2
- run: python -m pip install coveralls
- run: python -m coveralls --service=github --finish
env:
GITHUB_TOKEN: ${{ secrets.github_token }}
-164
View File
@@ -1,164 +0,0 @@
sudo: true
dist: trusty
language: python
addons:
apt:
packages:
- expect-dev # for unbuffer
env:
global:
- ETCDVERSION=3.3.13 ZKVERSION=3.4.14 CONSULVERSION=0.7.4
- PYVERSIONS="2.7 3.6"
- BOTO_CONFIG=/doesnotexist
matrix:
include:
- python: "3.5"
env: PYVERSIONS="2.7 3.5 3.6" TEST_SUITE="python setup.py"
- python: "3.6"
env: DCS="etcd" TEST_SUITE="behave"
- python: "3.6"
env: DCS="etcd3" TEST_SUITE="behave"
- python: "3.6"
env: DCS="exhibitor" TEST_SUITE="behave"
- python: "3.6"
env: DCS="consul" TEST_SUITE="behave"
- python: "3.6"
env: DCS="raft" TEST_SUITE="behave"
- python: "3.6"
env: DCS="kubernetes" TEST_SUITE="behave"
branches:
only:
- master
- /^v\d+\.\d+(\.\d+)?$/
cache:
directories:
- $HOME/mycache
before_cache:
- |
rm -fr $HOME/mycache/python*
for pv in $PYVERSIONS; do
fpv=$(basename $(readlink $HOME/virtualenv/python${pv}))
mv $HOME/virtualenv/${fpv} $HOME/mycache/${fpv}
done
install:
- |
set -e
if [[ $TEST_SUITE == "behave" ]]; then
function get_consul() {
CC=~/mycache/consul_${CONSULVERSION}
if [[ ! -x $CC ]]; then
curl -L https://releases.hashicorp.com/consul/${CONSULVERSION}/consul_${CONSULVERSION}_linux_amd64.zip \
| gunzip > $CC
[[ ${PIPESTATUS[0]} == 0 ]] || return 1
chmod +x $CC
fi
ln -s $CC consul
}
function get_etcd() {
EC=~/mycache/etcd_${ETCDVERSION}
if [[ ! -x $EC ]]; then
rm -rf ~/mycache/etcd_*
curl -L https://github.com/coreos/etcd/releases/download/v${ETCDVERSION}/etcd-v${ETCDVERSION}-linux-amd64.tar.gz \
| tar xz -C . --strip=1 --wildcards --no-anchored etcd
[[ ${PIPESTATUS[0]} == 0 ]] || return 1
mv etcd $EC
fi
ln -s $EC etcd
}
function get_etcd3() {
get_etcd
}
function get_kubernetes() {
wget -O localkube "https://storage.googleapis.com/minikube/k8sReleases/v1.7.0/localkube-linux-amd64"
chmod +x localkube
sudo nohup ./localkube --logtostderr=true --enable-dns=false > localkube.log 2>&1 &
echo "Waiting for localkube to start..."
if ! timeout 120 sh -c "while ! curl -ks http://127.0.0.1:8080/ >/dev/null; do sleep 1; done"; then
sudo cat localkube.log
echo "localkube did not start"
exit 1
fi
echo "Check certificate permissions"
sudo chmod 644 /var/lib/localkube/certs/*
sudo ls -altr /var/lib/localkube/certs/
echo "Set up .kube/config"
mkdir ~/.kube
echo -e "apiVersion: v1\nclusters:\n- cluster:\n certificate-authority: /var/lib/localkube/certs/ca.crt\n server: https://127.0.0.1:8443\n name: local\ncontexts:\n- context:\n cluster: local\n user: myself\n name: local\ncurrent-context: local\nkind: Config\npreferences: {}\nusers:\n- name: myself\n user:\n client-certificate: /var/lib/localkube/certs/apiserver.crt\n client-key: /var/lib/localkube/certs/apiserver.key\n" > ~/.kube/config
}
function get_exhibitor() {
ZC=~/mycache/zookeeper-${ZKVERSION}
if [[ ! -d $ZC ]]; then
curl -L http://www.apache.org/dist/zookeeper/zookeeper-${ZKVERSION}/zookeeper-${ZKVERSION}.tar.gz | tar xz
[[ ${PIPESTATUS[0]} == 0 ]] || return 1
mv zookeeper-${ZKVERSION}/conf/zoo_sample.cfg zookeeper-${ZKVERSION}/conf/zoo.cfg
mv zookeeper-${ZKVERSION} $ZC
fi
$ZC/bin/zkServer.sh start
# following lines are 'emulating' exhibitor REST API
while true; do
echo -e 'HTTP/1.0 200 OK\nContent-Type: application/json\n\n{"servers":["127.0.0.1"],"port":2181}' \
| nc -l 8181 &> /dev/null
done&
}
function get_raft() {
return 0
}
attempt_num=1
until get_${DCS}; do
[[ $attempt_num -ge 3 ]] && exit 1
echo "Attempt $attempt_num failed! Trying again in $attempt_num seconds..."
sleep $(( attempt_num++ ))
done
fi
for pv in $PYVERSIONS; do
fpv=$(basename $(readlink $HOME/virtualenv/python$pv))
if [[ -d ~/mycache/${fpv} ]]; then
mv ~/virtualenv/${fpv} ~/virtualenv/${fpv}.bckp
mv ~/mycache/${fpv} ~/virtualenv/${fpv}
fi
source ~/virtualenv/python${pv}/bin/activate
# explicitly install all needed python modules to cache them
for p in '-r requirements.txt' 'psycopg2-binary behave codacy-coverage coverage coveralls flake8 mock pytest-cov pytest setuptools'; do
pip install $p --upgrade
done
done
script:
- |
for pv in $PYVERSIONS; do
source ~/virtualenv/python${pv}/bin/activate
if [[ $TEST_SUITE = "behave" ]]; then
echo Running integration tests using python${pv}
if ! PATH=.:/usr/lib/postgresql/9.6/bin:$PATH unbuffer $TEST_SUITE; then
# output all log files when tests are failing
grep . features/output/*_failed/*postgres?.*
exit 1
fi
else
echo Running unit tests using python${pv}
unbuffer $TEST_SUITE test
$TEST_SUITE flake8
fi
mv .coverage /tmp/.coverage.$pv
done
mv /tmp/.coverage.* .
python -m coverage combine
set +e
after_success:
# before_cache is executed earlier than after_success, so we need to restore one of virtualenv directories
- fpv=$(basename $(readlink $HOME/virtualenv/python3.6)) && mv $HOME/mycache/${fpv} $HOME/virtualenv/${fpv}
- coveralls
- if [[ $TEST_SUITE != "behave" ]]; then python-codacy-coverage -r coverage.xml; fi
- if [[ $DCS == "exhibitor" ]]; then ~/mycache/zookeeper-${ZKVERSION}/bin/zkServer.sh stop; fi
- sudo kill $(jobs -p)
+2
View File
@@ -12,6 +12,8 @@ Patroni is a template for you to create your own customized, high-availability s
We call Patroni a "template" because it is far from being a one-size-fits-all or plug-and-play replication system. It will have its own caveats. Use wisely. We call Patroni a "template" because it is far from being a one-size-fits-all or plug-and-play replication system. It will have its own caveats. Use wisely.
Currently supported PostgreSQL versions: 9.3 to 13.
**Note to Kubernetes users**: Patroni can run natively on top of Kubernetes. Take a look at the `Kubernetes <https://github.com/zalando/patroni/blob/master/docs/kubernetes.rst>`__ chapter of the Patroni documentation. **Note to Kubernetes users**: Patroni can run natively on top of Kubernetes. Take a look at the `Kubernetes <https://github.com/zalando/patroni/blob/master/docs/kubernetes.rst>`__ chapter of the Patroni documentation.
.. contents:: .. contents::
+33
View File
@@ -10,6 +10,39 @@ Chatting
Just want to chat with other Patroni users? Looking for interactive troubleshooting help? Join us on channel #patroni in the `PostgreSQL Slack <https://postgres-slack.herokuapp.com/>`__. Just want to chat with other Patroni users? Looking for interactive troubleshooting help? Join us on channel #patroni in the `PostgreSQL Slack <https://postgres-slack.herokuapp.com/>`__.
Running tests
-------------
Requirements for running behave tests:
1. PostgreSQL packages need to be installed.
2. PostgreSQL binaries must be available in your `PATH`. You may need to add them to the path with something like `PATH=/usr/lib/postgresql/11/bin:$PATH python -m behave`.
3. If you'd like to test with external DCSs (e.g., Etcd, Consul, and Zookeeper) you'll need the packages installed and respective services running and accepting unencrypted/unprotected connections on localhost and default port. In the case of Etcd or Consul, the behave test suite could start them up if binaries are available in the `PATH`.
Install dependencies:
.. code-block:: bash
# You may want to use Virtualenv or specify pip3.
pip install -r requirements.txt
pip install -r requirements.dev.txt
After you have all dependencies installed, you can run the various test suites:
.. code-block:: bash
# You may want to use Virtualenv or specify python3.
# Run flake8 to check syntax and formatting:
python setup.py flake8
# Run the pytest suite in tests/:
python setup.py test
# Run the behave (https://behave.readthedocs.io/en/latest/) test suite in features/;
# modify DCS as desired (raft has no dependencies so is the easiest to start with):
DCS=raft python -m behave
Reporting issues Reporting issues
---------------- ----------------
+2
View File
@@ -162,7 +162,9 @@ REST API
- **PATRONI\_RESTAPI\_PASSWORD**: Basic-auth password to protect unsafe REST API endpoints. - **PATRONI\_RESTAPI\_PASSWORD**: Basic-auth password to protect unsafe REST API endpoints.
- **PATRONI\_RESTAPI\_CERTFILE**: Specifies the file with the certificate in the PEM format. If the certfile is not specified or is left empty, the API server will work without SSL. - **PATRONI\_RESTAPI\_CERTFILE**: Specifies the file with the certificate in the PEM format. If the certfile is not specified or is left empty, the API server will work without SSL.
- **PATRONI\_RESTAPI\_KEYFILE**: Specifies the file with the secret key in the PEM format. - **PATRONI\_RESTAPI\_KEYFILE**: Specifies the file with the secret key in the PEM format.
- **PATRONI\_RESTAPI\_KEYFILE\_PASSWORD**: Specifies a password for decrypting the keyfile.
- **PATRONI\_RESTAPI\_CAFILE**: Specifies the file with the CA_BUNDLE with certificates of trusted CAs to use while verifying client certs. - **PATRONI\_RESTAPI\_CAFILE**: Specifies the file with the CA_BUNDLE with certificates of trusted CAs to use while verifying client certs.
- **PATRONI\_RESTAPI\_CIPHERS**: (optional) Specifies the permitted cipher suites (e.g. "ECDHE-RSA-AES256-GCM-SHA384:DHE-RSA-AES256-GCM-SHA384:ECDHE-RSA-AES128-GCM-SHA256:DHE-RSA-AES128-GCM-SHA256:!SSLv1:!SSLv2:!SSLv3:!TLSv1:!TLSv1.1")
- **PATRONI\_RESTAPI\_VERIFY\_CLIENT**: ``none`` (default), ``optional`` or ``required``. When ``none`` REST API will not check client certificates. When ``required`` client certificates are required for all REST API calls. When ``optional`` client certificates are required for all unsafe REST API endpoints. When ``required`` is used, then client authentication succeeds, if the certificate signature verification succeeds. For ``optional`` the client cert will only be checked for ``PUT``, ``POST``, ``PATCH``, and ``DELETE`` requests. - **PATRONI\_RESTAPI\_VERIFY\_CLIENT**: ``none`` (default), ``optional`` or ``required``. When ``none`` REST API will not check client certificates. When ``required`` client certificates are required for all REST API calls. When ``optional`` client certificates are required for all unsafe REST API endpoints. When ``required`` is used, then client authentication succeeds, if the certificate signature verification succeeds. For ``optional`` the client cert will only be checked for ``PUT``, ``POST``, ``PATCH``, and ``DELETE`` requests.
- **PATRONI\_RESTAPI\_HTTP\_EXTRA\_HEADERS**: (optional) HTTP headers let the REST API server pass additional information with an HTTP response. - **PATRONI\_RESTAPI\_HTTP\_EXTRA\_HEADERS**: (optional) HTTP headers let the REST API server pass additional information with an HTTP response.
- **PATRONI\_RESTAPI\_HTTPS\_EXTRA\_HEADERS**: (optional) HTTPS headers let the REST API server pass additional information with an HTTP response when TLS is enabled. This will also pass additional information set in ``http_extra_headers``. - **PATRONI\_RESTAPI\_HTTPS\_EXTRA\_HEADERS**: (optional) HTTPS headers let the REST API server pass additional information with an HTTP response when TLS is enabled. This will also pass additional information set in ``http_extra_headers``.
+17
View File
@@ -163,3 +163,20 @@ When connecting from an application, always use a non-superuser. Patroni require
:target: https://travis-ci.org/zalando/patroni :target: https://travis-ci.org/zalando/patroni
.. |Coverage Status| image:: https://coveralls.io/repos/zalando/patroni/badge.svg?branch=master .. |Coverage Status| image:: https://coveralls.io/repos/zalando/patroni/badge.svg?branch=master
:target: https://coveralls.io/r/zalando/patroni?branch=master :target: https://coveralls.io/r/zalando/patroni?branch=master
Testing Your HA Solution
--------------------------------------
Testing an HA solution is a time consuming process, with many variables. This is particularly true considering a cross-platform application. You need a trained system administrator or a consultant to do this work. It is not something we can cover in depth in the documentaiton.
That said, here are some pieces of your infrastructure you should be sure to test:
* Network (the network in front of your system as well as the NICs [physical or virtual] themselves)
* Disk IO
* file limits (nofile in Linux)
* RAM. Even if you have oomkiller turned off as suggested, the unavailability of RAM could cause issues.
* CPU
* Virtualization Contention (overcommitting the hypervisor)
* Any cgroup limitation (likely to be related to the above)
* ``kill -9`` of any postgres process (except postmaster!). This is a decent simulation of a segfault.
One thing that you should not do is run ``kill -9`` on a postmaster process. This is because doing so does not mimic any real life scenario. If you are concerned your infrastructure is insecure and an attacker could run ``kill -9``, no amount of HA process is going to fix that. The attacker will simply kill the process again, or cause chaos in another way.
+30 -1
View File
@@ -15,6 +15,7 @@ Dynamic configuration is stored in the DCS (Distributed Configuration Store) and
- **ttl**: the TTL to acquire the leader lock (in seconds). Think of it as the length of time before initiation of the automatic failover process. Default value: 30 - **ttl**: the TTL to acquire the leader lock (in seconds). Think of it as the length of time before initiation of the automatic failover process. Default value: 30
- **retry\_timeout**: timeout for DCS and PostgreSQL operation retries (in seconds). DCS or network issues shorter than this will not cause Patroni to demote the leader. Default value: 10 - **retry\_timeout**: timeout for DCS and PostgreSQL operation retries (in seconds). DCS or network issues shorter than this will not cause Patroni to demote the leader. Default value: 10
- **maximum\_lag\_on\_failover**: the maximum bytes a follower may lag to be able to participate in leader election. - **maximum\_lag\_on\_failover**: the maximum bytes a follower may lag to be able to participate in leader election.
- **maximum\_lag\_on\_syncnode**: the maximum bytes a synchronous follower may lag before it is considered as an unhealthy candidate and swapped by healthy asynchronous follower. Patroni utilize the max replica lsn if there is more than one follower, otherwise it will use leader's current wal lsn. Default is -1, Patroni will not take action to swap synchronous unhealthy follower when the value is set to 0 or below. Please set the value high enough so Patroni won't swap synchrounous follower fequently during high transaction volume.
- **max\_timelines\_history**: maximum number of timeline history items kept in DCS. Default value: 0. When set to 0, it keeps the full history in DCS. - **max\_timelines\_history**: maximum number of timeline history items kept in DCS. Default value: 0. When set to 0, it keeps the full history in DCS.
- **master\_start\_timeout**: the amount of time a master is allowed to recover from failures before failover is triggered (in seconds). Default is 300 seconds. When set to 0 failover is done immediately after a crash is detected if possible. When using asynchronous replication a failover can cause lost transactions. Worst case failover time for master failure is: loop\_wait + master\_start\_timeout + loop\_wait, unless master\_start\_timeout is zero, in which case it's just loop\_wait. Set the value according to your durability/availability tradeoff. - **master\_start\_timeout**: the amount of time a master is allowed to recover from failures before failover is triggered (in seconds). Default is 300 seconds. When set to 0 failover is done immediately after a crash is detected if possible. When using asynchronous replication a failover can cause lost transactions. Worst case failover time for master failure is: loop\_wait + master\_start\_timeout + loop\_wait, unless master\_start\_timeout is zero, in which case it's just loop\_wait. Set the value according to your durability/availability tradeoff.
- **master\_stop\_timeout**: The number of seconds Patroni is allowed to wait when stopping Postgres and effective only when synchronous_mode is enabled. When set to > 0 and the synchronous_mode is enabled, Patroni sends SIGKILL to the postmaster if the stop operation is running for more than the value set by master_stop_timeout. Set the value according to your durability/availability tradeoff. If the parameter is not set or set <= 0, master_stop_timeout does not apply. - **master\_stop\_timeout**: The number of seconds Patroni is allowed to wait when stopping Postgres and effective only when synchronous_mode is enabled. When set to > 0 and the synchronous_mode is enabled, Patroni sends SIGKILL to the postmaster if the stop operation is running for more than the value set by master_stop_timeout. Set the value according to your durability/availability tradeoff. If the parameter is not set or set <= 0, master_stop_timeout does not apply.
@@ -38,6 +39,32 @@ Dynamic configuration is stored in the DCS (Distributed Configuration Store) and
- **type**: slot type. Could be ``physical`` or ``logical``. If the slot is logical, you have to additionally define ``database`` and ``plugin``. - **type**: slot type. Could be ``physical`` or ``logical``. If the slot is logical, you have to additionally define ``database`` and ``plugin``.
- **database**: the database name where logical slots should be created. - **database**: the database name where logical slots should be created.
- **plugin**: the plugin name for the logical slot. - **plugin**: the plugin name for the logical slot.
- **ignore_slots**: list of sets of replication slot properties for which Patroni should ignore matching slots. This configuration/feature/etc. is useful when some replication slots are managed outside of Patroni. Any subset of matching properties will cause a slot to be ignored.
- **name**: the name of the replication slot.
- **type**: slot type. Can be ``physical`` or ``logical``. If the slot is logical, you may additionally define ``database`` and/or ``plugin``.
- **database**: the database name (when matching a ``logical`` slot).
- **plugin**: the logical decoding plugin (when matching a ``logical`` slot).
Note: **slots** is a hashmap while **ignore_slots** is an array. For example:
.. code:: YAML
slots:
permanent_logical_slot_name:
type: logical
database: my_db
plugin: test_decoding
permanent_physical_slot_name:
type: physical
...
ignore_slots:
- name: ignored_logical_slot_name
type: logical
database: my_db
plugin: test_decoding
- name: ignored_physical_slot_name
type: physical
...
Global/Universal Global/Universal
---------------- ----------------
@@ -209,7 +236,7 @@ Raft
- Q: It is possible to run Patroni and PostgreSQL only on two nodes? - Q: It is possible to run Patroni and PostgreSQL only on two nodes?
A: Yes, on the third node you can run ``patroni_raft_controller`` (without Patroni and PostgreSQL). In such setup one can temporary loose one node without affecting the primary. A: Yes, on the third node you can run ``patroni_raft_controller`` (without Patroni and PostgreSQL). In such a setup, one can temporarily lose one node without affecting the primary.
.. _postgresql_settings: .. _postgresql_settings:
@@ -295,7 +322,9 @@ REST API
- **password**: Basic-auth password to protect unsafe REST API endpoints. - **password**: Basic-auth password to protect unsafe REST API endpoints.
- **certfile**: (optional): Specifies the file with the certificate in the PEM format. If the certfile is not specified or is left empty, the API server will work without SSL. - **certfile**: (optional): Specifies the file with the certificate in the PEM format. If the certfile is not specified or is left empty, the API server will work without SSL.
- **keyfile**: (optional): Specifies the file with the secret key in the PEM format. - **keyfile**: (optional): Specifies the file with the secret key in the PEM format.
- **keyfile_password**: (optional): Specifies a password for decrypting the keyfile.
- **cafile**: (optional): Specifies the file with the CA_BUNDLE with certificates of trusted CAs to use while verifying client certs. - **cafile**: (optional): Specifies the file with the CA_BUNDLE with certificates of trusted CAs to use while verifying client certs.
- **ciphers**: (optional): Specifies the permitted cipher suites (e.g. "ECDHE-RSA-AES256-GCM-SHA384:DHE-RSA-AES256-GCM-SHA384:ECDHE-RSA-AES128-GCM-SHA256:DHE-RSA-AES128-GCM-SHA256:!SSLv1:!SSLv2:!SSLv3:!TLSv1:!TLSv1.1")
- **verify\_client**: (optional): ``none`` (default), ``optional`` or ``required``. When ``none`` REST API will not check client certificates. When ``required`` client certificates are required for all REST API calls. When ``optional`` client certificates are required for all unsafe REST API endpoints. When ``required`` is used, then client authentication succeeds, if the certificate signature verification succeeds. For ``optional`` the client cert will only be checked for ``PUT``, ``POST``, ``PATCH``, and ``DELETE`` requests. - **verify\_client**: (optional): ``none`` (default), ``optional`` or ``required``. When ``none`` REST API will not check client certificates. When ``required`` client certificates are required for all REST API calls. When ``optional`` client certificates are required for all unsafe REST API endpoints. When ``required`` is used, then client authentication succeeds, if the certificate signature verification succeeds. For ``optional`` the client cert will only be checked for ``PUT``, ``POST``, ``PATCH``, and ``DELETE`` requests.
- **http\_extra\_headers**: (optional): HTTP headers let the REST API server pass additional information with an HTTP response. - **http\_extra\_headers**: (optional): HTTP headers let the REST API server pass additional information with an HTTP response.
- **https\_extra\_headers**: (optional): HTTPS headers let the REST API server pass additional information with an HTTP response when TLS is enabled. This will also pass additional information set in ``http_extra_headers``. - **https\_extra\_headers**: (optional): HTTPS headers let the REST API server pass additional information with an HTTP response when TLS is enabled. This will also pass additional information set in ``http_extra_headers``.
+3
View File
@@ -10,6 +10,8 @@ Patroni is a template for you to create your own customized, high-availability s
We call Patroni a "template" because it is far from being a one-size-fits-all or plug-and-play replication system. It will have its own caveats. Use wisely. There are many ways to run high availability with PostgreSQL; for a list, see the `PostgreSQL Documentation <https://wiki.postgresql.org/wiki/Replication,_Clustering,_and_Connection_Pooling>`__. We call Patroni a "template" because it is far from being a one-size-fits-all or plug-and-play replication system. It will have its own caveats. Use wisely. There are many ways to run high availability with PostgreSQL; for a list, see the `PostgreSQL Documentation <https://wiki.postgresql.org/wiki/Replication,_Clustering,_and_Connection_Pooling>`__.
Currently supported PostgreSQL versions: 9.3 to 13.
**Note to Kubernetes users**: Patroni can run natively on top of Kubernetes. Take a look at the :ref:`Kubernetes <kubernetes>` chapter of the Patroni documentation. **Note to Kubernetes users**: Patroni can run natively on top of Kubernetes. Take a look at the :ref:`Kubernetes <kubernetes>` chapter of the Patroni documentation.
@@ -20,6 +22,7 @@ We call Patroni a "template" because it is far from being a one-size-fits-all or
README README
dynamic_configuration dynamic_configuration
rest_api rest_api
existing_data
ENVIRONMENT ENVIRONMENT
SETTINGS SETTINGS
security security
+96
View File
@@ -3,6 +3,102 @@
Release notes Release notes
============= =============
Version 2.0.2
-------------
**New features**
- Ability to ignore externally managed replication slots (James Coleman)
Patroni is trying to remove any replication slot which is unknown to it, but there are certainly cases when replication slots should be managed externally. From now on it is possible to configure slots that should not be removed.
- Added support for cipher suite limitation for REST API (Gunnar "Nick" Bluth)
It could be configured via ``restapi.ciphers`` or the ``PATRONI_RESTAPI_CIPHERS`` environment variable.
- Added support for encrypted TLS keys for REST API (Jonathan S. Katz)
It could be configured via ``restapi.keyfile_password`` or the ``PATRONI_RESTAPI_KEYFILE_PASSWORD`` environment variable.
- Constant time comparison of REST API authentication credentials (Alex Brasetvik)
Use ``hmac.compare_digest()`` instead of ``==``, which is vulnerable to timing attack.
- Choose synchronous nodes based on replication lag (Krishna Sarabu)
If the replication lag on the synchronous node starts exceeding the configured threshold it could be demoted to asynchronous and/or replaced by the other node. Behaviour is controlled with ``maximum_lag_on_syncnode``.
**Stability improvements**
- Start postgres with ``hot_standby = off`` when doing custom bootstrap (Igor Yanchenko)
During custom bootstrap Patroni is restoring the basebackup, starting Postgres up, and waiting until recovery finishes. Some PostgreSQL parameters on the standby can't be smaller than on the primary and if the new value (restored from WAL) is higher than the configured one, Postgres panics and stops. In order to avoid such behavior we will do custom bootstrap without ``hot_standby`` mode.
- Warn the user if the required watchdog is not healthy (Nicolas Thauvin)
When the watchdog device is not writable or missing in required mode, the member cannot be promoted. Added a warning to show the user where to search for this misconfiguration.
- Better verbosity for single-user mode recovery (Alexander Kukushkin)
If Patroni notices that PostgreSQL wasn't shutdown clearly, in certain cases the crash-recovery is executed by starting Postgres in single-user mode. It could happen that the recovery failed (for example due to the lack of space on disk) but errors were swallowed.
- Added compatibility with ``python-consul2`` module (Alexander, Wilfried Roset)
The good old ``python-consul`` is not maintained since a few years, therefore someone created a fork with new features and bug-fixes.
- Don't use ``bypass_api_service`` when running ``patronictl`` (Alexander)
When a K8s pod is running in a non-``default`` namespace it does not necessarily have enough permissions to query the ``kubernetes`` endpoint. In this case Patroni shows the warning and ignores the ``bypass_api_service`` setting. In case of ``patronictl`` the warning was a bit annoying.
- Create ``raft.data_dir`` if it doesn't exists or make sure that it is writable (Mark Mercado)
Improves user-friendliness and usability.
**Bugfixes**
- Don't interrupt restart or promote if lost leader lock in pause (Alexander)
In pause it is allowed to run postgres as primary without lock.
- Fixed issue with ``shutdown_request()`` in the REST API (Nicolas Limage)
In order to improve handling of SSL connections and delay the handshake until thread is started Patroni overrides a few methods in the ``HTTPServer``. The ``shutdown_request()`` method was forgotten.
- Fixed issue with sleep time when using Zookeeper (Alexander)
There were chances that Patroni was sleeping up to twice longer between running HA code.
- Fixed invalid ``os.symlink()`` calls when moving data directory after failed bootstrap (Andrew L'Ecuyer)
If the bootstrap failed Patroni is renaming data directory, pg_wal, and all tablespaces. After that it updates symlinks so filesystem remains consistent. The symlink creation was failing due to the ``src`` and ``dst`` arguments being swapped.
- Fixed bug in the post_bootstrap() method (Alexander)
If the superuser password wasn't configured Patroni was failing to call the ``post_init`` script and therefore the whole bootstrap was failing.
- Fixed an issues with pg_rewind in the standby cluster (Alexander)
If the superuser name is different from Postgres, the ``pg_rewind`` in the standby cluster was failing because the connection string didn't contain the database name.
- Exit only if authentication with Etcd v3 explicitly failed (Alexander)
On start Patroni performs discovery of Etcd cluster topology and authenticates if it is necessarily. It could happen that one of etcd servers is not accessible, Patroni was trying to perform authentication on this server and failing instead of retrying with the next node.
- Handle case with psutil cmdline() returning empty list (Alexander)
Zombie processes are still postmasters children, but they don't have cmdline()
- Treat ``PATRONI_KUBERNETES_USE_ENDPOINTS`` environment variable as boolean (Alexander)
Not doing so was making impossible disabling ``kubernetes.use_endpoints`` via environment.
- Improve handling of concurrent endpoint update errors (Alexander)
Patroni will explicitly query the current endpoint object, verify that the current pod still holds the leader lock and repeat the update.
Version 2.0.1 Version 2.0.1
------------- -------------
+5
View File
@@ -143,6 +143,10 @@ possible to specify a ``basebackup`` configuration section. Same rules as with t
namely, only long (with --) options should be specified there. Not all parameters make sense, if you override a connection namely, only long (with --) options should be specified there. Not all parameters make sense, if you override a connection
string or provide an option to created tar-ed or compressed base backups, patroni won't be able to make a replica out string or provide an option to created tar-ed or compressed base backups, patroni won't be able to make a replica out
of it. There is no validation performed on the names or values of the parameters passed to the ``basebackup`` section. of it. There is no validation performed on the names or values of the parameters passed to the ``basebackup`` section.
Also note that in case symlinks are used for the WAL folder it is up to the user to specify the correct ``--waldir``
path as an option, so that after replica buildup or re-initialization the symlink would persist. This option is supported
only since v10 though.
You can specify basebackup parameters as either a map (key-value pairs) or a list of elements, where each element You can specify basebackup parameters as either a map (key-value pairs) or a list of elements, where each element
could be either a key-value pair or a single key (for options that does not receive any values, for instance, ``--verbose``). could be either a key-value pair or a single key (for options that does not receive any values, for instance, ``--verbose``).
Consider those 2 examples: Consider those 2 examples:
@@ -162,6 +166,7 @@ and
basebackup: basebackup:
- verbose - verbose
- max-rate: '100M' - max-rate: '100M'
- waldir: /pg-wal-mount/external-waldir
If all replica creation methods fail, Patroni will try again all methods in order during the next event loop cycle. If all replica creation methods fail, Patroni will try again all methods in order during the next event loop cycle.
+19
View File
@@ -28,6 +28,25 @@ Feature: basic replication
When I issue a GET request to http://127.0.0.1:8009/async When I issue a GET request to http://127.0.0.1:8009/async
Then I receive a response code 200 Then I receive a response code 200
Scenario: check stuck sync replica
Given I issue a PATCH request to http://127.0.0.1:8008/config with {"maximum_lag_on_syncnode": 15000000, "postgresql": {"parameters": {"synchronous_commit": "remote_apply"}}}
Then I receive a response code 200
And I create table on postgres0
And table mytest is present on postgres1 after 2 seconds
And table mytest is present on postgres2 after 2 seconds
When I pause wal replay on postgres2
And I load data on postgres0
Then "sync" key in DCS has sync_standby=postgres1 after 15 seconds
And I resume wal replay on postgres2
And I sleep for 2 seconds
And I issue a GET request to http://127.0.0.1:8009/sync
Then I receive a response code 200
When I issue a GET request to http://127.0.0.1:8010/async
Then I receive a response code 200
When I issue a PATCH request to http://127.0.0.1:8008/config with {"maximum_lag_on_syncnode": -1, "postgresql": {"parameters": {"synchronous_commit": "on"}}}
Then I receive a response code 200
And I drop table on postgres0
Scenario: check multi sync replication Scenario: check multi sync replication
Given I issue a PATCH request to http://127.0.0.1:8008/config with {"synchronous_node_count": 2} Given I issue a PATCH request to http://127.0.0.1:8008/config with {"synchronous_node_count": 2}
Then I receive a response code 200 Then I receive a response code 200
+25 -3
View File
@@ -13,6 +13,8 @@ import threading
import time import time
import yaml import yaml
from six.moves.BaseHTTPServer import BaseHTTPRequestHandler, HTTPServer
@six.add_metaclass(abc.ABCMeta) @six.add_metaclass(abc.ABCMeta)
class AbstractController(object): class AbstractController(object):
@@ -176,6 +178,7 @@ class PatroniController(AbstractController):
config['name'] = name config['name'] = name
config['postgresql']['data_dir'] = self._data_dir config['postgresql']['data_dir'] = self._data_dir
config['postgresql']['basebackup'] = [{'checkpoint': 'fast'}]
config['postgresql']['use_unix_socket'] = os.name != 'nt' # windows doesn't yet support unix-domain sockets config['postgresql']['use_unix_socket'] = os.name != 'nt' # windows doesn't yet support unix-domain sockets
config['postgresql']['pgpass'] = os.path.join(tempfile.gettempdir(), 'pgpass_' + name) config['postgresql']['pgpass'] = os.path.join(tempfile.gettempdir(), 'pgpass_' + name)
config['postgresql']['parameters'].update({ config['postgresql']['parameters'].update({
@@ -191,6 +194,8 @@ class PatroniController(AbstractController):
if custom_config is not None: if custom_config is not None:
self.recursive_update(config, custom_config) self.recursive_update(config, custom_config)
self.recursive_update(config, {
'bootstrap': {'dcs': {'postgresql': {'parameters': {'wal_keep_segments': 100}}}}})
if config['postgresql'].get('callbacks', {}).get('on_role_change'): if config['postgresql'].get('callbacks', {}).get('on_role_change'):
config['postgresql']['callbacks']['on_role_change'] += ' ' + str(self.__PORT) config['postgresql']['callbacks']['on_role_change'] += ' ' + str(self.__PORT)
@@ -558,11 +563,28 @@ class ZooKeeperController(AbstractDcsController):
return False return False
class MockExhibitor(BaseHTTPRequestHandler):
def do_GET(self):
self.send_response(200)
self.end_headers()
self.wfile.write(b'{"servers":["127.0.0.1"],"port":2181}')
def log_message(self, fmt, *args):
pass
class ExhibitorController(ZooKeeperController): class ExhibitorController(ZooKeeperController):
def __init__(self, context): def __init__(self, context):
super(ExhibitorController, self).__init__(context, False) super(ExhibitorController, self).__init__(context, False)
os.environ.update({'PATRONI_EXHIBITOR_HOSTS': 'localhost', 'PATRONI_EXHIBITOR_PORT': '8181'}) port = 8181
exhibitor = HTTPServer(('', port), MockExhibitor)
exhibitor.daemon_thread = True
exhibitor_thread = threading.Thread(target=exhibitor.serve_forever)
exhibitor_thread.daemon = True
exhibitor_thread.start()
os.environ.update({'PATRONI_EXHIBITOR_HOSTS': 'localhost', 'PATRONI_EXHIBITOR_PORT': str(port)})
class RaftController(AbstractDcsController): class RaftController(AbstractDcsController):
@@ -849,8 +871,8 @@ class WatchdogMonitor(object):
# actions to execute on start/stop of the tests and before running invidual features # actions to execute on start/stop of the tests and before running invidual features
def before_all(context): def before_all(context):
os.environ.update({'PATRONI_RESTAPI_USERNAME': 'username', 'PATRONI_RESTAPI_PASSWORD': 'password'}) os.environ.update({'PATRONI_RESTAPI_USERNAME': 'username', 'PATRONI_RESTAPI_PASSWORD': 'password'})
context.ci = 'TRAVIS_BUILD_NUMBER' in os.environ or 'BUILD_NUMBER' in os.environ context.ci = any(a in os.environ for a in ('TRAVIS_BUILD_NUMBER', 'BUILD_NUMBER', 'GITHUB_ACTIONS'))
context.timeout_multiplier = 2 if context.ci else 1 context.timeout_multiplier = 5 if context.ci else 1 # MacOS sometimes is VERY slow
context.pctl = PatroniPoolController(context) context.pctl = PatroniPoolController(context)
context.dcs_ctl = context.pctl.known_dcs[context.pctl.dcs](context) context.dcs_ctl = context.pctl.known_dcs[context.pctl.dcs](context)
context.dcs_ctl.start() context.dcs_ctl.start()
+61
View File
@@ -0,0 +1,61 @@
Feature: ignored slots
Scenario: check ignored slots aren't removed on failover/switchover
Given I start postgres1
Then postgres1 is a leader after 10 seconds
And there is a non empty initialize key in DCS after 15 seconds
When I issue a PATCH request to http://127.0.0.1:8009/config with {"loop_wait": 2, "ignore_slots": [{"name": "unmanaged_slot_0", "database": "postgres", "plugin": "test_decoding", "type": "logical"}, {"name": "unmanaged_slot_1", "database": "postgres", "plugin": "test_decoding"}, {"name": "unmanaged_slot_2", "database": "postgres"}, {"name": "unmanaged_slot_3"}], "postgresql": {"parameters": {"wal_level": "logical"}}}
Then I receive a response code 200
And Response on GET http://127.0.0.1:8009/config contains ignore_slots after 10 seconds
# Make sure the wal_level has been changed.
When I shut down postgres1
And I start postgres1
Then postgres1 is a leader after 10 seconds
And "members/postgres1" key in DCS has role=master after 3 seconds
# Make sure Patroni has finished telling Postgres it should be accepting writes.
And postgres1 role is the primary after 20 seconds
# 1. Create our test logical replication slot.
# Test that ny subset of attributes in the ignore slots matcher is enough to match a slot
# by using 3 different slots.
When I create a logical replication slot unmanaged_slot_0 on postgres1 with the test_decoding plugin
And I create a logical replication slot unmanaged_slot_1 on postgres1 with the test_decoding plugin
And I create a logical replication slot unmanaged_slot_2 on postgres1 with the test_decoding plugin
And I create a logical replication slot unmanaged_slot_3 on postgres1 with the test_decoding plugin
And I create a logical replication slot dummy_slot on postgres1 with the test_decoding plugin
# It seems like it'd be obvious that these slots exist since we just created them,
# but Patroni can actually end up dropping them almost immediately, so it's helpful
# to verify they exist before we begin testing whether they persist through failover
# cycles.
Then postgres1 has a logical replication slot named unmanaged_slot_0 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_1 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_2 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_3 with the test_decoding plugin
When I start postgres0
Then "members/postgres0" key in DCS has role=replica after 3 seconds
And postgres0 role is the secondary after 20 seconds
# Verify that the replica has advanced beyond the point in the WAL
# where we created the replication slot so that on the next failover
# cycle we don't accidentally rewind to before the slot creation.
And replication works from postgres1 to postgres0 after 20 seconds
When I shut down postgres1
Then "members/postgres0" key in DCS has role=master after 3 seconds
# 2. After a failover the server (now a replica) still has the slot.
When I start postgres1
Then postgres1 role is the secondary after 20 seconds
And "members/postgres1" key in DCS has role=replica after 3 seconds
# give Patroni time to sync replication slots
And I sleep for 2 seconds
And postgres1 has a logical replication slot named unmanaged_slot_0 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_1 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_2 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_3 with the test_decoding plugin
And postgres1 does not have a logical replication slot named dummy_slot
# 3. After a failover the server (now a master) still has the slot.
When I shut down postgres0
Then "members/postgres1" key in DCS has role=master after 3 seconds
And postgres1 has a logical replication slot named unmanaged_slot_0 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_1 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_2 with the test_decoding plugin
And postgres1 has a logical replication slot named unmanaged_slot_3 with the test_decoding plugin
+31
View File
@@ -33,6 +33,37 @@ def add_table(context, table_name, pg_name):
assert False, "Error creating table {0} on {1}: {2}".format(table_name, pg_name, e) assert False, "Error creating table {0} on {1}: {2}".format(table_name, pg_name, e)
@step('I {action:w} wal replay on {pg_name:w}')
def toggle_wal_replay(context, action, pg_name):
# pause or resume the wal replay process
try:
version = context.pctl.query(pg_name, "select pg_catalog.pg_read_file('PG_VERSION', 0, 2)").fetchone()
wal = version and version[0] and int(version[0].split('.')[0]) < 10 and "xlog" or "wal"
context.pctl.query(pg_name, "SELECT pg_{0}_replay_{1}()".format(wal, action))
except pg.Error as e:
assert False, "Error during {0} wal recovery on {1}: {2}".format(action, pg_name, e)
@step('I {action:w} table on {pg_name:w}')
def crdr_mytest(context, action, pg_name):
try:
if (action == "create"):
context.pctl.query(pg_name, "create table if not exists mytest(id Numeric)")
else:
context.pctl.query(pg_name, "drop table if exists mytest")
except pg.Error as e:
assert False, "Error {0} table mytest on {1}: {2}".format(action, pg_name, e)
@step('I load data on {pg_name:w}')
def initiate_load(context, pg_name):
# perform dummy load
try:
context.pctl.query(pg_name, "begin; insert into mytest select r::numeric from generate_series(1, 350000) r; commit;")
except pg.Error as e:
assert False, "Error loading test data on {0}: {1}".format(pg_name, e)
@then('Table {table_name:w} is present on {pg_name:w} after {max_replication_delay:d} seconds') @then('Table {table_name:w} is present on {pg_name:w} after {max_replication_delay:d} seconds')
def table_is_present_on(context, table_name, pg_name, max_replication_delay): def table_is_present_on(context, table_name, pg_name, max_replication_delay):
max_replication_delay *= context.timeout_multiplier max_replication_delay *= context.timeout_multiplier
+8 -3
View File
@@ -12,8 +12,10 @@ def start_patroni_with_a_name_value_tag(context, name, tag_name, tag_value):
@then('There is a {label} with "{content}" in {name:w} data directory') @then('There is a {label} with "{content}" in {name:w} data directory')
def check_label(context, label, content, name): def check_label(context, label, content, name):
label = context.pctl.read_label(name, label) label = context.pctl.read_label(name, label)
if label is None:
label = ""
label = label.replace('\n', '\\n') label = label.replace('\n', '\\n')
assert content in label, "{0} doesn't contain {1}".format(label, content) assert content in label, "\"{0}\" doesn't contain {1}".format(label, content)
@step('I create label with "{content:w}" in {name:w} data directory') @step('I create label with "{content:w}" in {name:w} data directory')
@@ -25,15 +27,18 @@ def write_label(context, content, name):
def check_member(context, name, key, value, time_limit): def check_member(context, name, key, value, time_limit):
time_limit *= context.timeout_multiplier time_limit *= context.timeout_multiplier
max_time = time.time() + int(time_limit) max_time = time.time() + int(time_limit)
dcs_value = None
while time.time() < max_time: while time.time() < max_time:
try: try:
response = json.loads(context.dcs_ctl.query(name)) response = json.loads(context.dcs_ctl.query(name))
if response.get(key) == value: dcs_value = response.get(key)
if dcs_value == value:
return return
except Exception: except Exception:
pass pass
time.sleep(1) time.sleep(1)
assert False, "{0} does not have {1}={2} in dcs after {3} seconds".format(name, key, value, time_limit) assert False, "{0} does not have {1}={2} (found {3}) in dcs after {4} seconds".format(name, key, value,
dcs_value, time_limit)
@step('there is a non empty {key:w} key in DCS after {time_limit:d} seconds') @step('there is a non empty {key:w} key in DCS after {time_limit:d} seconds')
+36
View File
@@ -0,0 +1,36 @@
from behave import step, then
import psycopg2 as pg
@step('I create a logical replication slot {slot_name} on {pg_name:w} with the {plugin:w} plugin')
def create_logical_replication_slot(context, slot_name, pg_name, plugin):
try:
output = context.pctl.query(pg_name, ("SELECT pg_create_logical_replication_slot('{0}', '{1}'),"
" current_database()").format(slot_name, plugin))
print(output.fetchone())
except pg.Error as e:
print(e)
assert False, "Error creating slot {0} on {1} with plugin {2}".format(slot_name, pg_name, plugin)
@then('{pg_name:w} has a logical replication slot named {slot_name} with the {plugin:w} plugin')
def has_logical_replication_slot(context, pg_name, slot_name, plugin):
try:
row = context.pctl.query(pg_name, ("SELECT slot_type, plugin FROM pg_replication_slots"
" WHERE slot_name = '{0}'").format(slot_name)).fetchone()
assert row, "Couldn't find replication slot named {0}".format(slot_name)
assert row[0] == "logical", "Found replication slot named {0} but wasn't a logical slot".format(slot_name)
assert row[1] == plugin, ("Found replication slot named {0} but was using plugin "
"{1} rather than {2}").format(slot_name, row[1], plugin)
except pg.Error:
assert False, "Error looking for slot {0} on {1} with plugin {2}".format(slot_name, pg_name, plugin)
@then('{pg_name:w} does not have a logical replication slot named {slot_name}')
def does_not_have_logical_replication_slot(context, pg_name, slot_name):
try:
row = context.pctl.query(pg_name, ("SELECT 1 FROM pg_replication_slots"
" WHERE slot_name = '{0}'").format(slot_name)).fetchone()
assert not row, "Found unexpected replication slot named {0}".format(slot_name)
except pg.Error:
assert False, "Error looking for slot {0} on {1}".format(slot_name, pg_name)
+15 -5
View File
@@ -1,4 +1,5 @@
import base64 import base64
import hmac
import json import json
import logging import logging
import psycopg2 import psycopg2
@@ -197,7 +198,7 @@ class RestApiHandler(BaseHTTPRequestHandler):
def do_PATCH_config(self): def do_PATCH_config(self):
request = self._read_json_content() request = self._read_json_content()
if request: if request:
cluster = self.server.patroni.dcs.get_cluster() cluster = self.server.patroni.dcs.get_cluster(True)
if not (cluster.config and cluster.config.modify_index): if not (cluster.config and cluster.config.modify_index):
return self.send_error(503) return self.send_error(503)
data = cluster.config.data.copy() data = cluster.config.data.copy()
@@ -558,7 +559,7 @@ class RestApiServer(ThreadingMixIn, HTTPServer, Thread):
fcntl.fcntl(fd, fcntl.F_SETFD, flags | fcntl.FD_CLOEXEC) fcntl.fcntl(fd, fcntl.F_SETFD, flags | fcntl.FD_CLOEXEC)
def check_basic_auth_key(self, key): def check_basic_auth_key(self, key):
return self.__auth_key == key return hmac.compare_digest(self.__auth_key, key.encode('utf-8'))
def check_auth_header(self, auth_header): def check_auth_header(self, auth_header):
if self.__auth_key: if self.__auth_key:
@@ -634,7 +635,10 @@ class RestApiServer(ThreadingMixIn, HTTPServer, Thread):
if self.__protocol == 'https': if self.__protocol == 'https':
import ssl import ssl
ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH, cafile=ssl_options.get('cafile')) ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH, cafile=ssl_options.get('cafile'))
ctx.load_cert_chain(certfile=ssl_options['certfile'], keyfile=ssl_options.get('keyfile')) if ssl_options.get('ciphers'):
ctx.set_ciphers(ssl_options['ciphers'])
ctx.load_cert_chain(certfile=ssl_options['certfile'], keyfile=ssl_options.get('keyfile'),
password=ssl_options.get('keyfile_password'))
verify_client = ssl_options.get('verify_client') verify_client = ssl_options.get('verify_client')
if verify_client: if verify_client:
modes = {'none': ssl.CERT_NONE, 'optional': ssl.CERT_OPTIONAL, 'required': ssl.CERT_REQUIRED} modes = {'none': ssl.CERT_NONE, 'optional': ssl.CERT_OPTIONAL, 'required': ssl.CERT_REQUIRED}
@@ -664,11 +668,17 @@ class RestApiServer(ThreadingMixIn, HTTPServer, Thread):
newsock = (sock, newsock) newsock = (sock, newsock)
return newsock, addr return newsock, addr
def shutdown_request(self, request):
if isinstance(request, tuple):
_, request = request # SSLSocket
return super(RestApiServer, self).shutdown_request(request)
def reload_config(self, config): def reload_config(self, config):
if 'listen' not in config: # changing config in runtime if 'listen' not in config: # changing config in runtime
raise ValueError('Can not find "restapi.listen" config') raise ValueError('Can not find "restapi.listen" config')
ssl_options = {n: config[n] for n in ('certfile', 'keyfile', 'cafile') if n in config} ssl_options = {n: config[n] for n in ('certfile', 'keyfile', 'keyfile_password',
'cafile', 'ciphers') if n in config}
self.http_extra_headers = config.get('http_extra_headers') or {} self.http_extra_headers = config.get('http_extra_headers') or {}
self.http_extra_headers.update((config.get('https_extra_headers') or {}) if ssl_options.get('certfile') else {}) self.http_extra_headers.update((config.get('https_extra_headers') or {}) if ssl_options.get('certfile') else {})
@@ -679,7 +689,7 @@ class RestApiServer(ThreadingMixIn, HTTPServer, Thread):
if self.__listen != config['listen'] or self.__ssl_options != ssl_options: if self.__listen != config['listen'] or self.__ssl_options != ssl_options:
self.__initialize(config['listen'], ssl_options) self.__initialize(config['listen'], ssl_options)
self.__auth_key = base64.b64encode(config['auth'].encode('utf-8')).decode('utf-8') if 'auth' in config else None self.__auth_key = base64.b64encode(config['auth'].encode('utf-8')) if 'auth' in config else None
self.connection_string = uri(self.__protocol, config.get('connect_address') or self.__listen, 'patroni') self.connection_string = uri(self.__protocol, config.get('connect_address') or self.__listen, 'patroni')
@staticmethod @staticmethod
+6 -3
View File
@@ -60,6 +60,7 @@ class Config(object):
__DEFAULT_CONFIG = { __DEFAULT_CONFIG = {
'ttl': 30, 'loop_wait': 10, 'retry_timeout': 10, 'ttl': 30, 'loop_wait': 10, 'retry_timeout': 10,
'maximum_lag_on_failover': 1048576, 'maximum_lag_on_failover': 1048576,
'maximum_lag_on_syncnode': -1,
'check_timeline': False, 'check_timeline': False,
'master_start_timeout': 300, 'master_start_timeout': 300,
'master_stop_timeout': 0, 'master_stop_timeout': 0,
@@ -265,8 +266,9 @@ class Config(object):
if value: if value:
ret[section][param] = value ret[section][param] = value
_set_section_values('restapi', ['listen', 'connect_address', 'certfile', 'keyfile', 'cafile', 'verify_client', _set_section_values('restapi', ['listen', 'connect_address', 'certfile', 'keyfile', 'keyfile_password',
'http_extra_headers', 'https_extra_headers']) 'cafile', 'ciphers', 'verify_client', 'http_extra_headers',
'https_extra_headers'])
_set_section_values('ctl', ['insecure', 'cacert', 'certfile', 'keyfile']) _set_section_values('ctl', ['insecure', 'cacert', 'certfile', 'keyfile'])
_set_section_values('postgresql', ['listen', 'connect_address', 'config_dir', 'data_dir', 'pgpass', 'bin_dir']) _set_section_values('postgresql', ['listen', 'connect_address', 'config_dir', 'data_dir', 'pgpass', 'bin_dir'])
_set_section_values('log', ['level', 'traceback_level', 'format', 'dateformat', 'max_queue_size', _set_section_values('log', ['level', 'traceback_level', 'format', 'dateformat', 'max_queue_size',
@@ -337,7 +339,7 @@ class Config(object):
value = value and _parse_list(value) value = value and _parse_list(value)
elif suffix == 'LABELS': elif suffix == 'LABELS':
value = _parse_dict(value) value = _parse_dict(value)
elif suffix in ('USE_PROXIES', 'REGISTER_SERVICE', 'BYPASS_API_SERVICE'): elif suffix in ('USE_PROXIES', 'REGISTER_SERVICE', 'USE_ENDPOINTS', 'BYPASS_API_SERVICE'):
value = parse_bool(value) value = parse_bool(value)
if value: if value:
ret[name.lower()][suffix.lower()] = value ret[name.lower()][suffix.lower()] = value
@@ -412,6 +414,7 @@ class Config(object):
'synchronous_mode', 'synchronous_mode',
'synchronous_mode_strict', 'synchronous_mode_strict',
'synchronous_node_count', 'synchronous_node_count',
'maximum_lag_on_syncnode'
) )
pg_config.update({p: config[p] for p in updated_fields if p in config}) pg_config.update({p: config[p] for p in updated_fields if p in config})
+9 -25
View File
@@ -24,20 +24,22 @@ import yaml
from click import ClickException from click import ClickException
from collections import defaultdict from collections import defaultdict
from contextlib import contextmanager from contextlib import contextmanager
from patroni.dcs import get_dcs as _get_dcs
from patroni.exceptions import PatroniException
from patroni.postgresql import Postgresql
from patroni.postgresql.misc import postgres_version_to_int
from patroni.utils import cluster_as_json, patch_config, polling_loop
from patroni.request import PatroniRequest
from patroni.version import __version__
from prettytable import ALL, FRAME, PrettyTable from prettytable import ALL, FRAME, PrettyTable
from six.moves.urllib_parse import urlparse from six.moves.urllib_parse import urlparse
try: try:
from ydiff import markup_to_pager, PatchStream from ydiff import markup_to_pager, PatchStream
except ImportError: # pragma: no cover except ImportError: # pragma: no cover
from cdiff import markup_to_pager, PatchStream from cdiff import markup_to_pager, PatchStream
from .dcs import get_dcs as _get_dcs
from .exceptions import PatroniException
from .postgresql import Postgresql
from .postgresql.misc import postgres_version_to_int
from .utils import cluster_as_json, find_executable, patch_config, polling_loop
from .request import PatroniRequest
from .version import __version__
CONFIG_DIR_PATH = click.get_app_dir('patroni') CONFIG_DIR_PATH = click.get_app_dir('patroni')
CONFIG_FILE_PATH = os.path.join(CONFIG_DIR_PATH, 'patronictl.yaml') CONFIG_FILE_PATH = os.path.join(CONFIG_DIR_PATH, 'patronictl.yaml')
DCS_DEFAULTS = {'zookeeper': {'port': 2181, 'template': "zookeeper:\n hosts: ['{host}:{port}']"}, DCS_DEFAULTS = {'zookeeper': {'port': 2181, 'template': "zookeeper:\n hosts: ['{host}:{port}']"},
@@ -1171,24 +1173,6 @@ def apply_yaml_file(data, filename):
return format_config_for_editing(changed_data), changed_data return format_config_for_editing(changed_data), changed_data
def find_executable(executable, path=None):
_, ext = os.path.splitext(executable)
if (sys.platform == 'win32') and (ext != '.exe'):
executable = executable + '.exe'
if os.path.isfile(executable):
return executable
if path is None:
path = os.environ.get('PATH', os.defpath)
for p in path.split(os.pathsep):
f = os.path.join(p, executable)
if os.path.isfile(f):
return f
def invoke_editor(before_editing, cluster_name): def invoke_editor(before_editing, cluster_name):
"""Starts editor command to edit configuration in human readable format """Starts editor command to edit configuration in human readable format
+4
View File
@@ -341,6 +341,10 @@ class ClusterConfig(namedtuple('ClusterConfig', 'index,data,modify_index')):
self.data.get('permanent_slots') or self.data.get('slots') self.data.get('permanent_slots') or self.data.get('slots')
) or {} ) or {}
@property
def ignore_slots_matchers(self):
return isinstance(self.data, dict) and self.data.get('ignore_slots') or []
@property @property
def max_timelines_history(self): def max_timelines_history(self):
return self.data.get('max_timelines_history', 0) return self.data.get('max_timelines_history', 0)
+17 -8
View File
@@ -8,6 +8,7 @@ import ssl
import time import time
import urllib3 import urllib3
from collections import namedtuple
from consul import ConsulException, NotFound, base from consul import ConsulException, NotFound, base
from urllib3.exceptions import HTTPError from urllib3.exceptions import HTTPError
from six.moves.urllib.parse import urlencode, urlparse, quote from six.moves.urllib.parse import urlencode, urlparse, quote
@@ -36,6 +37,9 @@ class InvalidSession(ConsulException):
"""invalid session""" """invalid session"""
Response = namedtuple('Response', 'code,headers,body,content')
class HTTPClient(object): class HTTPClient(object):
def __init__(self, host='127.0.0.1', port=8500, token=None, scheme='http', verify=True, cert=None, ca_cert=None): def __init__(self, host='127.0.0.1', port=8500, token=None, scheme='http', verify=True, cert=None, ca_cert=None):
@@ -71,16 +75,17 @@ class HTTPClient(object):
@staticmethod @staticmethod
def response(response): def response(response):
data = response.data.decode('utf-8') content = response.data
body = content.decode('utf-8')
if response.status == 500: if response.status == 500:
msg = '{0} {1}'.format(response.status, data) msg = '{0} {1}'.format(response.status, body)
if data.startswith('Invalid Session TTL'): if body.startswith('Invalid Session TTL'):
raise InvalidSessionTTL(msg) raise InvalidSessionTTL(msg)
elif data.startswith('invalid session'): elif body.startswith('invalid session'):
raise InvalidSession(msg) raise InvalidSession(msg)
else: else:
raise ConsulInternalError(msg) raise ConsulInternalError(msg)
return base.Response(response.status, response.headers, data) return Response(response.status, response.headers, body, content)
def uri(self, path, params=None): def uri(self, path, params=None):
return '{0}{1}{2}'.format(self.base_uri, path, params and '?' + urlencode(params) or '') return '{0}{1}{2}'.format(self.base_uri, path, params and '?' + urlencode(params) or '')
@@ -89,7 +94,7 @@ class HTTPClient(object):
if method not in ('get', 'post', 'put', 'delete'): if method not in ('get', 'post', 'put', 'delete'):
raise AttributeError("HTTPClient instance has no attribute '{0}'".format(method)) raise AttributeError("HTTPClient instance has no attribute '{0}'".format(method))
def wrapper(callback, path, params=None, data=''): def wrapper(callback, path, params=None, data='', headers=None):
# python-consul doesn't allow to specify ttl smaller then 10 seconds # python-consul doesn't allow to specify ttl smaller then 10 seconds
# because session_ttl_min defaults to 10s, so we have to do this ugly dirty hack... # because session_ttl_min defaults to 10s, so we have to do this ugly dirty hack...
if method == 'put' and path == '/v1/session/create': if method == 'put' and path == '/v1/session/create':
@@ -110,8 +115,9 @@ class HTTPClient(object):
kwargs['timeout'] = timeout + max(timeout/15.0, 1) kwargs['timeout'] = timeout + max(timeout/15.0, 1)
else: else:
kwargs['timeout'] = self._read_timeout kwargs['timeout'] = self._read_timeout
kwargs['headers'] = (headers or {}).copy()
kwargs['headers'].update(urllib3.make_headers(user_agent=USER_AGENT))
token = params.pop('token', self.token) if isinstance(params, dict) else self.token token = params.pop('token', self.token) if isinstance(params, dict) else self.token
kwargs['headers'] = urllib3.make_headers(user_agent=USER_AGENT)
if token: if token:
kwargs['headers']['X-Consul-Token'] = token kwargs['headers']['X-Consul-Token'] = token
return callback(self.response(self.http.request(method.upper(), self.uri(path, params), **kwargs))) return callback(self.response(self.http.request(method.upper(), self.uri(path, params), **kwargs)))
@@ -126,7 +132,7 @@ class ConsulClient(base.Consul):
self.token = kwargs.get('token') self.token = kwargs.get('token')
super(ConsulClient, self).__init__(*args, **kwargs) super(ConsulClient, self).__init__(*args, **kwargs)
def connect(self, *args, **kwargs): def http_connect(self, *args, **kwargs):
kwargs.update(dict(zip(['host', 'port', 'scheme', 'verify'], args))) kwargs.update(dict(zip(['host', 'port', 'scheme', 'verify'], args)))
if self._cert: if self._cert:
kwargs['cert'] = self._cert kwargs['cert'] = self._cert
@@ -136,6 +142,9 @@ class ConsulClient(base.Consul):
kwargs['token'] = self.token kwargs['token'] = self.token
return HTTPClient(**kwargs) return HTTPClient(**kwargs)
def connect(self, *args, **kwargs):
return self.http_connect(*args, **kwargs)
def reload_config(self, config): def reload_config(self, config):
self.http.token = self.token = config.get('token') self.http.token = self.token = config.get('token')
self.consistency = config.get('consistency', 'default') self.consistency = config.get('consistency', 'default')
+1 -1
View File
@@ -404,7 +404,7 @@ class EtcdClient(AbstractEtcdClientWithFailover):
if self.http is not None: if self.http is not None:
try: try:
self.http.clear() self.http.clear()
except (ReferenceError, TypeError): except (ReferenceError, TypeError, AttributeError):
pass pass
def _prepare_get_members(self, etcd_nodes): def _prepare_get_members(self, etcd_nodes):
+5 -1
View File
@@ -98,6 +98,10 @@ class UserEmpty(InvalidArgument):
error = "etcdserver: user name is empty" error = "etcdserver: user name is empty"
class AuthFailed(InvalidArgument):
error = "etcdserver: authentication failed, invalid user ID or password"
class PermissionDenied(Etcd3ClientError): class PermissionDenied(Etcd3ClientError):
code = GRPCCode.PermissionDenied code = GRPCCode.PermissionDenied
error = "etcdserver: permission denied" error = "etcdserver: permission denied"
@@ -186,7 +190,7 @@ class Etcd3Client(AbstractEtcdClientWithFailover):
try: try:
self.authenticate() self.authenticate()
except Exception as e: except AuthFailed as e:
logger.fatal('Etcd3 authentication failed: %r', e) logger.fatal('Etcd3 authentication failed: %r', e)
sys.exit(1) sys.exit(1)
+16 -6
View File
@@ -631,7 +631,8 @@ class Kubernetes(AbstractDCS):
port.update({n: p[n] for n in ('name', 'protocol') if p.get(n)}) port.update({n: p[n] for n in ('name', 'protocol') if p.get(n)})
self.__ports.append(k8s_client.V1EndpointPort(**port)) self.__ports.append(k8s_client.V1EndpointPort(**port))
self._api = CoreV1ApiProxy(config.get('use_endpoints'), config.get('bypass_api_service')) bypass_api_service = not config.get('patronictl') and config.get('bypass_api_service')
self._api = CoreV1ApiProxy(config.get('use_endpoints'), bypass_api_service)
self._should_create_config_service = self._api.use_endpoints self._should_create_config_service = self._api.use_endpoints
self.reload_config(config) self.reload_config(config)
# leader_observed_record, leader_resource_version, and leader_observed_time are used only for leader race! # leader_observed_record, leader_resource_version, and leader_observed_time are used only for leader race!
@@ -904,13 +905,23 @@ class Kubernetes(AbstractDCS):
except (RetryFailedError, K8sException): except (RetryFailedError, K8sException):
return False return False
deadline = retry.stoptime - time.time() retry.deadline = retry.stoptime - time.time()
if deadline < 2: if retry.deadline < 1:
return False return False
retry.sleep_func(1) # Give a chance for ObjectCache to receive the latest version # Try to get the latest version directly from K8s API instead of relying on async cache
try:
kind = retry(self._api.read_namespaced_kind, self.leader_path, self._namespace)
except Exception as e:
logger.error('Failed to get the leader object "%s": %r', self.leader_path, e)
return False
self._kinds.set(self.leader_path, kind)
retry.deadline = retry.stoptime - time.time()
if retry.deadline < 0.5:
return False
kind = self._kinds.get(self.leader_path)
kind_annotations = kind and kind.metadata.annotations or {} kind_annotations = kind and kind.metadata.annotations or {}
kind_resource_version = kind and kind.metadata.resource_version kind_resource_version = kind and kind.metadata.resource_version
@@ -918,7 +929,6 @@ class Kubernetes(AbstractDCS):
if kind and (kind_annotations.get(self._LEADER) != self._name or kind_resource_version == resource_version): if kind and (kind_annotations.get(self._LEADER) != self._name or kind_resource_version == resource_version):
return False return False
retry.deadline = deadline - 1 # Update deadline and retry
return self.patch_or_create(self.leader_path, annotations, kind_resource_version, ips=ips, retry=_retry) return self.patch_or_create(self.leader_path, annotations, kind_resource_version, ips=ips, retry=_retry)
def update_leader(self, last_operation, access_is_restricted=False): def update_leader(self, last_operation, access_is_restricted=False):
+6
View File
@@ -5,6 +5,7 @@ import threading
import time import time
from patroni.dcs import AbstractDCS, ClusterConfig, Cluster, Failover, Leader, Member, SyncState, TimelineHistory from patroni.dcs import AbstractDCS, ClusterConfig, Cluster, Failover, Leader, Member, SyncState, TimelineHistory
from ..utils import validate_directory
from pysyncobj import SyncObj, SyncObjConf, replicated, FAIL_REASON from pysyncobj import SyncObj, SyncObjConf, replicated, FAIL_REASON
from pysyncobj.transport import Node, TCPTransport, CONNECTION_STATE from pysyncobj.transport import Node, TCPTransport, CONNECTION_STATE
@@ -275,6 +276,11 @@ class Raft(AbstractDCS):
partner_addrs.append(self_addr) partner_addrs.append(self_addr)
self_addr = None self_addr = None
# Create raft data_dir if necessary
raft_data_dir = config.get('data_dir', '')
if raft_data_dir != '':
validate_directory(raft_data_dir)
ready_event = threading.Event() ready_event = threading.Event()
file_template = os.path.join(config.get('data_dir', ''), (self_addr or '')) file_template = os.path.join(config.get('data_dir', ''), (self_addr or ''))
conf = SyncObjConf(password=config.get('password'), appendEntriesUseBatch=False, conf = SyncObjConf(password=config.get('password'), appendEntriesUseBatch=False,
+6 -4
View File
@@ -123,13 +123,14 @@ class ZooKeeper(AbstractDCS):
# the same time, set_ttl method will reestablish connection and return # the same time, set_ttl method will reestablish connection and return
# `!True`, otherwise we will close existing connection and let kazoo # `!True`, otherwise we will close existing connection and let kazoo
# open the new one. # open the new one.
if not self.set_ttl(int(config['ttl'] * 1000)) and loop_wait_changed: if not self.set_ttl(config['ttl']) and loop_wait_changed:
self._client._connection._socket.close() self._client._connection._socket.close()
def set_ttl(self, ttl): def set_ttl(self, ttl):
"""It is not possible to change ttl (session_timeout) in zookeeper without """It is not possible to change ttl (session_timeout) in zookeeper without
destroying old session and creating the new one. This method returns `!True` destroying old session and creating the new one. This method returns `!True`
if session_timeout has been changed (`restart()` has been called).""" if session_timeout has been changed (`restart()` has been called)."""
ttl = int(ttl * 1000)
if self._client._session_timeout != ttl: if self._client._session_timeout != ttl:
self._client._session_timeout = ttl self._client._session_timeout = ttl
self._client.restart() self._client.restart()
@@ -137,7 +138,7 @@ class ZooKeeper(AbstractDCS):
@property @property
def ttl(self): def ttl(self):
return self._client._session_timeout return self._client._session_timeout / 1000.0
def set_retry_timeout(self, retry_timeout): def set_retry_timeout(self, retry_timeout):
retry = self._client.retry if isinstance(self._client.retry, KazooRetry) else self._client._retry retry = self._client.retry if isinstance(self._client.retry, KazooRetry) else self._client._retry
@@ -372,6 +373,7 @@ class ZooKeeper(AbstractDCS):
return self.set_sync_state_value("{}", index) return self.set_sync_state_value("{}", index)
def watch(self, leader_index, timeout): def watch(self, leader_index, timeout):
if super(ZooKeeper, self).watch(leader_index, timeout) and not self._fetch_optime: ret = super(ZooKeeper, self).watch(leader_index, timeout)
if ret and not self._fetch_optime:
self._fetch_cluster = True self._fetch_cluster = True
return self._fetch_cluster return ret or self._fetch_cluster
+8 -3
View File
@@ -460,7 +460,9 @@ class Ha(object):
if self.is_synchronous_mode(): if self.is_synchronous_mode():
sync_node_count = self.patroni.config['synchronous_node_count'] sync_node_count = self.patroni.config['synchronous_node_count']
current = self.cluster.sync.leader and self.cluster.sync.members or [] current = self.cluster.sync.leader and self.cluster.sync.members or []
picked, allow_promote = self.state_handler.pick_synchronous_standby(self.cluster, sync_node_count) picked, allow_promote = self.state_handler.pick_synchronous_standby(self.cluster, sync_node_count,
self.patroni.config[
'maximum_lag_on_syncnode'])
if set(picked) != set(current): if set(picked) != set(current):
# update synchronous standby list in dcs temporarily to point to common nodes in current and picked # update synchronous standby list in dcs temporarily to point to common nodes in current and picked
sync_common = list(set(current).intersection(set(allow_promote))) sync_common = list(set(current).intersection(set(allow_promote)))
@@ -484,7 +486,9 @@ class Ha(object):
# Wait for PostgreSQL to enable synchronous mode and see if we can immediately set sync_standby # Wait for PostgreSQL to enable synchronous mode and see if we can immediately set sync_standby
time.sleep(2) time.sleep(2)
_, allow_promote = self.state_handler.pick_synchronous_standby(self.cluster, _, allow_promote = self.state_handler.pick_synchronous_standby(self.cluster,
sync_node_count) sync_node_count,
self.patroni.config[
'maximum_lag_on_syncnode'])
if allow_promote and set(allow_promote) != set(sync_common): if allow_promote and set(allow_promote) != set(sync_common):
try: try:
cluster = self.dcs.get_cluster() cluster = self.dcs.get_cluster()
@@ -782,6 +786,7 @@ class Ha(object):
return self.manual_failover_process_no_leader() return self.manual_failover_process_no_leader()
if not self.watchdog.is_healthy: if not self.watchdog.is_healthy:
logger.warning('Watchdog device is not usable')
return False return False
# When in sync mode, only last known master and sync standby are allowed to promote automatically. # When in sync mode, only last known master and sync standby are allowed to promote automatically.
@@ -1182,7 +1187,7 @@ class Ha(object):
return 'terminated crash recovery because of startup timeout' return 'terminated crash recovery because of startup timeout'
return 'updated leader lock during ' + self._async_executor.scheduled_action return 'updated leader lock during ' + self._async_executor.scheduled_action
elif not self.state_handler.bootstrapping: elif not self.state_handler.bootstrapping and not self.is_paused():
# Don't have lock, make sure we are not promoting or starting up a master in the background # Don't have lock, make sure we are not promoting or starting up a master in the background
if self._async_executor.scheduled_action == 'promote': if self._async_executor.scheduled_action == 'promote':
with self._async_response: with self._async_response:
+21 -11
View File
@@ -915,7 +915,7 @@ class Postgresql(object):
new_name = '{0}_{1}'.format(pg_wal_realpath, postfix) new_name = '{0}_{1}'.format(pg_wal_realpath, postfix)
os.rename(pg_wal_realpath, new_name) os.rename(pg_wal_realpath, new_name)
os.unlink(source) os.unlink(source)
os.symlink(source, new_name) os.symlink(new_name, source)
# Move user defined tablespace directory # Move user defined tablespace directory
for (source, pg_tsp_rpath) in self.pg_tblspc_realpaths().items(): for (source, pg_tsp_rpath) in self.pg_tblspc_realpaths().items():
@@ -923,7 +923,7 @@ class Postgresql(object):
new_name = '{0}_{1}'.format(pg_tsp_rpath, postfix) new_name = '{0}_{1}'.format(pg_tsp_rpath, postfix)
os.rename(pg_tsp_rpath, new_name) os.rename(pg_tsp_rpath, new_name)
os.unlink(source) os.unlink(source)
os.symlink(source, new_name) os.symlink(new_name, source)
new_name = '{0}_{1}'.format(self._data_dir, postfix) new_name = '{0}_{1}'.format(self._data_dir, postfix)
logger.info('renaming data directory to %s', new_name) logger.info('renaming data directory to %s', new_name)
@@ -962,11 +962,15 @@ class Postgresql(object):
def _get_synchronous_commit_param(self): def _get_synchronous_commit_param(self):
return self.query("SHOW synchronous_commit").fetchone()[0] return self.query("SHOW synchronous_commit").fetchone()[0]
def pick_synchronous_standby(self, cluster, sync_node_count=1): def pick_synchronous_standby(self, cluster, sync_node_count=1, sync_node_maxlag=-1):
"""Finds the best candidate to be the synchronous standby. """Finds the best candidate to be the synchronous standby.
Current synchronous standby is always preferred, unless it has disconnected or does not want to be a Current synchronous standby is always preferred, unless it has disconnected or does not want to be a
synchronous standby any longer. synchronous standby any longer.
Parameter sync_node_maxlag(maximum_lag_on_syncnode) would help swapping unhealthy sync replica incase
if it stops responding (or hung). Please set the value high enough so it won't unncessarily swap sync
standbys during high loads. Any less or equal of 0 value keep the behavior backward compatible and
will not swap. Please note that it will not also swap sync standbys in case where all replicas are hung.
:returns tuple of candidates list and synchronous standby list. :returns tuple of candidates list and synchronous standby list.
""" """
@@ -975,6 +979,7 @@ class Postgresql(object):
members = {m.name.lower(): m for m in cluster.members} members = {m.name.lower(): m for m in cluster.members}
candidates = [] candidates = []
sync_nodes = [] sync_nodes = []
replica_list = []
# Pick candidates based on who has higher replay/remote_write/flush lsn. # Pick candidates based on who has higher replay/remote_write/flush lsn.
sync_commit_par = self._get_synchronous_commit_param() sync_commit_par = self._get_synchronous_commit_param()
sort_col = {'remote_apply': 'replay', 'remote_write': 'write'}.get(sync_commit_par, 'flush') sort_col = {'remote_apply': 'replay', 'remote_write': 'write'}.get(sync_commit_par, 'flush')
@@ -984,17 +989,22 @@ class Postgresql(object):
# receiving changes faster than the sync member (very rare but possible). Such cases would # receiving changes faster than the sync member (very rare but possible). Such cases would
# trigger sync standby member swapping frequently and the sort on sync_state desc should # trigger sync standby member swapping frequently and the sort on sync_state desc should
# help in keeping the query result consistent. # help in keeping the query result consistent.
for app_name, state, sync_state in self.query( for app_name, sync_state, replica_lsn in self.query(
"SELECT pg_catalog.lower(application_name), state, sync_state" "SELECT pg_catalog.lower(application_name), sync_state, pg_{2}_{1}_diff({0}_{1}, '0/0')::bigint"
" FROM pg_catalog.pg_stat_replication" " FROM pg_catalog.pg_stat_replication"
" WHERE state = 'streaming'" " WHERE state = 'streaming'"
" ORDER BY sync_state DESC, {0}_{1} DESC".format(sort_col, self.lsn_name)): " ORDER BY sync_state DESC, {0}_{1} DESC".format(sort_col, self.lsn_name, self.wal_name)):
member = members.get(app_name) member = members.get(app_name)
if not member or member.tags.get('nosync', False): if member and not member.tags.get('nosync', False):
continue replica_list.append((member.name, sync_state, replica_lsn))
candidates.append(member.name)
if sync_state == 'sync': max_lsn = max(replica_list, key=lambda x: x[2])[2] if len(replica_list) > 1 else int(str(self.last_operation()))
sync_nodes.append(member.name)
for app_name, sync_state, replica_lsn in replica_list:
if sync_node_maxlag <= 0 or max_lsn - replica_lsn <= sync_node_maxlag:
candidates.append(app_name)
if sync_state == 'sync':
sync_nodes.append(app_name)
if len(candidates) >= sync_node_count: if len(candidates) >= sync_node_count:
break break
+1 -1
View File
@@ -130,7 +130,7 @@ class Bootstrap(object):
# (pghost empty or the default socket directory) connections coming from the local machine. # (pghost empty or the default socket directory) connections coming from the local machine.
r['host'] = 'localhost' # set it to localhost to write into pgpass r['host'] = 'localhost' # set it to localhost to write into pgpass
env = self._postgresql.config.write_pgpass(r) if 'password' in r else None env = self._postgresql.config.write_pgpass(r)
env['PGOPTIONS'] = '-c synchronous_commit=local' env['PGOPTIONS'] = '-c synchronous_commit=local'
try: try:
+5
View File
@@ -387,6 +387,11 @@ class ConfigHandler(object):
if 'custom_conf' not in self._config and not os.path.exists(self._postgresql_base_conf): if 'custom_conf' not in self._config and not os.path.exists(self._postgresql_base_conf):
os.rename(self._postgresql_conf, self._postgresql_base_conf) os.rename(self._postgresql_conf, self._postgresql_base_conf)
# In case we are using custom bootstrap from spilo image with PITR it fails if it contains increasing
# values like Max_connections. We disable hot_standby so it will accept increasing values.
if self._postgresql.bootstrap.running_custom_bootstrap:
configuration['hot_standby'] = 'off'
with ConfigWriter(self._postgresql_conf) as f: with ConfigWriter(self._postgresql_conf) as f:
include = self._config.get('custom_conf') or self._postgresql_base_conf_name include = self._config.get('custom_conf') or self._postgresql_base_conf_name
f.writeline("include '{0}'\n".format(ConfigWriter.escape(include))) f.writeline("include '{0}'\n".format(ConfigWriter.escape(include)))
+3 -3
View File
@@ -187,10 +187,10 @@ class PostmasterProcess(psutil.Process):
user_backends_cmdlines = [] user_backends_cmdlines = []
for child in children: for child in children:
try: try:
cmdline = child.cmdline()[0] cmdline = child.cmdline()
if not aux_proc_re.match(cmdline): if cmdline and not aux_proc_re.match(cmdline[0]):
user_backends.append(child) user_backends.append(child)
user_backends_cmdlines.append(cmdline) user_backends_cmdlines.append(cmdline[0])
except psutil.NoSuchProcess: except psutil.NoSuchProcess:
pass pass
if user_backends: if user_backends:
+21 -10
View File
@@ -59,11 +59,10 @@ class Rewind(object):
if self.can_rewind_or_reinitialize_allowed and self._state != REWIND_STATUS.NEED: if self.can_rewind_or_reinitialize_allowed and self._state != REWIND_STATUS.NEED:
self._state = REWIND_STATUS.CHECK self._state = REWIND_STATUS.CHECK
def check_leader_is_not_in_recovery(self, **kwargs): @staticmethod
if not kwargs.get('database'): def check_leader_is_not_in_recovery(conn_kwargs):
kwargs['database'] = self._postgresql.database
try: try:
with get_connection_cursor(connect_timeout=3, options='-c statement_timeout=2000', **kwargs) as cur: with get_connection_cursor(connect_timeout=3, options='-c statement_timeout=2000', **conn_kwargs) as cur:
cur.execute('SELECT pg_catalog.pg_is_in_recovery()') cur.execute('SELECT pg_catalog.pg_is_in_recovery()')
if not cur.fetchone()[0]: if not cur.fetchone()[0]:
return True return True
@@ -165,6 +164,12 @@ class Rewind(object):
logger.info('master: history=%s', '\n'.join(history_show)) logger.info('master: history=%s', '\n'.join(history_show))
def _conn_kwargs(self, member, auth):
ret = member.conn_kwargs(auth)
if not ret.get('database'):
ret['database'] = self._postgresql.database
return ret
def _check_timeline_and_lsn(self, leader): def _check_timeline_and_lsn(self, leader):
in_recovery, local_timeline, local_lsn = self._get_local_timeline_lsn() in_recovery, local_timeline, local_lsn = self._get_local_timeline_lsn()
if local_timeline is None or local_lsn is None: if local_timeline is None or local_lsn is None:
@@ -174,7 +179,7 @@ class Rewind(object):
if leader.member.data.get('role') != 'master': if leader.member.data.get('role') != 'master':
return return
# standby cluster # standby cluster
elif not self.check_leader_is_not_in_recovery(**leader.conn_kwargs(self._postgresql.config.replication)): elif not self.check_leader_is_not_in_recovery(self._conn_kwargs(leader, self._postgresql.config.replication)):
return return
history = need_rewind = None history = need_rewind = None
@@ -333,7 +338,7 @@ class Rewind(object):
return logger.warning('Can not run pg_rewind because postgres is still running') return logger.warning('Can not run pg_rewind because postgres is still running')
# prepare pg_rewind connection # prepare pg_rewind connection
r = leader.conn_kwargs(self._postgresql.config.rewind_credentials) r = self._conn_kwargs(leader, self._postgresql.config.rewind_credentials)
# 1. make sure that we are really trying to rewind from the master # 1. make sure that we are really trying to rewind from the master
# 2. make sure that pg_control contains the new timeline by: # 2. make sure that pg_control contains the new timeline by:
@@ -341,17 +346,17 @@ class Rewind(object):
# waiting until Patroni on the master will expose checkpoint_after_promote=True # waiting until Patroni on the master will expose checkpoint_after_promote=True
checkpoint_status = leader.checkpoint_after_promote if isinstance(leader, Leader) else None checkpoint_status = leader.checkpoint_after_promote if isinstance(leader, Leader) else None
if checkpoint_status is None: # master still runs the old Patroni if checkpoint_status is None: # master still runs the old Patroni
leader_status = self._postgresql.checkpoint(leader.conn_kwargs(self._postgresql.config.superuser)) leader_status = self._postgresql.checkpoint(self._conn_kwargs(leader, self._postgresql.config.superuser))
if leader_status: if leader_status:
return logger.warning('Can not use %s for rewind: %s', leader.name, leader_status) return logger.warning('Can not use %s for rewind: %s', leader.name, leader_status)
elif not checkpoint_status: elif not checkpoint_status:
return logger.info('Waiting for checkpoint on %s before rewind', leader.name) return logger.info('Waiting for checkpoint on %s before rewind', leader.name)
elif not self.check_leader_is_not_in_recovery(**r): elif not self.check_leader_is_not_in_recovery(r):
return return
if self.pg_rewind(r): if self.pg_rewind(r):
self._state = REWIND_STATUS.SUCCESS self._state = REWIND_STATUS.SUCCESS
elif not self.check_leader_is_not_in_recovery(**r): elif not self.check_leader_is_not_in_recovery(r):
logger.warning('Failed to rewind because master %s become unreachable', leader.name) logger.warning('Failed to rewind because master %s become unreachable', leader.name)
else: else:
logger.error('Failed to rewind from healty master: %s', leader.name) logger.error('Failed to rewind from healty master: %s', leader.name)
@@ -428,4 +433,10 @@ class Rewind(object):
opts = self.read_postmaster_opts() opts = self.read_postmaster_opts()
opts.update({'archive_mode': 'on', 'archive_command': 'false'}) opts.update({'archive_mode': 'on', 'archive_command': 'false'})
self._postgresql.config.remove_recovery_conf() self._postgresql.config.remove_recovery_conf()
return self.single_user_mode(communicate={}, options=opts) == 0 or None output = {}
ret = self.single_user_mode(communicate=output, options=opts)
if ret != 0:
logger.error('Crash recovery finished with code=%s', ret)
logger.info(' stdout=%s', output['stdout'].decode('utf-8'))
logger.info(' stderr=%s', output['stderr'].decode('utf-8'))
return ret == 0 or None
+10 -2
View File
@@ -33,6 +33,14 @@ class SlotsHandler(object):
self._replication_slots = replication_slots self._replication_slots = replication_slots
self._schedule_load_slots = False self._schedule_load_slots = False
def ignore_replication_slot(self, cluster, name):
slot = self._replication_slots[name]
for matcher in cluster.config.ignore_slots_matchers:
if ((matcher.get("name") is None or matcher["name"] == name)
and all(not matcher.get(a) or matcher[a] == slot.get(a) for a in ('database', 'plugin', 'type'))):
return True
return False
def drop_replication_slot(self, name): def drop_replication_slot(self, name):
cursor = self._query(('SELECT pg_catalog.pg_drop_replication_slot(%s) WHERE EXISTS (SELECT 1 ' + cursor = self._query(('SELECT pg_catalog.pg_drop_replication_slot(%s) WHERE EXISTS (SELECT 1 ' +
'FROM pg_catalog.pg_replication_slots WHERE slot_name = %s AND NOT active)'), name, name) 'FROM pg_catalog.pg_replication_slots WHERE slot_name = %s AND NOT active)'), name, name)
@@ -40,7 +48,7 @@ class SlotsHandler(object):
return cursor.rowcount == 1 return cursor.rowcount == 1
def sync_replication_slots(self, cluster): def sync_replication_slots(self, cluster):
if self._postgresql.major_version >= 90400: if self._postgresql.major_version >= 90400 and cluster.config:
try: try:
self.load_replication_slots() self.load_replication_slots()
@@ -48,7 +56,7 @@ class SlotsHandler(object):
# drop old replication slots which are not presented in desired slots # drop old replication slots which are not presented in desired slots
for name in set(self._replication_slots) - set(slots): for name in set(self._replication_slots) - set(slots):
if not self.drop_replication_slot(name): if not self.ignore_replication_slot(cluster, name) and not self.drop_replication_slot(name):
logger.error("Failed to drop replication slot '%s'", name) logger.error("Failed to drop replication slot '%s'", name)
self._schedule_load_slots = True self._schedule_load_slots = True
+21 -1
View File
@@ -1,3 +1,4 @@
import errno
import json.decoder as json_decoder import json.decoder as json_decoder
import logging import logging
import os import os
@@ -461,7 +462,8 @@ def validate_directory(d, msg="{} {}"):
os.makedirs(d) os.makedirs(d)
except OSError as e: except OSError as e:
logger.error(e) logger.error(e)
raise PatroniException(msg.format(d, "couldn't create the directory")) if e.errno != errno.EEXIST:
raise PatroniException(msg.format(d, "couldn't create the directory"))
elif os.path.isdir(d): elif os.path.isdir(d):
try: try:
fd, tmpfile = tempfile.mkstemp(dir=d) fd, tmpfile = tempfile.mkstemp(dir=d)
@@ -512,3 +514,21 @@ def enable_keepalive(sock, timeout, idle, cnt=3):
for opt in keepalive_socket_options(timeout, idle, cnt): for opt in keepalive_socket_options(timeout, idle, cnt):
sock.setsockopt(*opt) sock.setsockopt(*opt)
def find_executable(executable, path=None):
_, ext = os.path.splitext(executable)
if (sys.platform == 'win32') and (ext != '.exe'):
executable = executable + '.exe'
if os.path.isfile(executable):
return executable
if path is None:
path = os.environ.get('PATH', os.defpath)
for p in path.split(os.pathsep):
f = os.path.join(p, executable)
if os.path.isfile(f):
return f
+28 -14
View File
@@ -4,12 +4,12 @@ import socket
import re import re
import subprocess import subprocess
from patroni.utils import split_host_port, data_directory_is_empty
from patroni.ctl import find_executable
from patroni.dcs import dcs_modules
from patroni.exceptions import ConfigParseError
from six import string_types from six import string_types
from .utils import find_executable, split_host_port, data_directory_is_empty
from .dcs import dcs_modules
from .exceptions import ConfigParseError
def data_directory_empty(data_dir): def data_directory_empty(data_dir):
if os.path.isfile(os.path.join(data_dir, "global", "pg_control")): if os.path.isfile(os.path.join(data_dir, "global", "pg_control")):
@@ -53,11 +53,15 @@ def validate_host_port(host_port, listen=False, multiple_hosts=False):
return True return True
def comma_separated_host_port(string): def validate_host_port_list(value):
assert all([validate_host_port(s.strip()) for s in string.split(",")]), "didn't pass the validation" assert all([validate_host_port(v) for v in value]), "didn't pass the validation"
return True return True
def comma_separated_host_port(string):
return validate_host_port_list([s.strip() for s in string.split(",")])
def validate_host_port_listen(host_port): def validate_host_port_listen(host_port):
return validate_host_port(host_port, listen=True) return validate_host_port(host_port, listen=True)
@@ -291,11 +295,20 @@ def assert_(condition, message="Wrong value"):
userattributes = {"username": "", Optional("password"): ""} userattributes = {"username": "", Optional("password"): ""}
available_dcs = [m.split(".")[-1] for m in dcs_modules()] available_dcs = [m.split(".")[-1] for m in dcs_modules()]
validate_host_port_list.expected_type = list
comma_separated_host_port.expected_type = string_types comma_separated_host_port.expected_type = string_types
validate_connect_address.expected_type = string_types validate_connect_address.expected_type = string_types
validate_host_port_listen.expected_type = string_types validate_host_port_listen.expected_type = string_types
validate_host_port_listen_multiple_hosts.expected_type = string_types validate_host_port_listen_multiple_hosts.expected_type = string_types
validate_data_dir.expected_type = string_types validate_data_dir.expected_type = string_types
validate_etcd = {
Or("host", "hosts", "srv", "url", "proxy"): Case({
"host": validate_host_port,
"hosts": Or(comma_separated_host_port, [validate_host_port]),
"srv": str,
"url": str,
"proxy": str})
}
schema = Schema({ schema = Schema({
"name": str, "name": str,
@@ -320,19 +333,20 @@ schema = Schema({
"host": validate_host_port, "host": validate_host_port,
"url": str}) "url": str})
}, },
"etcd": { "etcd": validate_etcd,
Or("host", "hosts", "srv", "url", "proxy"): Case({ "etcd3": validate_etcd,
"host": validate_host_port,
"hosts": Or(comma_separated_host_port, [validate_host_port]),
"srv": str,
"url": str,
"proxy": str})
},
"exhibitor": { "exhibitor": {
"hosts": [str], "hosts": [str],
"port": lambda i: assert_(int(i) <= 65535), "port": lambda i: assert_(int(i) <= 65535),
Optional("pool_interval"): int Optional("pool_interval"): int
}, },
"raft": {
"self_addr": validate_connect_address,
Optional("bind_addr"): validate_host_port_listen,
"partner_addrs": validate_host_port_list,
Optional("data_dir"): str,
Optional("password"): str
},
"zookeeper": { "zookeeper": {
"hosts": Or(comma_separated_host_port, [validate_host_port]), "hosts": Or(comma_separated_host_port, [validate_host_port]),
}, },
+1 -1
View File
@@ -1 +1 @@
__version__ = '2.0.1' __version__ = '2.0.2'
+1
View File
@@ -112,6 +112,7 @@ postgresql:
basebackup: basebackup:
- verbose - verbose
- max-rate: 100M - max-rate: 100M
# - waldir: /pg-wal-mount/external-waldir # only needed in case pg_wal is symlinked outside of data_dir
# Additional fencing script executed after acquiring the leader lock but before promoting the replica # Additional fencing script executed after acquiring the leader lock but before promoting the replica
#pre_promote: /path/to/pre_promote.sh #pre_promote: /path/to/pre_promote.sh
+8
View File
@@ -0,0 +1,8 @@
psycopg2-binary
behave
coverage
flake8
mock
pytest-cov
pytest
setuptools
+2
View File
@@ -46,6 +46,8 @@ CLASSIFIERS = [
'Programming Language :: Python :: 3.5', 'Programming Language :: Python :: 3.5',
'Programming Language :: Python :: 3.6', 'Programming Language :: Python :: 3.6',
'Programming Language :: Python :: 3.7', 'Programming Language :: Python :: 3.7',
'Programming Language :: Python :: 3.8',
'Programming Language :: Python :: 3.9',
'Programming Language :: Python :: Implementation :: CPython', 'Programming Language :: Python :: Implementation :: CPython',
] ]
+16 -1
View File
@@ -446,10 +446,12 @@ class TestRestApiHandler(unittest.TestCase):
class TestRestApiServer(unittest.TestCase): class TestRestApiServer(unittest.TestCase):
@patch('ssl.SSLContext.load_cert_chain', Mock()) @patch('ssl.SSLContext.load_cert_chain', Mock())
@patch('ssl.SSLContext.set_ciphers', Mock())
@patch('ssl.SSLContext.wrap_socket', Mock(return_value=0)) @patch('ssl.SSLContext.wrap_socket', Mock(return_value=0))
@patch.object(BaseHTTPServer.HTTPServer, '__init__', Mock()) @patch.object(BaseHTTPServer.HTTPServer, '__init__', Mock())
def setUp(self): def setUp(self):
self.srv = MockRestApiServer(Mock(), '', {'listen': '*:8008', 'certfile': 'a', 'verify_client': 'required'}) self.srv = MockRestApiServer(Mock(), '', {'listen': '*:8008', 'certfile': 'a', 'verify_client': 'required',
'ciphers': '!SSLv1:!SSLv2:!SSLv3:!TLSv1:!TLSv1.1'})
@patch.object(BaseHTTPServer.HTTPServer, '__init__', Mock()) @patch.object(BaseHTTPServer.HTTPServer, '__init__', Mock())
def test_reload_config(self): def test_reload_config(self):
@@ -488,3 +490,16 @@ class TestRestApiServer(unittest.TestCase):
mock_accept.return_value = (newsock, '2') mock_accept.return_value = (newsock, '2')
self.srv.socket = Mock() self.srv.socket = Mock()
self.assertEqual(self.srv.get_request(), ((self.srv.socket, newsock), '2')) self.assertEqual(self.srv.get_request(), ((self.srv.socket, newsock), '2'))
@patch.object(MockRestApiServer, 'process_request', Mock(side_effect=RuntimeError))
def test_process_request_error(self):
mock_address = ('127.0.0.1', 55555)
mock_socket = Mock()
mock_ssl_socket = (Mock(), Mock())
for mock_request in (mock_socket, mock_ssl_socket):
with patch.object(
MockRestApiServer,
'get_request',
Mock(return_value=(mock_request, mock_address))
):
self.srv._handle_request_noblock()
+2
View File
@@ -14,6 +14,7 @@ class TestCallbackExecutor(unittest.TestCase):
ce = CallbackExecutor() ce = CallbackExecutor()
ce._kill_children = Mock(side_effect=Exception) ce._kill_children = Mock(side_effect=Exception)
ce._invoke_excepthook = Mock()
self.assertIsNone(ce.call([])) self.assertIsNone(ce.call([]))
ce.join() ce.join()
@@ -30,5 +31,6 @@ class TestCallbackExecutor(unittest.TestCase):
mock_popen.side_effect = Exception mock_popen.side_effect = Exception
ce = CallbackExecutor() ce = CallbackExecutor()
ce._condition.wait = Mock(side_effect=[None, Exception]) ce._condition.wait = Mock(side_effect=[None, Exception])
ce._invoke_excepthook = Mock()
self.assertIsNone(ce.call([])) self.assertIsNone(ce.call([]))
ce.join() ce.join()
+3 -2
View File
@@ -4,7 +4,7 @@ import unittest
from consul import ConsulException, NotFound from consul import ConsulException, NotFound
from mock import Mock, patch from mock import Mock, patch
from patroni.dcs.consul import AbstractDCS, Cluster, Consul, ConsulInternalError, \ from patroni.dcs.consul import AbstractDCS, Cluster, Consul, ConsulInternalError, \
ConsulError, HTTPClient, InvalidSessionTTL, InvalidSession ConsulError, ConsulClient, HTTPClient, InvalidSessionTTL, InvalidSession
from . import SleepException from . import SleepException
@@ -41,7 +41,8 @@ def kv_get(self, key, **kwargs):
class TestHTTPClient(unittest.TestCase): class TestHTTPClient(unittest.TestCase):
def setUp(self): def setUp(self):
self.client = HTTPClient('127.0.0.1', '8500', 'http', False) c = ConsulClient()
self.client = c.http
self.client.http.request = Mock() self.client.http.request = Mock()
def test_get(self): def test_get(self):
+1 -10
View File
@@ -7,7 +7,7 @@ from datetime import datetime, timedelta
from mock import patch, Mock from mock import patch, Mock
from patroni.ctl import ctl, store_config, load_config, output_members, get_dcs, parse_dcs, \ from patroni.ctl import ctl, store_config, load_config, output_members, get_dcs, parse_dcs, \
get_all_members, get_any_member, get_cursor, query_member, configure, PatroniCtlException, apply_config_changes, \ get_all_members, get_any_member, get_cursor, query_member, configure, PatroniCtlException, apply_config_changes, \
format_config_for_editing, show_diff, invoke_editor, format_pg_version, find_executable, CONFIG_FILE_PATH format_config_for_editing, show_diff, invoke_editor, format_pg_version, CONFIG_FILE_PATH
from patroni.dcs.etcd import AbstractEtcdClientWithFailover, Failover from patroni.dcs.etcd import AbstractEtcdClientWithFailover, Failover
from patroni.utils import tzutc from patroni.utils import tzutc
from psycopg2 import OperationalError from psycopg2 import OperationalError
@@ -629,15 +629,6 @@ class TestCtl(unittest.TestCase):
self.assertEqual(format_pg_version(100001), '10.1') self.assertEqual(format_pg_version(100001), '10.1')
self.assertEqual(format_pg_version(90605), '9.6.5') self.assertEqual(format_pg_version(90605), '9.6.5')
@patch('sys.platform', 'win32')
def test_find_executable(self):
with patch('os.path.isfile', Mock(return_value=True)):
self.assertEqual(find_executable('vim'), 'vim.exe')
with patch('os.path.isfile', Mock(return_value=False)):
self.assertIsNone(find_executable('vim'))
with patch('os.path.isfile', Mock(side_effect=[False, True])):
self.assertEqual(find_executable('vim', '/'), '/vim.exe')
@patch('patroni.ctl.get_dcs') @patch('patroni.ctl.get_dcs')
def test_get_members(self, mock_get_dcs): def test_get_members(self, mock_get_dcs):
mock_get_dcs.return_value = self.e mock_get_dcs.return_value = self.e
+1
View File
@@ -110,6 +110,7 @@ class TestDnsCachingResolver(unittest.TestCase):
@patch('socket.getaddrinfo', Mock(side_effect=socket.gaierror)) @patch('socket.getaddrinfo', Mock(side_effect=socket.gaierror))
def test_run(self): def test_run(self):
r = DnsCachingResolver() r = DnsCachingResolver()
r._invoke_excepthook = Mock()
self.assertIsNone(r.resolve_async('', 0)) self.assertIsNone(r.resolve_async('', 0))
r.join() r.join()
+2 -2
View File
@@ -5,7 +5,7 @@ import urllib3
from mock import Mock, patch from mock import Mock, patch
from patroni.dcs.etcd3 import PatroniEtcd3Client, Cluster, Etcd3, Etcd3Error, Etcd3ClientError, RetryFailedError,\ from patroni.dcs.etcd3 import PatroniEtcd3Client, Cluster, Etcd3, Etcd3Error, Etcd3ClientError, RetryFailedError,\
InvalidAuthToken, Unavailable, Unknown, UnsupportedEtcdVersion, UserEmpty, base64_encode InvalidAuthToken, Unavailable, Unknown, UnsupportedEtcdVersion, UserEmpty, AuthFailed, base64_encode
from threading import Thread from threading import Thread
from . import SleepException, MockResponse from . import SleepException, MockResponse
@@ -90,7 +90,7 @@ class TestKVCache(BaseTestEtcd3):
class TestPatroniEtcd3Client(BaseTestEtcd3): class TestPatroniEtcd3Client(BaseTestEtcd3):
@patch('patroni.dcs.etcd3.Etcd3Client.authenticate', Mock(side_effect=Exception)) @patch('patroni.dcs.etcd3.Etcd3Client.authenticate', Mock(side_effect=AuthFailed))
def test__init__(self): def test__init__(self):
self.assertRaises(SystemExit, self.setUp) self.assertRaises(SystemExit, self.setUp)
+3
View File
@@ -501,6 +501,9 @@ class TestHa(PostgresInit):
self.assertEqual(self.ha.run_cycle(), 'lost leader lock during restart') self.assertEqual(self.ha.run_cycle(), 'lost leader lock during restart')
mock_terminate.assert_called() mock_terminate.assert_called()
self.ha.is_paused = true
self.assertEqual(self.ha.run_cycle(), 'PAUSE: restart in progress')
def test_manual_failover_from_leader(self): def test_manual_failover_from_leader(self):
self.ha.fetch_node_status = get_node_status() self.ha.fetch_node_status = get_node_status()
self.ha.has_lock = true self.ha.has_lock = true
+18 -10
View File
@@ -26,7 +26,7 @@ def mock_list_namespaced_config_map(*args, **kwargs):
return k8s_client.V1ConfigMapList(metadata=metadata, items=items, kind='ConfigMapList') return k8s_client.V1ConfigMapList(metadata=metadata, items=items, kind='ConfigMapList')
def mock_list_namespaced_endpoints(*args, **kwargs): def mock_read_namespaced_endpoints(*args, **kwargs):
target_ref = k8s_client.V1ObjectReference(kind='Pod', resource_version='10', name='p-0', target_ref = k8s_client.V1ObjectReference(kind='Pod', resource_version='10', name='p-0',
namespace='default', uid='964dfeae-e79b-4476-8a5a-1920b5c2a69d') namespace='default', uid='964dfeae-e79b-4476-8a5a-1920b5c2a69d')
address0 = k8s_client.V1EndpointAddress(ip='10.0.0.0', target_ref=target_ref) address0 = k8s_client.V1EndpointAddress(ip='10.0.0.0', target_ref=target_ref)
@@ -35,9 +35,12 @@ def mock_list_namespaced_endpoints(*args, **kwargs):
subset = k8s_client.V1EndpointSubset(addresses=[address1, address0], ports=[port]) subset = k8s_client.V1EndpointSubset(addresses=[address1, address0], ports=[port])
metadata = k8s_client.V1ObjectMeta(resource_version='1', labels={'f': 'b'}, name='test', metadata = k8s_client.V1ObjectMeta(resource_version='1', labels={'f': 'b'}, name='test',
annotations={'optime': '1234', 'leader': 'p-0', 'ttl': '30s'}) annotations={'optime': '1234', 'leader': 'p-0', 'ttl': '30s'})
endpoint = k8s_client.V1Endpoints(subsets=[subset], metadata=metadata) return k8s_client.V1Endpoints(subsets=[subset], metadata=metadata)
metadata = k8s_client.V1ObjectMeta(resource_version='1')
return k8s_client.V1EndpointsList(metadata=metadata, items=[endpoint], kind='V1EndpointsList')
def mock_list_namespaced_endpoints(*args, **kwargs):
return k8s_client.V1EndpointsList(metadata=k8s_client.V1ObjectMeta(resource_version='1'),
items=[mock_read_namespaced_endpoints()], kind='V1EndpointsList')
def mock_list_namespaced_pod(*args, **kwargs): def mock_list_namespaced_pod(*args, **kwargs):
@@ -261,20 +264,25 @@ class TestKubernetesEndpoints(BaseTestKubernetes):
def test_update_leader_with_restricted_access(self): def test_update_leader_with_restricted_access(self):
self.assertIsNotNone(self.k.update_leader('123', True)) self.assertIsNotNone(self.k.update_leader('123', True))
@patch.object(k8s_client.CoreV1Api, 'read_namespaced_endpoints', create=True)
@patch.object(k8s_client.CoreV1Api, 'patch_namespaced_endpoints', create=True) @patch.object(k8s_client.CoreV1Api, 'patch_namespaced_endpoints', create=True)
def test__update_leader_with_retry(self, mock_patch): def test__update_leader_with_retry(self, mock_patch, mock_read):
mock_read.return_value = mock_read_namespaced_endpoints()
mock_patch.side_effect = k8s_client.rest.ApiException(502, '') mock_patch.side_effect = k8s_client.rest.ApiException(502, '')
self.assertFalse(self.k.update_leader('123')) self.assertFalse(self.k.update_leader('123'))
mock_patch.side_effect = RetryFailedError('') mock_patch.side_effect = RetryFailedError('')
self.assertFalse(self.k.update_leader('123')) self.assertFalse(self.k.update_leader('123'))
mock_patch.side_effect = k8s_client.rest.ApiException(409, '') mock_patch.side_effect = k8s_client.rest.ApiException(409, '')
with patch('time.time', Mock(side_effect=[0, 100, 200])): with patch('time.time', Mock(side_effect=[0, 100, 200, 0, 0, 0, 0, 100, 200])):
self.assertFalse(self.k.update_leader('123')) self.assertFalse(self.k.update_leader('123'))
with patch('time.sleep', Mock()):
self.assertFalse(self.k.update_leader('123')) self.assertFalse(self.k.update_leader('123'))
mock_patch.side_effect = [k8s_client.rest.ApiException(409, ''), mock_namespaced_kind()] self.assertFalse(self.k.update_leader('123'))
self.k._kinds._object_cache['test'].metadata.resource_version = '2' mock_patch.side_effect = [k8s_client.rest.ApiException(409, ''), mock_namespaced_kind()]
self.assertIsNotNone(self.k._update_leader_with_retry({}, '1', [])) mock_read.return_value.metadata.resource_version = '2'
self.assertIsNotNone(self.k._update_leader_with_retry({}, '1', []))
mock_patch.side_effect = k8s_client.rest.ApiException(409, '')
mock_read.side_effect = Exception
self.assertFalse(self.k.update_leader('123'))
@patch.object(k8s_client.CoreV1Api, 'create_namespaced_endpoints', @patch.object(k8s_client.CoreV1Api, 'create_namespaced_endpoints',
Mock(side_effect=[k8s_client.rest.ApiException(500, ''), Mock(side_effect=[k8s_client.rest.ApiException(500, ''),
+13 -12
View File
@@ -306,7 +306,8 @@ class TestPostgresql(BaseTestPostgresql):
def test_sync_replication_slots(self): def test_sync_replication_slots(self):
self.p.start() self.p.start()
config = ClusterConfig(1, {'slots': {'test_3': {'database': 'a', 'plugin': 'b'}, config = ClusterConfig(1, {'slots': {'test_3': {'database': 'a', 'plugin': 'b'},
'A': 0, 'ls': 0, 'b': {'type': 'logical', 'plugin': '1'}}}, 1) 'A': 0, 'ls': 0, 'b': {'type': 'logical', 'plugin': '1'}},
'ignore_slots': [{'name': 'blabla'}]}, 1)
cluster = Cluster(True, config, self.leader, 0, [self.me, self.other, self.leadermem], None, None, None) cluster = Cluster(True, config, self.leader, 0, [self.me, self.other, self.leadermem], None, None, None)
with mock.patch('patroni.postgresql.Postgresql._query', Mock(side_effect=psycopg2.OperationalError)): with mock.patch('patroni.postgresql.Postgresql._query', Mock(side_effect=psycopg2.OperationalError)):
self.p.slots_handler.sync_replication_slots(cluster) self.p.slots_handler.sync_replication_slots(cluster)
@@ -614,32 +615,32 @@ class TestPostgresql(BaseTestPostgresql):
with patch.object(Postgresql, "query", side_effect=[ with patch.object(Postgresql, "query", side_effect=[
mock_cursor, mock_cursor,
[(self.leadermem.name, 'streaming', 'sync'), [(self.leadermem.name, 'sync', 1),
(self.me.name, 'streaming', 'async'), (self.me.name, 'async', 2),
(self.other.name, 'streaming', 'async')] (self.other.name, 'async', 2)]
]): ]):
self.assertEqual(self.p.pick_synchronous_standby(cluster), ([self.leadermem.name], [self.leadermem.name])) self.assertEqual(self.p.pick_synchronous_standby(cluster), ([self.leadermem.name], [self.leadermem.name]))
with patch.object(Postgresql, "query", side_effect=[ with patch.object(Postgresql, "query", side_effect=[
mock_cursor, mock_cursor,
[(self.leadermem.name, 'streaming', 'potential'), [(self.leadermem.name, 'potential', 1),
(self.me.name, 'streaming', 'async'), (self.me.name, 'async', 2),
(self.other.name, 'streaming', 'async')] (self.other.name, 'async', 2)]
]): ]):
self.assertEqual(self.p.pick_synchronous_standby(cluster), ([self.leadermem.name], [])) self.assertEqual(self.p.pick_synchronous_standby(cluster), ([self.leadermem.name], []))
with patch.object(Postgresql, "query", side_effect=[ with patch.object(Postgresql, "query", side_effect=[
mock_cursor, mock_cursor,
[(self.me.name, 'streaming', 'async'), [(self.me.name, 'async', 1),
(self.other.name, 'streaming', 'async')] (self.other.name, 'async', 2)]
]): ]):
self.assertEqual(self.p.pick_synchronous_standby(cluster), ([self.me.name], [])) self.assertEqual(self.p.pick_synchronous_standby(cluster), ([self.me.name], []))
with patch.object(Postgresql, "query", side_effect=[ with patch.object(Postgresql, "query", side_effect=[
mock_cursor, mock_cursor,
[('missing', 'streaming', 'sync'), [('missing', 'sync', 1),
(self.me.name, 'streaming', 'async'), (self.me.name, 'async', 2),
(self.other.name, 'streaming', 'async')] (self.other.name, 'async', 3)]
]): ]):
self.assertEqual(self.p.pick_synchronous_standby(cluster), ([self.me.name], [])) self.assertEqual(self.p.pick_synchronous_standby(cluster), ([self.me.name], []))
+13 -6
View File
@@ -1,5 +1,6 @@
import os import os
import unittest import unittest
import tempfile
import time import time
from mock import Mock, patch from mock import Mock, patch
@@ -86,8 +87,6 @@ class TestKVStoreTTL(unittest.TestCase):
self.so.set('foo', 'bar') self.so.set('foo', 'bar')
self.so.set('fooo', 'bar') self.so.set('fooo', 'bar')
self.assertFalse(self.so.delete('foo', prevValue='buz')) self.assertFalse(self.so.delete('foo', prevValue='buz'))
self.assertFalse(self.so.delete('foo', prevValue='bar', timeout=0.00001))
self.assertFalse(self.so.delete('foo', prevValue='bar'))
self.assertTrue(self.so.delete('foo', recursive=True)) self.assertTrue(self.so.delete('foo', recursive=True))
self.assertFalse(self.so.retry(self.so._delete, 'foo', prevValue='')) self.assertFalse(self.so.retry(self.so._delete, 'foo', prevValue=''))
@@ -99,10 +98,14 @@ class TestKVStoreTTL(unittest.TestCase):
@patch('time.sleep', Mock()) @patch('time.sleep', Mock())
def test_retry(self): def test_retry(self):
return_values = [FAIL_REASON.QUEUE_FULL, FAIL_REASON.SUCCESS, FAIL_REASON.REQUEST_DENIED] return_values = [FAIL_REASON.QUEUE_FULL] * 2 + [FAIL_REASON.SUCCESS, FAIL_REASON.REQUEST_DENIED]
def test(callback): def test(callback):
callback(True, return_values.pop(0)) callback(True, return_values.pop(0))
with patch('time.time', Mock(side_effect=[1, 100])):
self.assertFalse(self.so.retry(test))
self.assertTrue(self.so.retry(test)) self.assertTrue(self.so.retry(test))
self.assertFalse(self.so.retry(test)) self.assertFalse(self.so.retry(test))
@@ -119,8 +122,11 @@ class TestKVStoreTTL(unittest.TestCase):
class TestRaft(unittest.TestCase): class TestRaft(unittest.TestCase):
_TMP = tempfile.gettempdir()
def test_raft(self): def test_raft(self):
raft = Raft({'ttl': 30, 'scope': 'test', 'name': 'pg', 'self_addr': '127.0.0.1:1234', 'retry_timeout': 10}) raft = Raft({'ttl': 30, 'scope': 'test', 'name': 'pg', 'self_addr': '127.0.0.1:1234',
'retry_timeout': 10, 'data_dir': self._TMP})
raft.set_retry_timeout(20) raft.set_retry_timeout(20)
raft.set_ttl(60) raft.set_ttl(60)
self.assertTrue(raft.touch_member('')) self.assertTrue(raft.touch_member(''))
@@ -142,7 +148,7 @@ class TestRaft(unittest.TestCase):
raft._sync_obj._SyncObj__thread.join() raft._sync_obj._SyncObj__thread.join()
def tearDown(self): def tearDown(self):
remove_files('127.0.0.1:1234.') remove_files(os.path.join(self._TMP, '127.0.0.1:1234.'))
def setUp(self): def setUp(self):
self.tearDown() self.tearDown()
@@ -152,4 +158,5 @@ class TestRaft(unittest.TestCase):
def test_init(self, mock_event, mock_kvstore): def test_init(self, mock_event, mock_kvstore):
mock_kvstore.return_value.applied_local_log = False mock_kvstore.return_value.applied_local_log = False
mock_event.return_value.isSet.side_effect = [False, True] mock_event.return_value.isSet.side_effect = [False, True]
self.assertIsNotNone(Raft({'ttl': 30, 'scope': 'test', 'name': 'pg', 'patronictl': True, 'self_addr': '1'})) self.assertIsNotNone(Raft({'ttl': 30, 'scope': 'test', 'name': 'pg', 'patronictl': True,
'self_addr': '1', 'data_dir': self._TMP}))
+10 -4
View File
@@ -40,6 +40,12 @@ def mock_cancellable_call1(*args, **kwargs):
return 1 return 1
def mock_single_user_mode(self, communicate, options):
communicate['stdout'] = b'foo'
communicate['stderr'] = b'bar'
return 1
@patch('subprocess.call', Mock(return_value=0)) @patch('subprocess.call', Mock(return_value=0))
@patch('psycopg2.connect', psycopg2_connect) @patch('psycopg2.connect', psycopg2_connect)
class TestRewind(BaseTestPostgresql): class TestRewind(BaseTestPostgresql):
@@ -170,8 +176,8 @@ class TestRewind(BaseTestPostgresql):
@patch.object(MockCursor, 'fetchone', Mock(side_effect=[(True,), Exception])) @patch.object(MockCursor, 'fetchone', Mock(side_effect=[(True,), Exception]))
def test_check_leader_is_not_in_recovery(self): def test_check_leader_is_not_in_recovery(self):
self.r.check_leader_is_not_in_recovery() self.r.check_leader_is_not_in_recovery({})
self.r.check_leader_is_not_in_recovery() self.r.check_leader_is_not_in_recovery({})
def test_read_postmaster_opts(self): def test_read_postmaster_opts(self):
m = mock_open(read_data='/usr/lib/postgres/9.6/bin/postgres "-D" "data/postgresql0" \ m = mock_open(read_data='/usr/lib/postgres/9.6/bin/postgres "-D" "data/postgresql0" \
@@ -206,9 +212,9 @@ class TestRewind(BaseTestPostgresql):
@patch('os.listdir', Mock(return_value=[])) @patch('os.listdir', Mock(return_value=[]))
@patch('os.path.isfile', Mock(return_value=True)) @patch('os.path.isfile', Mock(return_value=True))
@patch.object(Rewind, 'read_postmaster_opts', Mock(return_value={})) @patch.object(Rewind, 'read_postmaster_opts', Mock(return_value={}))
@patch.object(Rewind, 'single_user_mode', Mock(return_value=0)) @patch.object(Rewind, 'single_user_mode', mock_single_user_mode)
def test_ensure_clean_shutdown(self): def test_ensure_clean_shutdown(self):
self.assertTrue(self.r.ensure_clean_shutdown()) self.assertIsNone(self.r.ensure_clean_shutdown())
@patch('patroni.postgresql.rewind.Thread', MockThread) @patch('patroni.postgresql.rewind.Thread', MockThread)
@patch.object(Postgresql, 'controldata') @patch.object(Postgresql, 'controldata')
+10 -1
View File
@@ -2,7 +2,7 @@ import unittest
from mock import Mock, patch from mock import Mock, patch
from patroni.exceptions import PatroniException from patroni.exceptions import PatroniException
from patroni.utils import Retry, RetryFailedError, enable_keepalive, polling_loop, validate_directory from patroni.utils import Retry, RetryFailedError, enable_keepalive, find_executable, polling_loop, validate_directory
class TestUtils(unittest.TestCase): class TestUtils(unittest.TestCase):
@@ -41,6 +41,15 @@ class TestUtils(unittest.TestCase):
with patch('sys.platform', platform): with patch('sys.platform', platform):
self.assertIsNone(enable_keepalive(Mock(), 10, 5)) self.assertIsNone(enable_keepalive(Mock(), 10, 5))
@patch('sys.platform', 'win32')
def test_find_executable(self):
with patch('os.path.isfile', Mock(return_value=True)):
self.assertEqual(find_executable('vim'), 'vim.exe')
with patch('os.path.isfile', Mock(return_value=False)):
self.assertIsNone(find_executable('vim'))
with patch('os.path.isfile', Mock(side_effect=[False, True])):
self.assertEqual(find_executable('vim', '/'), '/vim.exe')
@patch('time.sleep', Mock()) @patch('time.sleep', Mock())
class TestRetrySleeper(unittest.TestCase): class TestRetrySleeper(unittest.TestCase):
+21 -9
View File
@@ -33,11 +33,21 @@ config = {
"etcd": { "etcd": {
"hosts": "127.0.0.1:2379,127.0.0.1:2380" "hosts": "127.0.0.1:2379,127.0.0.1:2380"
}, },
"etcd3": {
"url": "https://127.0.0.1:2379"
},
"exhibitor": { "exhibitor": {
"hosts": ["string"], "hosts": ["string"],
"port": 4000, "port": 4000,
"pool_interval": 1000 "pool_interval": 1000
}, },
"raft": {
"self_addr": "127.0.0.1:2222",
"bind_addr": "0.0.0.0:2222",
"partner_addrs": ["127.0.0.1:2223", "127.0.0.1:2224"],
"data_dir": "/",
"password": "12345"
},
"zookeeper": { "zookeeper": {
"hosts": "127.0.0.1:3379,127.0.0.1:3380" "hosts": "127.0.0.1:3379,127.0.0.1:3380"
}, },
@@ -139,7 +149,7 @@ class TestValidator(unittest.TestCase):
def test_complete_config(self, mock_out, mock_err): def test_complete_config(self, mock_out, mock_err):
schema(config) schema(config)
output = mock_out.getvalue() output = mock_out.getvalue()
self.assertEqual(['postgresql.bin_dir'], parse_output(output)) self.assertEqual(['postgresql.bin_dir', 'raft.bind_addr', 'raft.self_addr'], parse_output(output))
def test_bin_dir_is_file(self, mock_out, mock_err): def test_bin_dir_is_file(self, mock_out, mock_err):
files.append(config["postgresql"]["data_dir"]) files.append(config["postgresql"]["data_dir"])
@@ -151,7 +161,8 @@ class TestValidator(unittest.TestCase):
schema(c) schema(c)
output = mock_out.getvalue() output = mock_out.getvalue()
self.assertEqual(['etcd.hosts.1', 'etcd.hosts.2', 'kubernetes.pod_ip', 'postgresql.bin_dir', self.assertEqual(['etcd.hosts.1', 'etcd.hosts.2', 'kubernetes.pod_ip', 'postgresql.bin_dir',
'postgresql.data_dir', 'restapi.connect_address'], parse_output(output)) 'postgresql.data_dir', 'raft.bind_addr', 'raft.self_addr',
'restapi.connect_address'], parse_output(output))
@patch('socket.inet_pton', Mock(), create=True) @patch('socket.inet_pton', Mock(), create=True)
def test_bin_dir_is_empty(self, mock_out, mock_err): def test_bin_dir_is_empty(self, mock_out, mock_err):
@@ -167,8 +178,8 @@ class TestValidator(unittest.TestCase):
with patch('patroni.validator.open', mock_open(read_data='9')): with patch('patroni.validator.open', mock_open(read_data='9')):
schema(c) schema(c)
output = mock_out.getvalue() output = mock_out.getvalue()
self.assertEqual(['consul.host', 'etcd.host', 'postgresql.bin_dir', 'postgresql.data_dir', self.assertEqual(['consul.host', 'etcd.host', 'postgresql.bin_dir', 'postgresql.data_dir', 'postgresql.listen',
'postgresql.listen', 'restapi.connect_address'], parse_output(output)) 'raft.bind_addr', 'raft.self_addr', 'restapi.connect_address'], parse_output(output))
@patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1")) @patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1"))
def test_data_dir_contains_pg_version(self, mock_out, mock_err): def test_data_dir_contains_pg_version(self, mock_out, mock_err):
@@ -186,7 +197,7 @@ class TestValidator(unittest.TestCase):
with patch('patroni.validator.open', mock_open(read_data='12')): with patch('patroni.validator.open', mock_open(read_data='12')):
schema(config) schema(config)
output = mock_out.getvalue() output = mock_out.getvalue()
self.assertEqual([], parse_output(output)) self.assertEqual(['raft.bind_addr', 'raft.self_addr'], parse_output(output))
@patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1")) @patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1"))
def test_pg_version_missmatch(self, mock_out, mock_err): def test_pg_version_missmatch(self, mock_out, mock_err):
@@ -201,7 +212,8 @@ class TestValidator(unittest.TestCase):
with patch('patroni.validator.open', mock_open(read_data='11')): with patch('patroni.validator.open', mock_open(read_data='11')):
schema(c) schema(c)
output = mock_out.getvalue() output = mock_out.getvalue()
self.assertEqual(['etcd.hosts', 'postgresql.data_dir'], parse_output(output)) self.assertEqual(['etcd.hosts', 'postgresql.data_dir',
'raft.bind_addr', 'raft.self_addr'], parse_output(output))
@patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1")) @patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1"))
def test_pg_wal_doesnt_exist(self, mock_out, mock_err): def test_pg_wal_doesnt_exist(self, mock_out, mock_err):
@@ -214,7 +226,7 @@ class TestValidator(unittest.TestCase):
with patch('patroni.validator.open', mock_open(read_data='11')): with patch('patroni.validator.open', mock_open(read_data='11')):
schema(c) schema(c)
output = mock_out.getvalue() output = mock_out.getvalue()
self.assertEqual(['postgresql.data_dir'], parse_output(output)) self.assertEqual(['postgresql.data_dir', 'raft.bind_addr', 'raft.self_addr'], parse_output(output))
def test_data_dir_is_empty_string(self, mock_out, mock_err): def test_data_dir_is_empty_string(self, mock_out, mock_err):
directories.append(config["postgresql"]["data_dir"]) directories.append(config["postgresql"]["data_dir"])
@@ -226,5 +238,5 @@ class TestValidator(unittest.TestCase):
c["postgresql"]["bin_dir"] = "" c["postgresql"]["bin_dir"] = ""
schema(c) schema(c)
output = mock_out.getvalue() output = mock_out.getvalue()
self.assertEqual(['kubernetes', 'postgresql.bin_dir', self.assertEqual(['kubernetes', 'postgresql.bin_dir', 'postgresql.data_dir',
'postgresql.data_dir', 'postgresql.pg_hba'], parse_output(output)) 'postgresql.pg_hba', 'raft.bind_addr', 'raft.self_addr'], parse_output(output))
+1 -1
View File
@@ -17,7 +17,7 @@ class MockKazooClient(Mock):
def __init__(self, *args, **kwargs): def __init__(self, *args, **kwargs):
super(MockKazooClient, self).__init__() super(MockKazooClient, self).__init__()
self._session_timeout = 30 self._session_timeout = 30000
@property @property
def client_id(self): def client_id(self):