From c8b5003b86fa024090ea8836f5573bb97a644011 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 1 Jun 2016 13:41:49 +0200 Subject: [PATCH] Set __do_not_watch flag when ttl needs to be changed it's more readable comparing to `reset_cluster` --- patroni/dcs/consul.py | 31 +++++++++++++------------------ patroni/dcs/etcd.py | 11 ++++++----- patroni/dcs/zookeeper.py | 2 +- tests/test_consul.py | 6 ++++-- tests/test_etcd.py | 1 + 5 files changed, 25 insertions(+), 26 deletions(-) diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index f7736dfb..887787f2 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -79,32 +79,22 @@ class Consul(AbstractDCS): self._client = ConsulClient(host=host, port=port) self._client.http.patch_default_timeout(config['retry_timeout']/2.0) self._scope = config['scope'] - self.create_or_restore_session() + self.create_session() + self.__do_not_watch = False - def create_or_restore_session(self): + def create_session(self): while not self._session: try: - _, member = self._client.kv.get(self.member_path) - self._session = (member or {}).get('Session') - if self.refresh_session(): - self._client.kv.delete(self.member_path) - except (ConsulException, RequestException): + self.refresh_session() + except ConsulError: logger.info('waiting on consul') sleep(5) def set_ttl(self, ttl): ttl = ttl/2.0 # My experiments have shown that session expires after 2*ttl time if self._ttl != ttl: - if self._session: - try: - self._client.session.destroy(self._session) - except Exception: - logger.exception("Can not destroy session %s", self._session) self._session = None - # force `watch` method to call `AbstractDCS.watch` instead of watching for leader key - self.reset_cluster() - # fire up an event to wake up from `watch` and immediately run HA loop (to create the new session) - self.event.set() + self.__do_not_watch = True self._ttl = ttl def set_retry_timeout(self, retry_timeout): @@ -187,12 +177,13 @@ class Consul(AbstractDCS): raise ConsulError('Consul is not responding properly') def touch_member(self, data, **kwargs): - create_member = self.refresh_session() cluster = self.cluster member = cluster and ([m for m in cluster.members if m.name == self._name] or [None])[0] - if create_member and member: + create_member = self.refresh_session() + if member and (create_member or member.session != self._session): try: self._client.kv.delete(self.member_path) + create_member = True except Exception: return False @@ -253,6 +244,10 @@ class Consul(AbstractDCS): return self._client.kv.delete(self.leader_path, cas=cluster.leader.index) def watch(self, timeout): + if self.__do_not_watch: + self.__do_not_watch = False + return True + cluster = self.cluster if cluster and cluster.leader and cluster.leader.name != self._name and cluster.leader.index: end_time = time.time() + timeout diff --git a/patroni/dcs/etcd.py b/patroni/dcs/etcd.py index 467e1ddd..7b2b7477 100644 --- a/patroni/dcs/etcd.py +++ b/patroni/dcs/etcd.py @@ -200,6 +200,7 @@ class Etcd(AbstractDCS): etcd.EtcdWatcherCleared, etcd.EtcdEventIndexCleared)) self._client = self.get_etcd_client(config) + self.__do_not_watch = False def retry(self, *args, **kwargs): return self._retry.copy()(*args, **kwargs) @@ -217,11 +218,7 @@ class Etcd(AbstractDCS): def set_ttl(self, ttl): ttl = int(ttl) - if self._ttl != ttl: - # force `watch` method to call `AbstractDCS.watch` instead of watching for leader key - self.reset_cluster() - # fire up an event to wake up from `watch` and immediately run HA loop (to update TTL of leader and member) - self.event.set() + self.__do_not_watch = self._ttl != ttl self._ttl = ttl def set_retry_timeout(self, retry_timeout): @@ -320,6 +317,10 @@ class Etcd(AbstractDCS): return self.retry(self._client.delete, self.client_path(''), recursive=True) def watch(self, timeout): + if self.__do_not_watch: + self.__do_not_watch = False + return True + cluster = self.cluster # watch on leader key changes if it is defined and current node is not lock owner if cluster and cluster.leader and cluster.leader.name != self._name and cluster.leader.index: diff --git a/patroni/dcs/zookeeper.py b/patroni/dcs/zookeeper.py index 618cc9bd..e077f726 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -188,7 +188,7 @@ class ZooKeeper(AbstractDCS): self._client.retry(self._inner_load_cluster) except: logger.exception('get_cluster') - self.session_listener(KazooState.LOST) + self.cluster_watcher(None) raise ZooKeeperError('ZooKeeper in not responding properly') def _create(self, path, value, **kwargs): diff --git a/tests/test_consul.py b/tests/test_consul.py index baf28719..32e27462 100644 --- a/tests/test_consul.py +++ b/tests/test_consul.py @@ -56,9 +56,10 @@ class TestConsul(unittest.TestCase): self.c._load_cluster() @patch('time.sleep', Mock(side_effect=SleepException)) - def test_create_or_restore_session(self): + @patch.object(consul.Consul.Session, 'create', Mock(side_effect=ConsulException)) + def test_create_session(self): self.c._session = None - self.assertRaises(SleepException, self.c.create_or_restore_session) + self.assertRaises(SleepException, self.c.create_session) @patch.object(consul.Consul.Session, 'renew', Mock(side_effect=NotFound)) @patch.object(consul.Consul.Session, 'create', Mock(side_effect=ConsulException)) @@ -133,6 +134,7 @@ class TestConsul(unittest.TestCase): @patch.object(consul.Consul.Session, 'destroy', Mock(side_effect=ConsulException)) def test_set_ttl(self): self.c.set_ttl(20) + self.assertTrue(self.c.watch(1)) def test_set_retry_timeout(self): self.c.set_retry_timeout(10) diff --git a/tests/test_etcd.py b/tests/test_etcd.py index d675db43..ebdf6aa2 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -260,3 +260,4 @@ class TestEtcd(unittest.TestCase): def test_set_ttl(self): self.etcd.set_ttl(20) + self.assertTrue(self.etcd.watch(1))