From f1d7ccf36efcb4ae13360400c591d82520220d2b Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 3 Dec 2018 16:35:14 +0100 Subject: [PATCH] Make sure we refresh session at least once per HA loop (#880) Fixes https://github.com/zalando/patroni/issues/879 --- patroni/dcs/consul.py | 25 +++++++++++++++++++++---- tests/test_consul.py | 19 ++++++++++--------- 2 files changed, 31 insertions(+), 13 deletions(-) diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index be370c41..1547710b 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -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() diff --git a/tests/test_consul.py b/tests/test_consul.py index 7f752f18..a9d1cc93 100644 --- a/tests/test_consul.py +++ b/tests/test_consul.py @@ -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:rep-pass@127.0.0.1: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()