diff --git a/features/basic_replication.feature b/features/basic_replication.feature index bc1f6286..8aa112fa 100644 --- a/features/basic_replication.feature +++ b/features/basic_replication.feature @@ -21,6 +21,7 @@ Feature: basic replication Then "sync" key in DCS has sync_standby=postgres2 after 10 seconds When I start postgres1 And "members/postgres1" key in DCS has state=running after 10 seconds + And I sleep for 2 seconds When I issue a GET request to http://127.0.0.1:8010/sync Then I receive a response code 200 When I issue a GET request to http://127.0.0.1:8009/async diff --git a/patroni/__init__.py b/patroni/__init__.py index edf6cf10..ee4c5824 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -188,7 +188,7 @@ def main(): if ret == (0, 0): break elif ret[0] != pid: - logging.info('Reaped pid=%s, exit status=%s', *ret) + logger.info('Reaped pid=%s, exit status=%s', *ret) except OSError: pass diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index bde5542a..6bcd3216 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -144,7 +144,6 @@ class Consul(AbstractDCS): retry_exceptions=(ConsulInternalError, HTTPException, HTTPError, socket.error, socket.timeout)) - self._my_member_data = {} kwargs = {} if 'url' in config: r = urlparse(config['url']) @@ -318,13 +317,12 @@ class Consul(AbstractDCS): except Exception: return False - if not create_member and member and deep_compare(data, self._my_member_data): + if not create_member and member and deep_compare(data, member.data): return True try: args = {} if permanent else {'acquire': self._session} self._client.kv.put(self.member_path, json.dumps(data, separators=(',', ':')), **args) - self._my_member_data = data return True except Exception: logger.exception('touch_member') @@ -408,7 +406,7 @@ class Consul(AbstractDCS): idx, _ = self._client.kv.get(self.leader_path, index=leader_index, wait=str(timeout) + 's') return str(idx) != str(leader_index) except (ConsulException, HTTPException, HTTPError, socket.error, socket.timeout): - logging.exception('watch') + logger.exception('watch') timeout = end_time - time.time() diff --git a/patroni/dcs/etcd.py b/patroni/dcs/etcd.py index a5809d1c..5d664a72 100644 --- a/patroni/dcs/etcd.py +++ b/patroni/dcs/etcd.py @@ -496,7 +496,7 @@ class Etcd(AbstractDCS): @catch_etcd_errors def touch_member(self, data, ttl=None, permanent=False): data = json.dumps(data, separators=(',', ':')) - return self.retry(self._client.set, self.member_path, data, None if permanent else ttl or self._ttl) + return self._client.set(self.member_path, data, None if permanent else ttl or self._ttl) @catch_etcd_errors def take_leader(self): diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index 842ec5ea..423a77fa 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -399,7 +399,7 @@ class Kubernetes(AbstractDCS): except KeyboardInterrupt: raise except Exception: - logging.exception('watch') + logger.exception('watch') timeout = end_time - time.time() diff --git a/patroni/dcs/zookeeper.py b/patroni/dcs/zookeeper.py index df476a4e..315258ae 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -59,7 +59,6 @@ class ZooKeeper(AbstractDCS): max_delay=1, max_tries=-1, sleep_func=time.sleep)) self._client.add_listener(self.session_listener) - self._my_member_data = {} self._fetch_cluster = True self._orig_kazoo_connect = self._client._connection._connect @@ -208,45 +207,52 @@ class ZooKeeper(AbstractDCS): self.cluster_watcher(None) raise ZooKeeperError('ZooKeeper in not responding properly') - def _create(self, path, value, **kwargs): + def _create(self, path, value, retry=False, ephemeral=False): try: - self._client.retry(self._client.create, path, value.encode('utf-8'), **kwargs) + if retry: + self._client.retry(self._client.create, path, value, makepath=True, ephemeral=ephemeral) + else: + self._client.create_async(path, value, makepath=True, ephemeral=ephemeral).get(timeout=1) return True except Exception: - return False + logger.exception('Failed to create %s', path) + return False def attempt_to_acquire_leader(self, permanent=False): - ret = self._create(self.leader_path, self._name, makepath=True, ephemeral=not permanent) + ret = self._create(self.leader_path, self._name.encode('utf-8'), retry=True, ephemeral=not permanent) if not ret: logger.info('Could not take out TTL lock') return ret - def __set_failover_or_sync_state_value(self, key, value, index=None): + def _set_or_create(self, key, value, index=None, retry=False, do_not_create_empty=False): + value = value.encode('utf-8') try: - self._client.retry(self._client.set, key, value.encode('utf-8'), version=index or -1) + if retry: + self._client.retry(self._client.set, key, value, version=index or -1) + else: + self._client.set_async(key, value, version=index or -1).get(timeout=1) return True except NoNodeError: - return value == '' or (index is None and self._create(key, value)) + if do_not_create_empty and not value: + return True + elif index is None: + return self._create(key, value, retry) + else: + return False except Exception: - logging.exception('set_failover_value') - return False + logger.exception('Failed to update %s', key) + return False def set_failover_value(self, value, index=None): - return self.__set_failover_or_sync_state_value(self.failover_path, value, index) + return self._set_or_create(self.failover_path, value, index) def set_config_value(self, value, index=None): - try: - self._client.retry(self._client.set, self.config_path, value.encode('utf-8'), version=index or -1) - return True - except NoNodeError: - return index is None and self._create(self.config_path, value) - except Exception: - logging.exception('set_config_value') - return False + return self._set_or_create(self.config_path, value, index, retry=True) def initialize(self, create_new=True, sysid=""): - return self._create(self.initialize_path, sysid, makepath=True) if create_new \ - else self._client.retry(self._client.set, self.initialize_path, sysid.encode("utf-8")) + sysid = sysid.encode('utf-8') + return self._create(self.initialize_path, sysid, retry=True) if create_new \ + else self._client.retry(self._client.set, self.initialize_path, sysid) def touch_member(self, data, ttl=None, permanent=False): cluster = self.cluster @@ -262,13 +268,12 @@ class ZooKeeper(AbstractDCS): member = None if member: - if deep_compare(data, self._my_member_data): + if deep_compare(data, member.data): return True else: try: self._client.create_async(self.member_path, encoded_data, makepath=True, ephemeral=not permanent).get(timeout=1) - self._my_member_data = data return True except Exception as e: if not isinstance(e, NodeExistsError): @@ -276,7 +281,6 @@ class ZooKeeper(AbstractDCS): return False try: self._client.set_async(self.member_path, encoded_data).get(timeout=1) - self._my_member_data = data return True except Exception: logger.exception('touch_member') @@ -286,30 +290,14 @@ class ZooKeeper(AbstractDCS): def take_leader(self): return self.attempt_to_acquire_leader() - def __write_leader_optime_or_history_value(self, key, value): - value = value.encode('utf-8') - try: - self._client.set_async(key, value).get(timeout=1) - return True - except NoNodeError: - try: - self._client.create_async(key, value, makepath=True).get(timeout=1) - return True - except Exception: - logger.exception('Failed to create %s', key) - except Exception: - logger.exception('Failed to update %s', key) - return False - def _write_leader_optime(self, last_operation): - return self.__write_leader_optime_or_history_value(self.leader_optime_path, last_operation) + return self._set_or_create(self.leader_optime_path, last_operation) def _update_leader(self): return True def delete_leader(self): self._client.restart() - self._my_member_data = None return True def _cancel_initialization(self): @@ -330,10 +318,10 @@ class ZooKeeper(AbstractDCS): return True def set_history_value(self, value): - return self.__write_leader_optime_or_history_value(self.history_path, value) + return self._set_or_create(self.history_path, value) def set_sync_state_value(self, value, index=None): - return self.__set_failover_or_sync_state_value(self.sync_path, value, index) + return self._set_or_create(self.sync_path, value, index, retry=True, do_not_create_empty=True) def delete_sync_state(self, index=None): return self.set_sync_state_value("{}", index) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 84b8963c..b5c305b3 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -945,7 +945,7 @@ class Postgresql(object): return 'is_in_recovery=true' return cur.execute('CHECKPOINT') except psycopg2.Error: - logging.exception('Exception during CHECKPOINT') + logger.exception('Exception during CHECKPOINT') return 'not accessible or not healty' def stop(self, mode='fast', block_callbacks=False, checkpoint=None, on_safepoint=None): diff --git a/tests/test_consul.py b/tests/test_consul.py index 496caaa2..c26ca3c6 100644 --- a/tests/test_consul.py +++ b/tests/test_consul.py @@ -116,7 +116,8 @@ class TestConsul(unittest.TestCase): self.c.touch_member({'balbla': 'blabla'}) self.c.touch_member({'balbla': 'blabla'}) self.c.refresh_session = Mock(return_value=False) - self.c.touch_member({'balbla': 'blabla'}) + 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'}) @patch.object(consul.Consul.KV, 'put', Mock(return_value=False)) def test_take_leader(self): diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py index 452fd7df..d8de0680 100644 --- a/tests/test_zookeeper.py +++ b/tests/test_zookeeper.py @@ -156,7 +156,7 @@ class TestZooKeeper(unittest.TestCase): self.zk.set_failover_value('Exception') def test_set_config_value(self): - self.zk.set_config_value('') + self.zk.set_config_value('', 1) self.zk.set_config_value('ok') self.zk.set_config_value('Exception') @@ -179,7 +179,8 @@ class TestZooKeeper(unittest.TestCase): self.zk.touch_member({'retry': 'retry'}) self.zk._fetch_cluster = True self.zk.get_cluster() - self.zk.touch_member({'retry': 'retry'}) + self.zk.touch_member({'conn_url': 'postgres://repuser:rep-pass@localhost:5434/postgres', + 'api_url': 'http://127.0.0.1:8009/patroni'}) def test_take_leader(self): self.zk.take_leader()