mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-28 16:39:32 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
491f230711 | ||
|
|
f1d7ccf36e | ||
|
|
96ea01bee4 | ||
|
|
e684ca66e5 | ||
|
|
1a0876e5ca | ||
|
|
9bf074acfb |
@@ -3,6 +3,23 @@
|
||||
Release notes
|
||||
=============
|
||||
|
||||
Version 1.5.3
|
||||
-------------
|
||||
|
||||
Compatibility and bugfix release.
|
||||
|
||||
- Improve stability when running with python3 against zookeeper (Alexander Kukushkin)
|
||||
|
||||
Change of `loop_wait` was causing Patroni to disconnect from zookeeper and never reconnect back.
|
||||
|
||||
- Fix broken compatibility with postgres 9.3 (Alexander)
|
||||
|
||||
When opening a replication connection we should specify replication=1, beacuse 9.3 does not understand replication='database'
|
||||
|
||||
- Make sure we refresh Consul session at least once per HA loop and improve handling of consul sessions exceptions (Alexander)
|
||||
|
||||
Restart of local consul agent invalidates all sessions related to the node. Not calling session refresh on time and not doing proper handling of session errors was causing demote of the primary.
|
||||
|
||||
Version 1.5.2
|
||||
-------------
|
||||
|
||||
|
||||
+15
-9
@@ -117,12 +117,24 @@ class PatroniController(AbstractController):
|
||||
except IOError:
|
||||
return None
|
||||
|
||||
def add_tag_to_config(self, tag, value):
|
||||
@staticmethod
|
||||
def recursive_update(dst, src):
|
||||
for k, v in src.items():
|
||||
if k in dst and isinstance(dst[k], dict):
|
||||
PatroniController.recursive_update(dst[k], v)
|
||||
else:
|
||||
dst[k] = v
|
||||
|
||||
def update_config(self, custom_config):
|
||||
with open(self._config) as r:
|
||||
config = yaml.safe_load(r)
|
||||
config['tags']['tag'] = value
|
||||
self.recursive_update(config, custom_config)
|
||||
with open(self._config, 'w') as w:
|
||||
yaml.safe_dump(config, w, default_flow_style=False)
|
||||
self._scope = config.get('scope', 'batman')
|
||||
|
||||
def add_tag_to_config(self, tag, value):
|
||||
self.update_config({'tags': {tag: value}})
|
||||
|
||||
def _start(self):
|
||||
if self.watchdog:
|
||||
@@ -174,13 +186,7 @@ class PatroniController(AbstractController):
|
||||
config['bootstrap']['initdb'].extend([{'auth': 'md5'}, {'auth-host': 'md5'}])
|
||||
|
||||
if custom_config is not None:
|
||||
def recursive_update(dst, src):
|
||||
for k, v in src.items():
|
||||
if k in dst and isinstance(dst[k], dict):
|
||||
recursive_update(dst[k], v)
|
||||
else:
|
||||
dst[k] = v
|
||||
recursive_update(config, custom_config)
|
||||
self.recursive_update(config, custom_config)
|
||||
|
||||
if config['postgresql'].get('callbacks', {}).get('on_role_change'):
|
||||
config['postgresql']['callbacks']['on_role_change'] += ' ' + str(self.__PORT)
|
||||
|
||||
@@ -43,7 +43,7 @@ Scenario: check dynamic configuration change via DCS
|
||||
And I receive a response loop_wait 2
|
||||
When I issue a GET request to http://127.0.0.1:8008/patroni
|
||||
Then I receive a response code 200
|
||||
And I receive a response tags {'tag': 'new_value'}
|
||||
And I receive a response tags {'new_tag': 'new_value'}
|
||||
|
||||
Scenario: check API requests for the primary-replica pair in the pause mode
|
||||
Given I run patronictl.py pause batman
|
||||
|
||||
@@ -3,13 +3,15 @@ Feature: standby cluster
|
||||
Given I start postgres1
|
||||
Then postgres1 is a leader after 10 seconds
|
||||
And I sleep for 2 seconds
|
||||
When I issue a PATCH request to http://127.0.0.1:8009/config with {"slots": {"pm_1": {"type": "physical"}}, "postgresql": {"parameters": {"wal_level": "logical"}}}
|
||||
When I issue a PATCH request to http://127.0.0.1:8009/config with {"loop_wait": 2, "slots": {"pm_1": {"type": "physical"}}, "postgresql": {"parameters": {"wal_level": "logical"}}}
|
||||
Then I receive a response code 200
|
||||
And Response on GET http://127.0.0.1:8009/config contains slots after 10 seconds
|
||||
And I sleep for 2 seconds
|
||||
When I issue a PATCH request to http://127.0.0.1:8009/config with {"slots": {"test_logical": {"type": "logical", "database": "postgres", "plugin": "test_decoding"}}}
|
||||
Then I receive a response code 200
|
||||
When I start postgres0 with callback configured
|
||||
Then "members/postgres0" key in DCS has state=running after 10 seconds
|
||||
And replication works from postgres1 to postgres0 after 15 seconds
|
||||
When I shut down postgres1
|
||||
Then postgres0 is a leader after 10 seconds
|
||||
And I sleep for 2 seconds
|
||||
@@ -20,8 +22,7 @@ Feature: standby cluster
|
||||
Scenario: check replication of a single table in a standby cluster
|
||||
Given I start postgres1 in a standby cluster batman1 as a clone of postgres0
|
||||
Then postgres1 is a leader of batman1 after 10 seconds
|
||||
When I issue a PATCH request to http://127.0.0.1:8009/config with {"ttl": 20, "loop_wait": 2}
|
||||
And I add the table foo to postgres0
|
||||
When I add the table foo to postgres0
|
||||
Then table foo is present on postgres1 after 20 seconds
|
||||
When I start postgres2 in a cluster batman1
|
||||
Then postgres2 role is the replica after 24 seconds
|
||||
|
||||
@@ -30,12 +30,10 @@ def start_patroni(context, name, cluster_name):
|
||||
|
||||
@step('I start {name:w} in a standby cluster {cluster_name:w} as a clone of {name2:w}')
|
||||
def start_patroni_stanby_cluster(context, name, cluster_name, name2):
|
||||
ctl = context.pctl._processes.pop(name, None)
|
||||
# we need to remove patroni.dynamic.json in order to "bootstrap" standby cluster with existing PGDATA
|
||||
if ctl:
|
||||
os.unlink(os.path.join(ctl._data_dir, 'patroni.dynamic.json'))
|
||||
os.unlink(os.path.join(context.pctl._processes[name]._data_dir, 'patroni.dynamic.json'))
|
||||
port = context.pctl._processes[name2]._connkwargs.get('port')
|
||||
return context.pctl.start(name, custom_config={
|
||||
context.pctl._processes[name].update_config({
|
||||
"scope": cluster_name,
|
||||
"bootstrap": {
|
||||
"dcs": {
|
||||
@@ -47,6 +45,7 @@ def start_patroni_stanby_cluster(context, name, cluster_name, name2):
|
||||
}
|
||||
}
|
||||
})
|
||||
return context.pctl.start(name)
|
||||
|
||||
|
||||
@step('{pg_name1:w} is replicating from {pg_name2:w} after {timeout:d} seconds')
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
FROM postgres:10
|
||||
FROM postgres:11
|
||||
MAINTAINER Alexander Kukushkin <[email protected]>
|
||||
|
||||
RUN export DEBIAN_FRONTEND=noninteractive \
|
||||
@@ -7,23 +7,18 @@ RUN export DEBIAN_FRONTEND=noninteractive \
|
||||
&& apt-get upgrade -y \
|
||||
&& apt-cache depends patroni | sed -n -e 's/.* Depends: \(python3-.\+\)$/\1/p' \
|
||||
| grep -Ev '^python3-(sphinx|etcd|consul|kazoo|kubernetes)' \
|
||||
| xargs apt-get install -y curl jq locales git python3-pip python3-wheel \
|
||||
|
||||
| xargs apt-get install -y vim-tiny curl jq locales git python3-pip python3-wheel \
|
||||
## Make sure we have a en_US.UTF-8 locale available
|
||||
&& localedef -i en_US -c -f UTF-8 -A /usr/share/locale/locale.alias en_US.UTF-8 \
|
||||
|
||||
&& pip3 install setuptools \
|
||||
&& pip3 install 'git+https://github.com/zalando/patroni.git#egg=patroni[kubernetes]' \
|
||||
|
||||
&& PGHOME=/home/postgres \
|
||||
&& mkdir -p $PGHOME \
|
||||
&& chown postgres $PGHOME \
|
||||
&& sed -i "s|/var/lib/postgresql.*|$PGHOME:/bin/bash|" /etc/passwd \
|
||||
|
||||
# Set permissions for OpenShift
|
||||
&& chmod 775 $PGHOME \
|
||||
&& chmod 664 /etc/passwd \
|
||||
|
||||
# Clean up
|
||||
&& apt-get remove -y git python3-pip python3-wheel \
|
||||
&& apt-get autoremove -y \
|
||||
@@ -33,7 +28,7 @@ RUN export DEBIAN_FRONTEND=noninteractive \
|
||||
ADD entrypoint.sh /
|
||||
|
||||
EXPOSE 5432 8008
|
||||
ENV LC_ALL=en_US.UTF-8 LANG=en_US.UTF-8
|
||||
ENV LC_ALL=en_US.UTF-8 LANG=en_US.UTF-8 EDITOR=/usr/bin/editor
|
||||
USER postgres
|
||||
WORKDIR /home/postgres
|
||||
CMD ["/bin/bash", "/entrypoint.sh"]
|
||||
CMD ["/bin/bash", "/entrypoint.sh"]
|
||||
|
||||
@@ -20,11 +20,11 @@ bootstrap:
|
||||
- data-checksums
|
||||
pg_hba:
|
||||
- host all all 0.0.0.0/0 md5
|
||||
- host replication ${PATRONI_REPLICATION_USERNAME} ${POD_IP}/16 md5
|
||||
- host replication ${PATRONI_REPLICATION_USERNAME} ${PATRONI_KUBERNETES_POD_IP}/16 md5
|
||||
restapi:
|
||||
connect_address: '${POD_IP}:8008'
|
||||
connect_address: '${PATRONI_KUBERNETES_POD_IP}:8008'
|
||||
postgresql:
|
||||
connect_address: '${POD_IP}:5432'
|
||||
connect_address: '${PATRONI_KUBERNETES_POD_IP}:5432'
|
||||
authentication:
|
||||
superuser:
|
||||
password: '${PATRONI_SUPERUSER_PASSWORD}'
|
||||
|
||||
@@ -108,7 +108,7 @@ objects:
|
||||
spec:
|
||||
containers:
|
||||
- env:
|
||||
- name: POD_IP
|
||||
- name: PATRONI_KUBERNETES_POD_IP
|
||||
valueFrom:
|
||||
fieldRef:
|
||||
apiVersion: v1
|
||||
|
||||
@@ -108,7 +108,7 @@ objects:
|
||||
spec:
|
||||
containers:
|
||||
- env:
|
||||
- name: POD_IP
|
||||
- name: PATRONI_KUBERNETES_POD_IP
|
||||
valueFrom:
|
||||
fieldRef:
|
||||
apiVersion: v1
|
||||
|
||||
@@ -28,7 +28,7 @@ spec:
|
||||
- mountPath: /home/postgres/pgdata
|
||||
name: pgdata
|
||||
env:
|
||||
- name: POD_IP
|
||||
- name: PATRONI_KUBERNETES_POD_IP
|
||||
valueFrom:
|
||||
fieldRef:
|
||||
fieldPath: status.podIP
|
||||
|
||||
+21
-4
@@ -27,10 +27,14 @@ class ConsulInternalError(ConsulException):
|
||||
"""An internal Consul server error occurred"""
|
||||
|
||||
|
||||
class InvalidSessionTTL(ConsulInternalError):
|
||||
class InvalidSessionTTL(ConsulException):
|
||||
"""Session TTL is too small or too big"""
|
||||
|
||||
|
||||
class InvalidSession(ConsulException):
|
||||
"""invalid session"""
|
||||
|
||||
|
||||
class HTTPClient(object):
|
||||
|
||||
def __init__(self, host='127.0.0.1', port=8500, token=None, scheme='http', verify=True, cert=None, ca_cert=None):
|
||||
@@ -72,6 +76,8 @@ class HTTPClient(object):
|
||||
msg = '{0} {1}'.format(response.status, data)
|
||||
if data.startswith('Invalid Session TTL'):
|
||||
raise InvalidSessionTTL(msg)
|
||||
elif data.startswith('invalid session'):
|
||||
raise InvalidSession(msg)
|
||||
else:
|
||||
raise ConsulInternalError(msg)
|
||||
return base.Response(response.status, response.headers, data)
|
||||
@@ -357,6 +363,9 @@ class Consul(AbstractDCS):
|
||||
if self._register_service:
|
||||
self.update_service(not create_member and member and member.data or {}, data)
|
||||
return True
|
||||
except InvalidSession:
|
||||
self._session = None
|
||||
logger.error('Our session disappeared from Consul, can not "touch_member"')
|
||||
except Exception:
|
||||
logger.exception('touch_member')
|
||||
return False
|
||||
@@ -414,14 +423,21 @@ class Consul(AbstractDCS):
|
||||
return self._update_service(new_data)
|
||||
|
||||
@catch_consul_errors
|
||||
def _do_attempt_to_acquire_leader(self, kwargs):
|
||||
return self.retry(self._client.kv.put, self.leader_path, self._name, **kwargs)
|
||||
def _do_attempt_to_acquire_leader(self, permanent):
|
||||
try:
|
||||
kwargs = {} if permanent else {'acquire': self._session}
|
||||
return self.retry(self._client.kv.put, self.leader_path, self._name, **kwargs)
|
||||
except InvalidSession:
|
||||
self._session = None
|
||||
logger.error('Our session disappeared from Consul. Will try to get a new one and retry attempt')
|
||||
self.refresh_session()
|
||||
return self.retry(self._client.kv.put, self.leader_path, self._name, acquire=self._session)
|
||||
|
||||
def attempt_to_acquire_leader(self, permanent=False):
|
||||
if not self._session and not permanent:
|
||||
self.refresh_session()
|
||||
|
||||
ret = self._do_attempt_to_acquire_leader({} if permanent else {'acquire': self._session})
|
||||
ret = self._do_attempt_to_acquire_leader(permanent)
|
||||
if not ret:
|
||||
logger.info('Could not take out TTL lock')
|
||||
|
||||
@@ -499,4 +515,5 @@ class Consul(AbstractDCS):
|
||||
try:
|
||||
return super(Consul, self).watch(None, timeout)
|
||||
finally:
|
||||
self._last_session_refresh = 0
|
||||
self.event.clear()
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import json
|
||||
import logging
|
||||
import select
|
||||
import time
|
||||
|
||||
from kazoo.client import KazooClient, KazooState, KazooRetry
|
||||
@@ -45,6 +46,13 @@ class PatroniSequentialThreadingHandler(SequentialThreadingHandler):
|
||||
args[1] = max(self._connect_timeout, args[1]/10.0)
|
||||
return super(PatroniSequentialThreadingHandler, self).create_connection(*args, **kwargs)
|
||||
|
||||
def select(self, *args, **kwargs):
|
||||
"""Python3 raises `ValueError` if socket is closed, because fd == -1"""
|
||||
try:
|
||||
return super(PatroniSequentialThreadingHandler, self).select(*args, **kwargs)
|
||||
except ValueError as e:
|
||||
raise select.error(9, str(e))
|
||||
|
||||
|
||||
class ZooKeeper(AbstractDCS):
|
||||
|
||||
|
||||
@@ -1243,7 +1243,8 @@ class Postgresql(object):
|
||||
@contextmanager
|
||||
def _get_replication_connection_cursor(self, host='localhost', port=5432, database=None, **kwargs):
|
||||
database = database or self._database
|
||||
with self._get_connection_cursor(host=host, port=int(port), database=database, replication='database',
|
||||
replication = 'database' if self._major_version >= 90400 else 1
|
||||
with self._get_connection_cursor(host=host, port=int(port), database=database, replication=replication,
|
||||
user=self._replication['username'], password=self._replication['password'],
|
||||
connect_timeout=3, options='-c statement_timeout=2000') as cur:
|
||||
yield cur
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
__version__ = '1.5.2'
|
||||
__version__ = '1.5.3'
|
||||
|
||||
+10
-9
@@ -4,7 +4,7 @@ import unittest
|
||||
from consul import ConsulException, NotFound
|
||||
from mock import Mock, patch
|
||||
from patroni.dcs.consul import AbstractDCS, Cluster, Consul, ConsulInternalError, \
|
||||
ConsulError, HTTPClient, InvalidSessionTTL
|
||||
ConsulError, HTTPClient, InvalidSessionTTL, InvalidSession
|
||||
from test_etcd import SleepException
|
||||
|
||||
|
||||
@@ -52,6 +52,8 @@ class TestHTTPClient(unittest.TestCase):
|
||||
self.assertRaises(ConsulInternalError, self.client.get, Mock(), '')
|
||||
self.client.http.request.return_value.data = b"Invalid Session TTL '3000000000', must be between [10s=24h0m0s]"
|
||||
self.assertRaises(InvalidSessionTTL, self.client.get, Mock(), '')
|
||||
self.client.http.request.return_value.data = b"invalid session '16492f43-c2d6-5307-432f-e32d6f7bcbd0'"
|
||||
self.assertRaises(InvalidSession, self.client.get, Mock(), '')
|
||||
|
||||
def test_unknown_method(self):
|
||||
try:
|
||||
@@ -110,19 +112,18 @@ class TestConsul(unittest.TestCase):
|
||||
self.c._session = 'fd4f44fe-2cac-bba5-a60b-304b51ff39b8'
|
||||
self.assertIsInstance(self.c.get_cluster(), Cluster)
|
||||
|
||||
@patch.object(consul.Consul.KV, 'delete', Mock(side_effect=[ConsulException, True, True]))
|
||||
@patch.object(consul.Consul.KV, 'put', Mock(side_effect=[True, ConsulException]))
|
||||
@patch.object(consul.Consul.KV, 'delete', Mock(side_effect=[ConsulException, True, True, True]))
|
||||
@patch.object(consul.Consul.KV, 'put', Mock(side_effect=[True, ConsulException, InvalidSession]))
|
||||
def test_touch_member(self):
|
||||
self.c._register_service = True
|
||||
self.c.refresh_session = Mock(return_value=True)
|
||||
self.c.touch_member({'balbla': 'blabla'})
|
||||
self.c.touch_member({'balbla': 'blabla'})
|
||||
self.c.touch_member({'balbla': 'blabla'})
|
||||
self.c.refresh_session = Mock(return_value=False)
|
||||
self.c.touch_member({'conn_url': 'postgres://replicator:[email protected]:5433/postgres',
|
||||
'api_url': 'http://127.0.0.1:8009/patroni'})
|
||||
self.c._register_service = True
|
||||
self.c.refresh_session = Mock(return_value=True)
|
||||
for _ in range(0, 4):
|
||||
self.c.touch_member({'balbla': 'blabla'})
|
||||
|
||||
@patch.object(consul.Consul.KV, 'put', Mock(return_value=False))
|
||||
@patch.object(consul.Consul.KV, 'put', Mock(side_effect=InvalidSession))
|
||||
def test_take_leader(self):
|
||||
self.c.set_ttl(20)
|
||||
self.c.refresh_session = Mock()
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import select
|
||||
import six
|
||||
import unittest
|
||||
|
||||
@@ -115,6 +116,10 @@ class TestPatroniSequentialThreadingHandler(unittest.TestCase):
|
||||
self.assertIsNotNone(self.handler.create_connection((), 40))
|
||||
self.assertIsNotNone(self.handler.create_connection(timeout=40))
|
||||
|
||||
@patch.object(SequentialThreadingHandler, 'select', Mock(side_effect=ValueError))
|
||||
def test_select(self):
|
||||
self.assertRaises(select.error, self.handler.select)
|
||||
|
||||
|
||||
class TestZooKeeper(unittest.TestCase):
|
||||
|
||||
|
||||
Reference in New Issue
Block a user