diff --git a/patroni/api.py b/patroni/api.py index 0333e86c..dc83249f 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -48,8 +48,9 @@ class RestApiHandler(BaseHTTPRequestHandler): response = self.get_postgresql_status() patroni = self.server.patroni - if patroni.dcs.cluster: # dcs available - if patroni.dcs.cluster.leader and patroni.dcs.cluster.leader.name == patroni.postgresql.name: # is_leader + cluster = patroni.dcs.cluster + if cluster: # dcs available + if cluster.leader and cluster.leader.name == patroni.postgresql.name: # is_leader status_code = 200 if 'master' in path else 503 elif 'role' not in response: status_code = 503 diff --git a/patroni/dcs.py b/patroni/dcs.py index d007b392..25e44cd1 100644 --- a/patroni/dcs.py +++ b/patroni/dcs.py @@ -4,7 +4,7 @@ import json from collections import namedtuple from patroni.exceptions import DCSError from six.moves.urllib_parse import urlparse, urlunparse, parse_qsl -from threading import Event +from threading import Event, Lock def parse_connection_string(value): @@ -125,7 +125,8 @@ class AbstractDCS: self._scope = config['scope'] self._base_path = '/service/' + self._scope - self.cluster = None + self._cluster = None + self._cluster_thread_lock = Lock() self.event = Event() def client_path(self, path): @@ -156,10 +157,32 @@ class AbstractDCS: return self.client_path(self._LEADER_OPTIME) @abc.abstractmethod + def _load_cluster(self): + """Internally this method should build `Cluster` object which + represents current state and topology of the cluster in DCS. + this method supposed to be called only by `get_cluster` method. + + raise `~DCSError` in case of communication or other problems with DCS. + If the current node was running as a master and exception raised, + instance would be demoted.""" + def get_cluster(self): - """:returns: `Cluster` object which represent current state and topology of the cluster - raise `~DCSError` in case of communication or other problems with DCS. If current instance was - running as a master and exception raised instance would be demoted.""" + with self._cluster_thread_lock: + try: + self._load_cluster() + except: + self._cluster = None + raise + return self._cluster + + @property + def cluster(self): + with self._cluster_thread_lock: + return self._cluster + + def reset_cluster(self): + with self._cluster_thread_lock: + self._cluster = None @abc.abstractmethod def write_leader_optime(self, last_operation): diff --git a/patroni/etcd.py b/patroni/etcd.py index 76e00554..79379700 100644 --- a/patroni/etcd.py +++ b/patroni/etcd.py @@ -171,7 +171,7 @@ class Etcd(AbstractDCS): def member(node): return Member.from_node(node.modifiedIndex, os.path.basename(node.key), node.ttl, node.value) - def get_cluster(self): + def _load_cluster(self): try: result = self.retry(self.client.read, self.client_path(''), recursive=True) nodes = {os.path.relpath(node.key, result.key): node for node in result.leaves} @@ -198,14 +198,12 @@ class Etcd(AbstractDCS): if failover: failover = Failover.from_node(failover.modifiedIndex, failover.value) - self.cluster = Cluster(initialize, leader, last_leader_operation, members, failover) + self._cluster = Cluster(initialize, leader, last_leader_operation, members, failover) except etcd.EtcdKeyNotFound: - self.cluster = Cluster(False, None, None, [], None) + self._cluster = Cluster(False, None, None, [], None) except: - self.cluster = None logger.exception('get_cluster') raise EtcdError('Etcd is not responding properly') - return self.cluster @catch_etcd_errors def touch_member(self, connection_string, ttl=None): @@ -249,10 +247,11 @@ class Etcd(AbstractDCS): return self.retry(self.client.delete, self.initialize_path, prevValue=self._name) def watch(self, timeout): + cluster = self.cluster # watch on leader key changes if it is defined and current node is not lock owner - if self.cluster and self.cluster.leader and self.cluster.leader.name != self._name: + if cluster and cluster.leader and cluster.leader.name != self._name: end_time = time.time() + timeout - index = self.cluster.leader.index + index = cluster.leader.index while index and timeout >= 1: # when timeout is too small urllib3 doesn't have enough time to connect try: diff --git a/patroni/ha.py b/patroni/ha.py index ca79d8e3..9203e153 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -16,13 +16,17 @@ class Ha: self.patroni = patroni self.state_handler = patroni.postgresql self.dcs = patroni.dcs + self.cluster = None self.old_cluster = None self._async_executor = AsyncExecutor() def load_cluster_from_dcs(self): + cluster = self.dcs.get_cluster() + # We want to keep the state of cluster when it was healhy - if not self.dcs.get_cluster().is_unlocked() or not self.old_cluster: - self.old_cluster = self.dcs.cluster + if not cluster.is_unlocked() or not self.old_cluster: + self.old_cluster = cluster + self.cluster = cluster def acquire_lock(self): return self.dcs.attempt_to_acquire_leader() @@ -37,7 +41,7 @@ class Ha: return ret def has_lock(self): - lock_owner = self.dcs.cluster.leader and self.dcs.cluster.leader.name + lock_owner = self.cluster.leader and self.cluster.leader.name logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name) return lock_owner == self.state_handler.name @@ -55,8 +59,8 @@ class Ha: pass self.dcs.touch_member(json.dumps(data, separators=(',', ':'))) - def copy_backup_from_leader(self): - if self.state_handler.bootstrap(self.dcs.cluster.leader): + def copy_backup_from_leader(self, leader): + if self.state_handler.bootstrap(leader): logger.info('bootstrapped from leader') else: self.state_handler.stop('immediate') @@ -64,14 +68,11 @@ class Ha: logger.error('failed to bootstrap from leader') def bootstrap(self): - if not self.dcs.cluster.is_unlocked(): # cluster already has leader - if self._async_executor.busy: - self.copy_backup_from_leader() - else: - self._async_executor.schedule('bootstrap from leader') - self._async_executor.run_async(self.copy_backup_from_leader) + if not self.cluster.is_unlocked(): # cluster already has leader + self._async_executor.schedule('bootstrap from leader') + self._async_executor.run_async(self.copy_backup_from_leader, args=(self.cluster.leader, )) return 'trying to bootstrap from leader' - elif not self.dcs.cluster.initialize: # no initialize key + elif not self.cluster.initialize: # no initialize key if self.dcs.initialize(): # race for initialization try: self.state_handler.bootstrap() @@ -91,11 +92,12 @@ class Ha: def recover(self): has_lock = self.has_lock() - self.state_handler.write_recovery_conf(None if has_lock else self.dcs.cluster.leader) + self.state_handler.write_recovery_conf(None if has_lock else self.cluster.leader) if not self.state_handler.start(): if not has_lock: return 'failed to start postgres' self.dcs.delete_leader() + self.dcs.reset_cluster() return 'removed leader key after trying and failing to start postgres' if not has_lock: return 'started as a secondary' @@ -105,9 +107,11 @@ class Ha: def follow_the_leader(self, demote_reason, follow_reason, refresh=True): refresh and self.load_cluster_from_dcs() ret = demote_reason if self.state_handler.is_leader() else follow_reason - if not self.state_handler.check_recovery_conf(self.dcs.cluster.leader): + if not self.state_handler.check_recovery_conf(self.cluster.leader): self._async_executor.schedule('changing primary_conninfo and restarting') - self._async_executor.run_async(self.state_handler.follow_the_leader, (self.dcs.cluster.leader, )) + leader = self.cluster.leader + leader = None if (leader and leader.name) == self.state_handler.name else leader + self._async_executor.run_async(self.state_handler.follow_the_leader, (leader, )) return ret def enforce_master_role(self, message, promote_message): @@ -150,8 +154,7 @@ class Ha: if self.state_handler.is_leader(): return True - if check_replication_lag and \ - not self.state_handler.check_replication_lag(self.dcs.cluster.last_leader_operation): + if check_replication_lag and not self.state_handler.check_replication_lag(self.cluster.last_leader_operation): return False # Too far behind last reported xlog location on master # Prepare list of nodes to run check against @@ -182,13 +185,13 @@ class Ha: return ret def manual_failover_process_no_leader(self): - failover = self.dcs.cluster.failover + failover = self.cluster.failover if failover.member: # manual failover to specific member if failover.member == self.state_handler.name: # manual failover to me return True # find specific node and check that it is healthy - members = [m for m in self.dcs.cluster.members if m.name == failover.member] + members = [m for m in self.cluster.members if m.name == failover.member] if members: member, reachable, in_recovery, xlog_location = self.fetch_node_status(members[0]) if reachable: # node is healthy @@ -204,7 +207,7 @@ class Ha: if failover.leader: if self.state_handler.name == failover.leader: # I was the leader # exclude me and desired member which is unhealthy (failover.member can be None) - members = [m for m in self.dcs.cluster.members if m.name != failover.member] + members = [m for m in self.cluster.members if m.name != failover.member] if self.is_failover_possible(members): # check that there are healthy members return False else: # I was the leader and it looks like currently I am the only healthy member @@ -213,28 +216,29 @@ class Ha: # at this point we assume that our node is a candidate for a failover among all nodes except former leader # exclude former leader from the list (failover.leader can be None) - members = [m for m in self.dcs.cluster.members if m.name != failover.leader] + members = [m for m in self.cluster.members if m.name != failover.leader] return self._is_healthiest_node(members, check_replication_lag=False) def is_healthiest_node(self): - if self.dcs.cluster.failover: + if self.cluster.failover: return self.manual_failover_process_no_leader() # run usual health check - members = {m.name: m for m in self.dcs.cluster.members + self.old_cluster.members} + members = {m.name: m for m in self.cluster.members + self.old_cluster.members} return self._is_healthiest_node(members.values()) def demote(self, delete_leader=True): if delete_leader: self.state_handler.stop() self.dcs.delete_leader() + self.dcs.reset_cluster() self.state_handler.follow_the_leader(None) def process_manual_failover_from_leader(self): - failover = self.dcs.cluster.failover + failover = self.cluster.failover if not failover.leader or failover.leader == self.state_handler.name: if not failover.member or failover.member != self.state_handler.name: - members = [m for m in self.dcs.cluster.members if not failover.member or m.name == failover.member] + members = [m for m in self.cluster.members if not failover.member or m.name == failover.member] if self.is_failover_possible(members): # check that there are healthy members self._async_executor.schedule('manual failover: demote') self._async_executor.run_async(self.demote) @@ -245,15 +249,15 @@ class Ha: logger.warning('manual failover: I am already the leader, no need to failover') else: logger.warning('manual failover: leader name does not match: %s != %s', - self.dcs.cluster.failover.leader, self.state_handler.name) + self.cluster.failover.leader, self.state_handler.name) logger.info('Trying to clean up failover key') - self.dcs.manual_failover('', '', self.dcs.cluster.failover.index) + self.dcs.manual_failover('', '', self.cluster.failover.index) def process_unhealthy_cluster(self): if self.is_healthiest_node(): if self.acquire_lock(): - if self.dcs.cluster.failover: + if self.cluster.failover: logger.info('Cleanning up failover key after acquiring leader lock...') self.dcs.manual_failover('', '') return self.enforce_master_role('acquired session lock as a leader', @@ -267,7 +271,7 @@ class Ha: def process_healthy_cluster(self): if self.has_lock(): - if self.dcs.cluster.failover: + if self.cluster.failover: msg = self.process_manual_failover_from_leader() if msg is not None: return msg @@ -307,22 +311,21 @@ class Ha: else: return (False, 'restart failed') - def reinitialize(self): + def reinitialize(self, cluster): self.state_handler.stop('immediate') self.state_handler.remove_data_directory() - self.load_cluster_from_dcs() - self.bootstrap() + self.copy_backup_from_leader(cluster.leader) def process_scheduled_action(self): if self.reinitialize_scheduled(): - if self.dcs.cluster.is_unlocked(): + if self.cluster.is_unlocked(): logger.error('Cluster has no leader, can not reinitialize') self._async_executor.reset_scheduled_action() elif self.has_lock(): logger.error('I am the leader, can not reinitialize') self._async_executor.reset_scheduled_action() else: - self._async_executor.run_async(self.reinitialize) + self._async_executor.run_async(self.reinitialize, args=(self.cluster, )) return 'reinitialize started' def handle_long_action_in_progress(self): @@ -331,7 +334,7 @@ class Ha: return 'updated leader lock during ' + self._async_executor.scheduled_action else: return 'failed to update leader lock during ' + self._async_executor.scheduled_action - elif self.dcs.cluster.is_unlocked(): + elif self.cluster.is_unlocked(): return 'not healthy enough for leader race' else: return self._async_executor.scheduled_action + ' in progress' @@ -343,7 +346,7 @@ class Ha: self.touch_member() # cluster has leader key but not initialize key - if not self.dcs.cluster.is_unlocked() and not self.dcs.cluster.initialize: + if not self.cluster.is_unlocked() and not self.cluster.initialize: self.dcs.initialize() # fix it if self._async_executor.busy: @@ -358,7 +361,7 @@ class Ha: if self.state_handler.data_directory_empty(): return self.bootstrap() # new node # "bootstrap", but data directory is not empty - elif not self.dcs.cluster.initialize and self.dcs.cluster.is_unlocked(): + elif not self.cluster.initialize and self.cluster.is_unlocked(): self.dcs.initialize() # try to start dead postgres @@ -368,12 +371,12 @@ class Ha: return msg try: - if self.dcs.cluster.is_unlocked(): + if self.cluster.is_unlocked(): return self.process_unhealthy_cluster() else: return self.process_healthy_cluster() finally: - self.state_handler.sync_replication_slots(self.dcs.cluster) + self.state_handler.sync_replication_slots(self.cluster) except DCSError: logger.error('Error communicating with DCS') if self.state_handler.is_running() and self.state_handler.is_leader(): diff --git a/patroni/zookeeper.py b/patroni/zookeeper.py index 68786485..c328ae52 100644 --- a/patroni/zookeeper.py +++ b/patroni/zookeeper.py @@ -164,9 +164,9 @@ class ZooKeeper(AbstractDCS): # get last leader operation self.last_leader_operation = self.get_node(self.leader_optime_path) if self.fetch_cluster else None self.last_leader_operation = 0 if self.last_leader_operation is None else int(self.last_leader_operation[0]) - self.cluster = Cluster(initialize, leader, self.last_leader_operation, members, failover) + self._cluster = Cluster(initialize, leader, self.last_leader_operation, members, failover) - def get_cluster(self): + def _load_cluster(self): if self.exhibitor and self.exhibitor.poll(): self.client.set_hosts(self.exhibitor.zookeeper_hosts) @@ -174,11 +174,9 @@ class ZooKeeper(AbstractDCS): try: self.client.retry(self._inner_load_cluster) except: - self.cluster = None logger.exception('get_cluster') self.session_listener(KazooState.LOST) raise ZooKeeperError('ZooKeeper in not responding properly') - return self.cluster def _create(self, path, value, **kwargs): try: @@ -206,7 +204,8 @@ class ZooKeeper(AbstractDCS): return self._create(self.initialize_path, self._name, makepath=True) def touch_member(self, data, ttl=None): - me = self.cluster and ([m for m in self.cluster.members if m.name == self._name] or [None])[0] + cluster = self.cluster + me = cluster and ([m for m in cluster.members if m.name == self._name] or [None])[0] path = self.member_path data = data.encode('utf-8') create = not me diff --git a/tests/test_ha.py b/tests/test_ha.py index 392ad663..0a986280 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -85,7 +85,7 @@ class MockPatroni: def run_async(func, args=()): - func(args) if args else func() + func(*args) if args else func() class TestHa(unittest.TestCase): @@ -101,7 +101,7 @@ class TestHa(unittest.TestCase): self.ha = Ha(MockPatroni(self.p, self.e)) self.ha._async_executor.run_async = run_async self.ha.old_cluster = self.e.get_cluster() - self.e.cluster = get_cluster_not_initialized_without_leader() + self.ha.cluster = get_cluster_not_initialized_without_leader() self.ha.load_cluster_from_dcs = Mock() def test_update_lock(self): @@ -161,28 +161,28 @@ class TestHa(unittest.TestCase): self.assertEquals(self.ha.run_cycle(), 'following a different leader because i am not the healthiest node') def test_promote_because_have_lock(self): - self.e.cluster.is_unlocked = false + self.ha.cluster.is_unlocked = false self.ha.has_lock = true self.p.is_leader = false self.assertEquals(self.ha.run_cycle(), 'promoted self to leader because i had the session lock') def test_leader_with_lock(self): - self.e.cluster.is_unlocked = false + self.ha.cluster.is_unlocked = false self.ha.has_lock = true self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') def test_demote_because_not_having_lock(self): - self.e.cluster.is_unlocked = false + self.ha.cluster.is_unlocked = false self.assertEquals(self.ha.run_cycle(), 'demoting self because i do not have the lock and i was a leader') def test_demote_because_update_lock_failed(self): - self.e.cluster.is_unlocked = false + self.ha.cluster.is_unlocked = false self.ha.has_lock = true self.ha.update_lock = false self.assertEquals(self.ha.run_cycle(), 'demoting self because i do not have the lock and i was a leader') def test_follow_the_leader(self): - self.e.cluster.is_unlocked = false + self.ha.cluster.is_unlocked = false self.p.is_leader = false self.assertEquals(self.ha.run_cycle(), 'no action. i am a secondary and i am following a leader') @@ -191,27 +191,25 @@ class TestHa(unittest.TestCase): self.assertEquals(self.ha.run_cycle(), 'demoted self because DCS is not accessible and i was a leader') def test_bootstrap_from_leader(self): - self.e.cluster = get_cluster_initialized_with_leader() + self.ha.cluster = get_cluster_initialized_with_leader() self.p.bootstrap = false self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap from leader') - self.ha._async_executor._busy = True - self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap from leader') def test_bootstrap_waiting_for_leader(self): - self.e.cluster = get_cluster_initialized_without_leader() + self.ha.cluster = get_cluster_initialized_without_leader() self.assertEquals(self.ha.bootstrap(), 'waiting for leader to bootstrap') def test_bootstrap_initialize_lock_failed(self): - self.e.cluster = get_cluster_not_initialized_without_leader() + self.ha.cluster = get_cluster_not_initialized_without_leader() self.assertEquals(self.ha.bootstrap(), 'failed to acquire initialize lock') def test_bootstrap_initialized_new_cluster(self): - self.e.cluster = get_cluster_not_initialized_without_leader() + self.ha.cluster = get_cluster_not_initialized_without_leader() self.e.initialize = true self.assertEquals(self.ha.bootstrap(), 'initialized a new cluster') def test_bootstrap_release_initialize_key_on_failure(self): - self.e.cluster = get_cluster_not_initialized_without_leader() + self.ha.cluster = get_cluster_not_initialized_without_leader() self.e.initialize = true self.p.bootstrap = Mock(side_effect=PostgresException("Could not bootstrap master PostgreSQL")) self.assertRaises(PostgresException, self.ha.bootstrap) @@ -222,7 +220,7 @@ class TestHa(unittest.TestCase): self.ha.run_cycle() self.assertIsNone(self.ha._async_executor.scheduled_action) - self.e.cluster = get_cluster_initialized_with_leader() + self.ha.cluster = get_cluster_initialized_with_leader() self.ha.has_lock = true self.ha.schedule_reinitialize() self.ha.run_cycle() @@ -244,7 +242,7 @@ class TestHa(unittest.TestCase): self.assertTrue(self.ha.restart_scheduled()) self.assertEquals(self.ha.run_cycle(), 'not healthy enough for leader race') - self.e.cluster = get_cluster_initialized_with_leader() + self.ha.cluster = get_cluster_initialized_with_leader() self.assertEquals(self.ha.run_cycle(), 'restart in progress') self.ha.has_lock = true @@ -256,26 +254,26 @@ class TestHa(unittest.TestCase): @patch('requests.get', requests_get) def test_manual_failover_from_leader(self): self.ha.has_lock = true - self.e.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', '')) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', '')) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') - self.e.cluster = get_cluster_initialized_with_leader(Failover(0, '', MockPostgresql.name)) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', MockPostgresql.name)) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') - self.e.cluster = get_cluster_initialized_with_leader(Failover(0, '', 'blabla')) + self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', 'blabla')) self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock') f = Failover(0, MockPostgresql.name, '') - self.e.cluster = get_cluster_initialized_with_leader(f) + self.ha.cluster = get_cluster_initialized_with_leader(f) self.assertEquals(self.ha.run_cycle(), 'manual failover: demoting myself') @patch('requests.get', requests_get) def test_manual_failover_process_no_leader(self): self.p.is_leader = false - self.e.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', MockPostgresql.name)) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', MockPostgresql.name)) self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') - self.e.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'leader')) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'leader')) self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') self.ha.fetch_node_status = lambda e: (e, True, True, 0) # accessible, in_recovery self.assertEquals(self.ha.run_cycle(), 'following a different leader because i am not the healthiest node') - self.e.cluster = get_cluster_initialized_without_leader(failover=Failover(0, MockPostgresql.name, '')) + self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, MockPostgresql.name, '')) self.assertEquals(self.ha.run_cycle(), 'following a different leader because i am not the healthiest node') self.ha.fetch_node_status = lambda e: (e, False, True, 0) # accessible, in_recovery self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock')