diff --git a/patroni/dcs/zookeeper.py b/patroni/dcs/zookeeper.py index ac0167c9..f3b9ac38 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -4,8 +4,9 @@ import select import time from kazoo.client import KazooClient, KazooState, KazooRetry -from kazoo.exceptions import NoNodeError, NodeExistsError +from kazoo.exceptions import NoNodeError, NodeExistsError, SessionExpiredError from kazoo.handlers.threading import SequentialThreadingHandler +from kazoo.protocol.states import KeeperState from . import AbstractDCS, ClusterConfig, Cluster, Failover, Leader, Member, SyncState, TimelineHistory from ..exceptions import DCSError @@ -55,6 +56,20 @@ class PatroniSequentialThreadingHandler(SequentialThreadingHandler): raise select.error(9, str(e)) +class PatroniKazooClient(KazooClient): + + def _call(self, request, async_object): + # Before kazoo==2.7.0 it wasn't possible to send requests to zookeeper if + # the connection is in the SUSPENDED state and Patroni was strongly relying on it. + # The https://github.com/python-zk/kazoo/pull/588 changed it, and now such requests are queued. + # We override the `_call()` method in order to keep the old behavior. + + if self._state == KeeperState.CONNECTING: + async_object.set_exception(SessionExpiredError()) + return False + return super(PatroniKazooClient, self)._call(request, async_object) + + class ZooKeeper(AbstractDCS): def __init__(self, config): @@ -68,10 +83,10 @@ class ZooKeeper(AbstractDCS): 'cert': 'certfile', 'key': 'keyfile', 'key_password': 'keyfile_password'} kwargs = {v: config[k] for k, v in mapping.items() if k in config} - self._client = KazooClient(hosts, handler=PatroniSequentialThreadingHandler(config['retry_timeout']), - timeout=config['ttl'], connection_retry=KazooRetry(max_delay=1, max_tries=-1, - sleep_func=time.sleep), command_retry=KazooRetry(deadline=config['retry_timeout'], - max_delay=1, max_tries=-1, sleep_func=time.sleep), **kwargs) + self._client = PatroniKazooClient(hosts, handler=PatroniSequentialThreadingHandler(config['retry_timeout']), + timeout=config['ttl'], connection_retry=KazooRetry(max_delay=1, max_tries=-1, + sleep_func=time.sleep), command_retry=KazooRetry(max_delay=1, max_tries=-1, + deadline=config['retry_timeout'], sleep_func=time.sleep), **kwargs) self._client.add_listener(self.session_listener) self._fetch_cluster = True diff --git a/patroni/ha.py b/patroni/ha.py index 25935564..7d08dcae 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -1308,8 +1308,12 @@ class Ha(object): def _run_cycle(self): dcs_failed = False try: - self.load_cluster_from_dcs() - self.state_handler.reset_cluster_info_state(self.cluster, self.patroni.nofailover) + try: + self.load_cluster_from_dcs() + self.state_handler.reset_cluster_info_state(self.cluster, self.patroni.nofailover) + except Exception: + self.state_handler.reset_cluster_info_state(None, self.patroni.nofailover) + raise if self.is_paused(): self.watchdog.disable() diff --git a/tests/test_exhibitor.py b/tests/test_exhibitor.py index 9d4d460d..fdd3e993 100644 --- a/tests/test_exhibitor.py +++ b/tests/test_exhibitor.py @@ -24,7 +24,7 @@ class TestExhibitor(unittest.TestCase): @patch('urllib3.PoolManager.request', Mock(return_value=urllib3.HTTPResponse( status=200, body=b'{"servers":["127.0.0.1","127.0.0.2","127.0.0.3"],"port":2181}'))) - @patch('patroni.dcs.zookeeper.KazooClient', MockKazooClient) + @patch('patroni.dcs.zookeeper.PatroniKazooClient', MockKazooClient) def setUp(self): self.e = Exhibitor({'hosts': ['localhost', 'exhibitor'], 'port': 8181, 'scope': 'test', 'name': 'foo', 'ttl': 30, 'retry_timeout': 10}) diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py index ef368826..ed718a21 100644 --- a/tests/test_zookeeper.py +++ b/tests/test_zookeeper.py @@ -2,12 +2,13 @@ import select import six import unittest -from kazoo.client import KazooState +from kazoo.client import KazooClient, KazooState from kazoo.exceptions import NoNodeError, NodeExistsError from kazoo.handlers.threading import SequentialThreadingHandler -from kazoo.protocol.states import ZnodeStat +from kazoo.protocol.states import KeeperState, ZnodeStat from mock import Mock, patch -from patroni.dcs.zookeeper import Leader, PatroniSequentialThreadingHandler, ZooKeeper, ZooKeeperError +from patroni.dcs.zookeeper import Leader, PatroniKazooClient,\ + PatroniSequentialThreadingHandler, ZooKeeper, ZooKeeperError class MockKazooClient(Mock): @@ -128,9 +129,19 @@ class TestPatroniSequentialThreadingHandler(unittest.TestCase): self.assertRaises(select.error, self.handler.select) +class TestPatroniKazooClient(unittest.TestCase): + + def test__call(self): + c = PatroniKazooClient() + with patch.object(KazooClient, '_call', Mock()): + self.assertIsNotNone(c._call(None, Mock())) + c._state = KeeperState.CONNECTING + self.assertFalse(c._call(None, Mock())) + + class TestZooKeeper(unittest.TestCase): - @patch('patroni.dcs.zookeeper.KazooClient', MockKazooClient) + @patch('patroni.dcs.zookeeper.PatroniKazooClient', MockKazooClient) def setUp(self): self.zk = ZooKeeper({'hosts': ['localhost:2181'], 'scope': 'test', 'name': 'foo', 'ttl': 30, 'retry_timeout': 10, 'loop_wait': 10})