diff --git a/patroni/ha.py b/patroni/ha.py index b50884f1..fae3535e 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -200,6 +200,13 @@ class Ha(object): with self._is_leader_lock: self._is_leader = time.time() + self.dcs.ttl if value else 0 + def sync_mode_is_active(self) -> bool: + """Check whether synchronous replication is requested and already active. + + :returns: ``True`` if the primary already put its name into the ``/sync`` in DCS. + """ + return self.is_synchronous_mode() and not self.cluster.sync.is_empty + def load_cluster_from_dcs(self) -> None: cluster = self.dcs.get_cluster() @@ -465,7 +472,7 @@ class Ha(object): if timeout == 0: # We are requested to prefer failing over to restarting primary. But see first if there # is anyone to fail over to. - if self.is_failover_possible(self.cluster.members): + if self.is_failover_possible(): self.watchdog.disable() logger.info("Primary crashed. Failing over.") self.demote('immediate') @@ -810,6 +817,8 @@ class Ha(object): return _MemberStatus.unknown(member) def fetch_nodes_statuses(self, members: List[Member]) -> List[_MemberStatus]: + if not members: + return [] pool = ThreadPool(len(members)) results = pool.map(self.fetch_node_status, members) # Run API calls on members in parallel pool.close() @@ -888,51 +897,50 @@ class Ha(object): # Prepare list of nodes to run check against members = [m for m in members if m.name != self.state_handler.name and not m.nofailover and m.api_url] - if members: - for st in self.fetch_nodes_statuses(members): - if st.failover_limitation() is None: - if st.in_recovery is False: - logger.warning('Primary (%s) is still alive', st.member.name) + for st in self.fetch_nodes_statuses(members): + if st.failover_limitation() is None: + if st.in_recovery is False: + logger.warning('Primary (%s) is still alive', st.member.name) + return False + if my_wal_position < st.wal_position: + logger.info('Wal position of %s is ahead of my wal position', st.member.name) + # In synchronous mode the former leader might be still accessible and even be ahead of us. + # We should not disqualify himself from the leader race in such a situation. + if not self.sync_mode_is_active() or not self.cluster.sync.leader_matches(st.member.name): return False - if my_wal_position < st.wal_position: - logger.info('Wal position of %s is ahead of my wal position', st.member.name) - # In synchronous mode the former leader might be still accessible and even be ahead of us. - # We should not disqualify himself from the leader race in such a situation. - if not self.is_synchronous_mode() or self.cluster.sync.is_empty\ - or not self.cluster.sync.leader_matches(st.member.name): - return False - logger.info('Ignoring the former leader being ahead of us') + logger.info('Ignoring the former leader being ahead of us') return True - def is_failover_possible(self, members: List[Member], check_synchronous: Optional[bool] = True, - cluster_lsn: Optional[int] = 0) -> bool: - """Checks whether one of the members from the list can possibly win the leader race. + def is_failover_possible(self, *, cluster_lsn: int = 0, exclude_failover_candidate: bool = False) -> bool: + """Checks whether any of the cluster members is allowed to promote and is healthy enough for that. - :param members: list of members to check - :param check_synchronous: consider only members that are known to be listed in /sync key when sync replication. - :param cluster_lsn: to calculate replication lag and exclude member if it is laggin - :returns: `True` if there are members eligible to be the new leader + :param cluster_lsn: to calculate replication lag and exclude member if it is lagging. + :param exclude_failover_candidate: if ``True``, exclude :attr:`failover.candidate` from the members + list against which the failover possibility checks are run. + :returns: `True` if there are members eligible to become the new leader. """ + candidates = self.get_failover_candidates(exclude_failover_candidate) + + if self.is_synchronous_mode() and self.cluster.failover and self.cluster.failover.candidate and not candidates: + logger.warning('Failover candidate=%s does not match with sync_standbys=%s', + self.cluster.failover.candidate, self.cluster.sync.sync_standby) + elif not candidates: + logger.warning('manual failover: candidates list is empty') + ret = False cluster_timeline = self.cluster.timeline - members = [m for m in members if m.name != self.state_handler.name and not m.nofailover and m.api_url] - if check_synchronous and self.is_synchronous_mode() and not self.cluster.sync.is_empty: - members = [m for m in members if self.cluster.sync.matches(m.name)] - if members: - for st in self.fetch_nodes_statuses(members): - not_allowed_reason = st.failover_limitation() - if not_allowed_reason: - logger.info('Member %s is %s', st.member.name, not_allowed_reason) - elif cluster_lsn and st.wal_position < cluster_lsn or\ - not cluster_lsn and self.is_lagging(st.wal_position): - logger.info('Member %s exceeds maximum replication lag', st.member.name) - elif self.check_timeline() and (not st.timeline or st.timeline < cluster_timeline): - logger.info('Timeline %s of member %s is behind the cluster timeline %s', - st.timeline, st.member.name, cluster_timeline) - else: - ret = True - else: - logger.warning('manual failover: members list is empty') + for st in self.fetch_nodes_statuses(candidates): + not_allowed_reason = st.failover_limitation() + if not_allowed_reason: + logger.info('Member %s is %s', st.member.name, not_allowed_reason) + elif cluster_lsn and st.wal_position < cluster_lsn or \ + not cluster_lsn and self.is_lagging(st.wal_position): + logger.info('Member %s exceeds maximum replication lag', st.member.name) + elif self.check_timeline() and (not st.timeline or st.timeline < cluster_timeline): + logger.info('Timeline %s of member %s is behind the cluster timeline %s', + st.timeline, st.member.name, cluster_timeline) + else: + ret = True return ret def manual_failover_process_no_leader(self) -> Optional[bool]: @@ -981,9 +989,8 @@ class Ha(object): # try to pick some other members to failover and check that they are healthy if failover.leader: if self.state_handler.name == failover.leader: # I was the leader - # exclude me and desired member which is unhealthy (failover.candidate can be None) - members = [m for m in self.cluster.members if m.name not in (failover.candidate, failover.leader)] - if self.is_failover_possible(members): # check that there are healthy members + # exclude desired member which is unhealthy if it was specified + if self.is_failover_possible(exclude_failover_candidate=bool(failover.candidate)): return False else: # I was the leader and it looks like currently I am the only healthy member return True @@ -1036,8 +1043,8 @@ class Ha(object): if self.cluster.failover: # When doing a switchover in synchronous mode only synchronous nodes and former leader are allowed to race - if self.is_synchronous_mode() and self.cluster.failover.leader and \ - not self.cluster.sync.is_empty and not self.cluster.sync.matches(self.state_handler.name, True): + if self.sync_mode_is_active() and not self.cluster.sync.matches(self.state_handler.name, True) and \ + self.cluster.failover.leader: return False return self.manual_failover_process_no_leader() or False @@ -1058,7 +1065,7 @@ class Ha(object): all_known_members += self.cluster.members # When in sync mode, only last known primary and sync standby are allowed to promote automatically. - if self.is_synchronous_mode() and not self.cluster.sync.is_empty: + if self.sync_mode_is_active(): if not self.cluster.sync.matches(self.state_handler.name, True): return False # pick between synchronous candidates so we minimize unnecessary failovers/demotions @@ -1109,9 +1116,7 @@ class Ha(object): # It could happen if Postgres is still archiving the backlog of WAL files. # If we know that there are replicas that received the shutdown checkpoint # location, we can remove the leader key and allow them to start leader race. - - # for a manual failover/switchover with a candidate, we should check the requested candidate only - if self.is_failover_possible(self.get_failover_candidates(), cluster_lsn=checkpoint_location): + if self.is_failover_possible(cluster_lsn=checkpoint_location): self.state_handler.set_role('demoted') with self._async_executor: self.release_leader_key_voluntarily(checkpoint_location) @@ -1213,19 +1218,11 @@ class Ha(object): if not failover.candidate or failover.candidate != self.state_handler.name: if not failover.candidate and self.is_paused(): logger.warning('Failover is possible only to a specific candidate in a paused state') + elif self.is_failover_possible(): + ret = self._async_executor.try_run_async('manual failover: demote', self.demote, ('graceful',)) + return ret or 'manual failover: demoting myself' else: - if self.is_synchronous_mode(): - members = self.get_failover_candidates(check_sync=True) - if failover.candidate and not members: - logger.warning('Failover candidate=%s does not match with sync_standbys=%s', - failover.candidate, self.cluster.sync.sync_standby) - else: - members = self.get_failover_candidates() - if self.is_failover_possible(members, False): # check that there are healthy members - ret = self._async_executor.try_run_async('manual failover: demote', self.demote, ('graceful',)) - return ret or 'manual failover: demoting myself' - else: - logger.warning('manual failover: no healthy members found, failover is not possible') + logger.warning('manual failover: no healthy members found, failover is not possible') else: logger.warning('manual failover: I am already the leader, no need to failover') else: @@ -1486,7 +1483,7 @@ class Ha(object): if self.has_lock() and self.update_lock(): if self._async_executor.scheduled_action == 'doing crash recovery in a single user mode': time_left = self.global_config.primary_start_timeout - (time.time() - self._crash_recovery_started) - if time_left <= 0 and self.is_failover_possible(self.cluster.members): + if time_left <= 0 and self.is_failover_possible(): logger.info("Demoting self because crash recovery is taking too long") self.state_handler.cancellable.cancel(True) self.demote('immediate') @@ -1598,7 +1595,7 @@ class Ha(object): time_left = timeout - self.state_handler.time_in_state() if time_left <= 0: - if self.is_failover_possible(self.cluster.members): + if self.is_failover_possible(): logger.info("Demoting self because primary startup is taking too long") self.demote('immediate') return 'stopped PostgreSQL because of startup timeout' @@ -1876,8 +1873,7 @@ class Ha(object): # If we know that there are replicas that received the shutdown checkpoint # location, we can remove the leader key and allow them to start leader race. - # for a manual failover/switchover with a candidate, we should check the requested candidate only - if self.is_failover_possible(self.get_failover_candidates(), cluster_lsn=checkpoint_location): + if self.is_failover_possible(cluster_lsn=checkpoint_location): self.dcs.delete_leader(checkpoint_location) status['deleted'] = True else: @@ -1938,23 +1934,30 @@ class Ha(object): name = member.name if member else 'remote_member:{}'.format(uuid.uuid1()) return RemoteMember(name, data) - def get_failover_candidates(self, check_sync: bool = False) -> List[Member]: - """Return list of candidates for either manual or automatic failover. + def get_failover_candidates(self, exclude_failover_candidate: bool) -> List[Member]: + """Return a list of candidates for either manual or automatic failover. - Mainly used to later be passed to ``Ha.is_failover_possible()``. + Exclude non-sync members when in synchronous mode, the current node (its checks are always performed earlier) + and the candidate if required. If failover candidate exclusion is not requested and a candidate is specified + in the /failover key, return the candidate only. + The result is further evaluated in the caller :func:`Ha.is_failover_possible` to check if any member is actually + healthy enough and is allowed to poromote. - :param check_sync: if ``True``, also check against the sync key members + :param exclude_failover_candidate: if ``True``, exclude :attr:`failover.candidate` from the candidates. - :returns: a list of ``Member`` ojects or an empty list if there is no candidate available + :returns: a list of :class:`Member` ojects or an empty list if there is no candidate available. """ failover = self.cluster.failover - if check_sync: + exclude = [self.state_handler.name] + ([failover.candidate] if failover and exclude_failover_candidate else []) + + def is_eligible(node: Member) -> bool: # TODO: allow manual failover (=no leader specified) to async node - # every sync_standby or the candidate specified if is in sync_standbys - return [m for m in self.cluster.members - if self.cluster.sync.matches(m.name) - and (not failover or not failover.candidate or m.name == failover.candidate)] - else: - # every member or the candidate specified - return [m for m in self.cluster.members - if not failover or not failover.candidate or m.name == failover.candidate] + if self.sync_mode_is_active() and not self.cluster.sync.matches(node.name): + return False + # Don't spend time on "nofailover" nodes checking. + # We also don't need nodes which we can't query with the api in the list. + return node.name not in exclude and \ + not node.nofailover and bool(node.api_url) and \ + (not failover or not failover.candidate or node.name == failover.candidate) + + return list(filter(is_eligible, self.cluster.members)) diff --git a/tests/test_ha.py b/tests/test_ha.py index d3a75ed3..6eb38204 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -748,7 +748,6 @@ class TestHa(PostgresInit): self.p.is_leader = true self.ha.has_lock = true self.ha.is_synchronous_mode = true - self.ha.is_failover_possible = false self.ha.process_sync_replication = Mock() self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, self.p.name, 'a', None), (self.p.name, None)) self.assertEqual('no action. I am (postgresql0), the leader with the lock', self.ha.run_cycle())