diff --git a/patroni/ha.py b/patroni/ha.py index 68e07bb2..3b4483fc 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -572,12 +572,33 @@ class Ha(object): """:returns: `True` if synchronous replication is requested.""" return self.global_config.is_synchronous_mode + def is_synchronous_mode_active(self) -> bool: + """:returns: `True` is synchronous replication requested is active (/sync key has a valid "leader" field).""" + return self.is_synchronous_mode() and not self.cluster.sync.is_empty + + def is_quorum_commit_mode(self) -> bool: + """:returns: `True` if quorum commit replication is requested and "supported".""" + return self.global_config.is_quorum_commit_mode and self.state_handler.supports_multiple_sync + + def is_quorum_commit_mode_active(self) -> bool: + """:returns: `True` if quorum replication is requested and active (/sync key has a valid "leader" field).""" + return self.is_quorum_commit_mode() and not self.cluster.sync.is_empty + def is_failsafe_mode(self) -> bool: """:returns: `True` if failsafe_mode is enabled in global configuration.""" return self.global_config.check_mode('failsafe_mode') - def process_sync_replication(self) -> None: - """Process synchronous standby beahvior. + def disable_synchronous_replication(self) -> None: + """Cleans up /sync key in DCS if synchronous replication is disabled.""" + if not self.cluster.sync.is_empty and self.dcs.delete_sync_state(index=self.cluster.sync.index): + logger.info("Disabled synchronous replication") + self.state_handler.sync_handler.set_synchronous_standby_names(CaseInsensitiveSet()) + + def _process_quorum_replication(self) -> None: + pass + + def _process_multisync_replication(self) -> None: + """Process synchronous replication state with one or more sync standbys. Synchronous standbys are registered in two places postgresql.conf and DCS. The order of updating them must be right. The invariant that should be kept is that if a node is primary and sync_standby is set in DCS, @@ -585,44 +606,89 @@ class Ha(object): and then in DCS. When removing, first remove in DCS, then in postgresql.conf. This is so we only consider promoting standbys that were guaranteed to be replicating synchronously. """ - if self.is_synchronous_mode(): - current_state = self.state_handler.sync_handler.current_state(self.cluster) - picked = current_state.active - allow_promote = current_state.sync - voters = CaseInsensitiveSet(self.cluster.sync.voters) + current_state = self.state_handler.sync_handler.current_state(self.cluster) + picked = current_state.active + allow_promote = current_state.sync + voters = CaseInsensitiveSet(self.cluster.sync.voters) - if picked != voters: - sync = self.cluster.sync - # update synchronous standby list in dcs temporarily to point to common nodes in current and picked - sync_common = voters & allow_promote - if sync_common != voters: - logger.info("Updating synchronous privilege temporarily from %s to %s", - list(voters), list(sync_common)) - sync = self.dcs.write_sync_state(self.state_handler.name, sync_common, 0, index=sync.index) - if not sync: - return logger.info('Synchronous replication key updated by someone else.') + if picked == voters: + return - # When strict mode and no suitable replication connections put "*" to synchronous_standby_names - if self.global_config.is_synchronous_mode_strict and not picked: - picked = CaseInsensitiveSet('*') - logger.warning("No standbys available!") + sync = self.cluster.sync - # Update postgresql.conf and wait 2 secs for changes to become active - logger.info("Assigning synchronous standby status to %s", list(picked)) - self.state_handler.sync_handler.set_synchronous_standby_names(picked) + # update synchronous standby list in dcs temporarily to point to common nodes in current and picked + sync_common = voters & allow_promote + if sync_common != voters: + logger.info("Updating synchronous privilege temporarily from %s to %s", + list(voters), list(sync_common)) + sync = self.dcs.write_sync_state(self.state_handler.name, sync_common, 0, index=sync.index) + if not sync: + return logger.info('Synchronous replication key updated by someone else.') - if picked and picked != CaseInsensitiveSet('*') and allow_promote != picked: - # Wait for PostgreSQL to enable synchronous mode and see if we can immediately set sync_standby - time.sleep(2) - allow_promote = self.state_handler.sync_handler.current_state(self.cluster).sync - if allow_promote and allow_promote != sync_common: - if not self.dcs.write_sync_state(self.state_handler.name, allow_promote, 0, index=sync.index): - return logger.info("Synchronous replication key updated by someone else") - logger.info("Synchronous standby status assigned to %s", list(allow_promote)) + # When strict mode and no suitable replication connections put "*" to synchronous_standby_names + if self.global_config.is_synchronous_mode_strict and not picked: + picked = CaseInsensitiveSet('*') + logger.warning("No standbys available!") + + # Update postgresql.conf and wait 2 secs for changes to become active + logger.info("Assigning synchronous standby status to %s", list(picked)) + self.state_handler.sync_handler.set_synchronous_standby_names(picked) + + if picked and picked != CaseInsensitiveSet('*') and allow_promote != picked: + # Wait for PostgreSQL to enable synchronous mode and see if we can immediately set sync_standby + time.sleep(2) + allow_promote = self.state_handler.sync_handler.current_state(self.cluster).sync + + if allow_promote and allow_promote != sync_common: + if self.dcs.write_sync_state(self.state_handler.name, allow_promote, 0, index=sync.index): + logger.info("Synchronous standby status assigned to %s", list(allow_promote)) + else: + logger.info("Synchronous replication key updated by someone else") + + def process_sync_replication(self) -> None: + """Process synchronous replication beahvior on the primary.""" + if self.is_quorum_commit_mode(): + self._process_quorum_replication() + elif self.is_synchronous_mode(): + self._process_multisync_replication() else: - if not self.cluster.sync.is_empty and self.dcs.delete_sync_state(index=self.cluster.sync.index): - logger.info("Disabled synchronous replication") - self.state_handler.sync_handler.set_synchronous_standby_names(CaseInsensitiveSet()) + self.disable_synchronous_replication() + + def process_sync_replication_prepromote(self) -> bool: + """Handle sync replication state before promote. + + If quorum replication is requested and we can keep syncing to enough nodes satisfying the quorum invariant + we can promote immediately and let normal quorum resolver process handle any membership changes later. + Otherwise we will just reset DCS state to ourselves and add replicas as they connect. + + :returns: `True` if on success or `False` if failed to update /sync key in DCS. + """ + if not self.is_synchronous_mode(): + self.disable_synchronous_replication() + return True + + if self.is_quorum_commit_mode_active(): + sync = CaseInsensitiveSet(self.cluster.sync.members) + numsync = len(sync) - self.cluster.sync.quorum - 1 + if self.state_handler.name not in sync: # Node outside voters achieved quorum and got leader + numsync += 1 + else: + sync.discard(self.state_handler.name) + else: + sync = CaseInsensitiveSet() + numsync = self.global_config.min_synchronous_nodes + + if not self.is_quorum_commit_mode() or not self.state_handler.supports_multiple_sync and numsync > 1: + sync = CaseInsensitiveSet() + numsync = self.global_config.min_synchronous_nodes + + # Just set ourselves as the authoritative source of truth for now. We don't want to wait for standbys + # to connect. We will try finding a synchronous standby in the next cycle. + if not self.dcs.write_sync_state(self.state_handler.name, None, 0, index=self.cluster.sync.index): + return False + + self.state_handler.sync_handler.set_synchronous_standby_names(sync, numsync) + return True def is_sync_standby(self, cluster: Cluster) -> bool: """:returns: `True` if the current node is a synchronous standby.""" @@ -726,15 +792,10 @@ class Ha(object): self.process_sync_replication() return message else: - if self.is_synchronous_mode(): - # Just set ourselves as the authoritative source of truth for now. We don't want to wait for standbys - # to connect. We will try finding a synchronous standby in the next cycle. - if not self.dcs.write_sync_state(self.state_handler.name, None, 0, index=self.cluster.sync.index): - # Somebody else updated sync state, it may be due to us losing the lock. To be safe, postpone - # promotion until next cycle. TODO: trigger immediate retry of run_cycle - return 'Postponing promotion because synchronous replication state was updated by somebody else' - self.state_handler.sync_handler.set_synchronous_standby_names( - CaseInsensitiveSet('*') if self.global_config.is_synchronous_mode_strict else CaseInsensitiveSet()) + if not self.process_sync_replication_prepromote(): + # Somebody else updated sync state, it may be due to us losing the lock. To be safe, + # postpone promotion until next cycle. TODO: trigger immediate retry of run_cycle. + return 'Postponing promotion because synchronous replication state was updated by somebody else' if self.state_handler.role not in ('master', 'promoted', 'primary'): def on_success(): self._rewind.reset_state() @@ -751,10 +812,14 @@ class Ha(object): return promote_message def fetch_node_status(self, member: Member) -> _MemberStatus: - """This function perform http get request on member.api_url and fetches its status - :returns: `_MemberStatus` object - """ + """This function performs http get request on member.api_url and fetches its status. + Usually it happens during the leader race and we can't afford wating for a response indefinite time, + therefore the request timeout is hardcoded to 2 seconds, which seems to be a good compromise. + The node which is slow to respond most likely will not be healthy. + + :returns: :class:`_MemberStatus` object + """ try: response = self.patroni.request(member, timeout=2, retries=0) data = response.data.decode('utf-8') @@ -815,17 +880,24 @@ class Ha(object): return all(results) def is_lagging(self, wal_position: int) -> bool: - """Returns if instance with an wal should consider itself unhealthy to be promoted due to replication lag. + """Checks if node should consider itself unhealthy to be promoted due to replication lag. :param wal_position: Current wal position. - :returns True when node is lagging + :returns `True` when node is lagging """ lag = (self.cluster.last_lsn or 0) - wal_position return lag > self.global_config.maximum_lag_on_failover def _is_healthiest_node(self, members: Collection[Member], check_replication_lag: bool = True) -> bool: - """This method tries to determine whether I am healthy enough to became a new leader candidate or not.""" + """This method tries to determine whether the current node is healthy enough to became a new leader candidate. + :param members: the list of nodes to check against + :param check_replication_lag: whether to take the replication lag into account. + If the lag exceeds configured threshold the node disqualifies itself. + :returns: `True` in case if the node is eligible to become the new leader. Since this method is executed + on multiple nodes independently it could happen that many nodes will count themselves as + healthiest because they received/replayed up to the same LSN, but it is totally fine. + """ my_wal_position = self.state_handler.last_operation() if check_replication_lag and self.is_lagging(my_wal_position): logger.info('My wal position exceeds maximum replication lag') @@ -841,8 +913,26 @@ class Ha(object): logger.info('My timeline %s is behind last known cluster timeline %s', my_timeline, cluster_timeline) return False - # 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 self.is_quorum_commit_mode_active(): + quorum = self.cluster.sync.quorum + voting_set = CaseInsensitiveSet(self.cluster.sync.members) + else: + quorum = 0 + voting_set = CaseInsensitiveSet() + + # Prepare list of nodes to run check against. If quorum commit is enabled + # we also include members with nofailover tag if they are listed in voters. + members = [m for m in members if m.name != self.state_handler.name + and m.api_url and (not m.nofailover or m.name in voting_set)] + + # If there is a quorum active then at least one of the quorum contains latest commit. A quorum member saying + # their WAL position is not ahead counts as a vote saying we may become new leader. Note that a node doesn't + # have to be a member of the voting set to gather the necessary votes. + + # Regardless of voting, if we observe a node that can become a leader and is ahead, we defer to that node. + # This can lead to failure to act on quorum if there is asymmetric connectivity. + quorum_votes = 0 if self.state_handler.name in voting_set else -1 + nodes_ahead = 0 if members: for st in self.fetch_nodes_statuses(members): @@ -851,14 +941,24 @@ class Ha(object): logger.warning('Primary (%s) is still alive', st.member.name) return False if my_wal_position < st.wal_position: + nodes_ahead += 1 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\ + if not self.is_synchronous_mode_active()\ or not self.cluster.sync.leader_matches(st.member.name): return False logger.info('Ignoring the former leader being ahead of us') - return True + # we want to count votes only from nodes with postgres up and running! + elif st.member.name in voting_set and st.wal_position > 0: + logger.info('Got quorum vote from %s', st.member.name) + quorum_votes += 1 + + # When not in quorum commit we just want to return `True`. + # In quorum commit the former leader is special and counted healthy even when there are no other nodes. + # Otherwise check that the number of votes exceeds the quorum field from the /sync key. + return not self.is_quorum_commit_mode_active() or quorum_votes >= quorum\ + or nodes_ahead == 0 and self.cluster.sync.leader == self.state_handler.name def is_failover_possible(self, members: List[Member], check_synchronous: Optional[bool] = True, cluster_lsn: Optional[int] = 0) -> bool: @@ -872,8 +972,10 @@ class Ha(object): 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 check_synchronous and self.is_synchronous_mode_active(): + # If quorum commit is requested we want to check all nodes (even not voters), + # because they could get enough votes and reach necessary quorum + 1. + members = [m for m in members if self.is_quorum_commit_mode() or self.cluster.sync.matches(m.name)] if members: for st in self.fetch_nodes_statuses(members): not_allowed_reason = st.failover_limitation() @@ -913,9 +1015,10 @@ class Ha(object): return None return False - # in synchronous mode when our name is not in the /sync key - # we shouldn't take any action even if the candidate is unhealthy - if self.is_synchronous_mode() and not self.cluster.sync.matches(self.state_handler.name, True): + # in synchronous mode (except quorum commit!) when our name is not in the + # /sync key we shouldn't take any action even if the candidate is unhealthy + if self.is_synchronous_mode() and not self.is_quorum_commit_mode()\ + and not self.cluster.sync.matches(self.state_handler.name, True): return False # find specific node and check that it is healthy @@ -937,7 +1040,7 @@ 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) + # exclude me (leader) 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 return False @@ -1003,9 +1106,11 @@ class Ha(object): all_known_members += [RemoteMember(name, {'api_url': url}) for name, url in failsafe_members.items()] 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 not self.cluster.sync.matches(self.state_handler.name, True): + # Special handling if synchronous mode was requested and activated (the leader in /sync is not empty) + if self.is_synchronous_mode_active(): + # In quorum commit mode we allow nodes outside of "voters" to take part in + # the leader race. They just need to get enough votes to `reach quorum + 1`. + if not self.is_quorum_commit_mode() and not self.cluster.sync.matches(self.state_handler.name, True): return False # pick between synchronous candidates so we minimize unnecessary failovers/demotions members = {m.name: m for m in all_known_members if self.cluster.sync.matches(m.name, True)} diff --git a/tests/test_ha.py b/tests/test_ha.py index 59015784..0a1c01d5 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -1206,7 +1206,7 @@ class TestHa(PostgresInit): # When we just became primary nobody is sync self.assertEqual(self.ha.enforce_primary_role('msg', 'promote msg'), 'promote msg') - mock_set_sync.assert_called_once_with(CaseInsensitiveSet()) + mock_set_sync.assert_called_once_with(frozenset(), 0) mock_write_sync.assert_called_once_with('leader', None, 0, index=0) mock_set_sync.reset_mock() @@ -1433,3 +1433,49 @@ class TestHa(PostgresInit): self.assertEqual(self.ha.patroni.request.call_args[1]['timeout'], 2) mock_logger.assert_called() self.assertTrue(mock_logger.call_args[0][0].startswith('Request to Citus coordinator')) + + def test_process_sync_replication_prepromote(self): + self.p._major_version = 90500 + self.ha.cluster = get_cluster_initialized_without_leader(sync=('other', self.p.name + ',foo')) + self.ha.cluster.config.data.update({'synchronous_mode': 'quorum'}) + self.ha.global_config = self.ha.patroni.config.get_global_config(self.ha.cluster) + self.p.is_leader = false + self.p.set_role('replica') + mock_write_sync = self.ha.dcs.write_sync_state = Mock(return_value=None) + # Postgres 9.5, write_sync_state to DCS failed + self.assertEqual(self.ha.run_cycle(), + 'Postponing promotion because synchronous replication state was updated by somebody else') + self.assertEqual(self.ha.dcs.write_sync_state.call_count, 1) + self.assertEqual(mock_write_sync.call_args_list[0][0], (self.p.name, None, 0)) + self.assertEqual(mock_write_sync.call_args_list[0][1], {'index': 0}) + + mock_set_sync = self.p.config.set_synchronous_standby_names = Mock() + mock_write_sync = self.ha.dcs.write_sync_state = Mock(return_value=True) + # Postgres 9.5, our name is written to leader of the /sync key, while voters list and ssn is empty + self.assertEqual(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') + self.assertEqual(self.ha.dcs.write_sync_state.call_count, 1) + self.assertEqual(mock_write_sync.call_args_list[0][0], (self.p.name, None, 0)) + self.assertEqual(mock_write_sync.call_args_list[0][1], {'index': 0}) + self.assertEqual(mock_set_sync.call_count, 1) + self.assertEqual(mock_set_sync.call_args_list[0][0], (None,)) + + self.p._major_version = 90600 + mock_set_sync.reset_mock() + mock_write_sync.reset_mock() + self.p.set_role('replica') + # Postgres 9.6, with quorum commit we avoid updating /sync key and put some nodes to ssn + self.assertEqual(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') + self.assertEqual(mock_write_sync.call_count, 0) + self.assertEqual(mock_set_sync.call_count, 1) + self.assertEqual(mock_set_sync.call_args_list[0][0], ('2 (foo,other)',)) + + self.p._major_version = 150000 + mock_set_sync.reset_mock() + self.p.set_role('replica') + self.p.name = 'nonsync' + self.ha.fetch_node_status = get_node_status() + # Postgres 15, with quorum commit. Non-sync node promoted we avoid updating /sync key and put some nodes to ssn + self.assertEqual(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock') + self.assertEqual(mock_write_sync.call_count, 0) + self.assertEqual(mock_set_sync.call_count, 1) + self.assertEqual(mock_set_sync.call_args_list[0][0], ('ANY 3 (foo,other,postgresql0)',))