From e97d2f09993032b0f0089805be6ec8ac30f79590 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 11 May 2023 12:18:08 +0200 Subject: [PATCH] Implement synchronous_mode=quorum --- patroni/config.py | 2 +- patroni/ha.py | 56 ++++- patroni/quorum.py | 305 ++++++++++++++++++++++++++++ tests/test_ha.py | 65 +++++- tests/test_quorum.py | 473 +++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 898 insertions(+), 3 deletions(-) create mode 100644 patroni/quorum.py create mode 100644 tests/test_quorum.py diff --git a/patroni/config.py b/patroni/config.py index 4486fdd3..63a34b7e 100644 --- a/patroni/config.py +++ b/patroni/config.py @@ -563,7 +563,7 @@ class Config(object): if 'citus' in config: bootstrap = config.setdefault('bootstrap', {}) dcs = bootstrap.setdefault('dcs', {}) - dcs.setdefault('synchronous_mode', True) + dcs.setdefault('synchronous_mode', 'quorum') updated_fields = ( 'name', diff --git a/patroni/ha.py b/patroni/ha.py index 3b4483fc..a8cc8e95 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -20,6 +20,7 @@ from .postgresql.callback_executor import CallbackAction from .postgresql.misc import postgres_version_to_int from .postgresql.postmaster import PostmasterProcess from .postgresql.rewind import Rewind +from .quorum import QuorumStateResolver from .utils import polling_loop, tzutc logger = logging.getLogger(__name__) @@ -595,7 +596,60 @@ class Ha(object): self.state_handler.sync_handler.set_synchronous_standby_names(CaseInsensitiveSet()) def _process_quorum_replication(self) -> None: - pass + """Process synchronous replication state when quorum commit is requested. + + Synchronous standbys are registered in two places postgresql.conf and DCS. The order of updating them must + keep the invariant that `quorum + sync >= len(set(quorum pool)|set(sync pool))`. This is done using + :class:`QuorumStateResolver` that given a current state and set of desired synchronous nodes and replication + level outputs changes to DCS and synchronous replication in correct order to reach the desired state. + In case any of those steps causes an error we can just bail out and let next iteration rediscover the state + and retry necessary transitions. + """ + min_sync = self.global_config.min_synchronous_nodes + sync_wanted = self.global_config.synchronous_node_count + + sync = self.cluster.sync + leader = sync.leader or self.state_handler.name + if sync.is_empty: + sync = self.dcs.write_sync_state(leader, None, 0, index=sync.index) + if not sync: + return logger.warning("Updating sync state failed") + + while True: + transition = 'break' # we need define transition value if `QuorumStateResolver` produced no changes + sync_state = self.state_handler.sync_handler.current_state(self.cluster) + for transition, leader, num, nodes in QuorumStateResolver(leader=leader, + quorum=sync.quorum, + voters=sync.voters, + numsync=sync_state.numsync, + sync=sync_state.sync, + numsync_confirmed=sync_state.numsync_confirmed, + active=sync_state.active, + sync_wanted=sync_wanted, + leader_wanted=self.state_handler.name): + if transition == 'quorum': + logger.info("Setting leader to %s, quorum to %d of %d (%s)", + leader, num, len(nodes), ", ".join(sorted(nodes))) + sync = self.dcs.write_sync_state(leader, nodes, num, index=sync.index) + if not sync: + return logger.info('Synchronous replication key updated by someone else.') + elif transition == 'sync': + logger.info("Setting synchronous replication to %d of %d (%s)", + num, len(nodes), ", ".join(sorted(nodes))) + # Bump up number of num nodes to meet minimum replication factor. Commits will have to wait until + # we have enough nodes to meet replication target. + if num < min_sync: + logger.warning("Replication factor %d requested, but %d synchronous standbys available." + " Commits will be delayed.", min_sync + 1, num) + num = min_sync + self.state_handler.sync_handler.set_synchronous_standby_names(nodes, num) + if transition != 'restart': + break + # synchronous_standby_names was transitioned from empty to non-empty and it may take + # some time for nodes to become synchronous. In this case we want to restart state machine + # hoping that we can update /sync key earlier than in loop_wait seconds. + time.sleep(1) + self.state_handler.reset_cluster_info_state(None) def _process_multisync_replication(self) -> None: """Process synchronous replication state with one or more sync standbys. diff --git a/patroni/quorum.py b/patroni/quorum.py new file mode 100644 index 00000000..a41597b5 --- /dev/null +++ b/patroni/quorum.py @@ -0,0 +1,305 @@ +import logging + +from typing import Collection, Iterator, Optional, Tuple + +from .collections import CaseInsensitiveSet + +logger = logging.getLogger(__name__) + + +class QuorumError(Exception): + pass + + +class QuorumStateResolver(object): + """Calculates a list of state transition tuples of the form `('sync'/'quorum'/'restart',leader,number,set_of_names)` + + Synchronous replication state is set in two places. PostgreSQL configuration sets how many and which nodes are + needed for a commit to succeed, abbreviated as `numsync` and `sync` set here. DCS contains information about how + many and which nodes need to be interrogated to be sure to see an xlog position containing latest confirmed commit, + abbreviated as `quorum` and `voters` set. Both pairs have the meaning "ANY n OF set". + + The number of nodes needed for commit to succeed, `numsync`, is also called the replication factor. + + To guarantee zero lost transactions on failover we need to keep the invariant that at all times any subset of + nodes that can acknowledge a commit overlaps with any subset of nodes that can achieve quorum to promote a new + leader. Given a desired replication factor and a set of nodes able to participate in sync replication there + is one optimal state satisfying this condition. Given the node set `active`, the optimal state is: + + sync = voters = active + numsync = min(sync_wanted, len(active)) + quorum = len(active) - numsync + + We need to be able to produce a series of state changes that take the system to this desired state from any + other state arbitrary given arbitrary changes is node availability, configuration and interrupted transitions. + + To keep the invariant the rule to follow is that when increasing `numsync` or `quorum`, we need to perform the + increasing operation first. When decreasing either, the decreasing operation needs to be performed later. + + Order of adding or removing nodes from sync and voters depends on the state of synchronous_standby_names: + When adding new nodes: + if sync (synchronous_standby_names) is empty: + add new nodes first to sync and then to voters when numsync_confirmed > 0 + else: + add new nodes first to voters and than to sync + When removing nodes: + if sync (synchronous_standby_names) will become empty after removal: + first remove nodes from voters and than from sync + else: + first remove nodes from sync and than from voters. make voters empty if numsync_confirmed == 0""" + + def __init__(self, leader: str, quorum: int, voters: Collection[str], + numsync: int, sync: Collection[str], numsync_confirmed: int, + active: Collection[str], sync_wanted: int, leader_wanted: str) -> None: + self.leader = leader # The leader according to the `/sync` key + self.quorum = quorum # The number of nodes we need to check when doing leader race + self.voters = CaseInsensitiveSet(voters) # Set of nodes we need to check (both stored in the /sync key) + self.numsync = min(numsync, len(sync)) # The number of sync nodes in synchronous_standby_names + self.sync = CaseInsensitiveSet(sync) # Set of nodes in synchronous_standby_names + # The number of nodes that are confirmed to reach safe LSN after adding them to `synchronous_standby_names`. + # We don't list them because it is known that they are always included into active. + self.numsync_confirmed = numsync_confirmed + self.active = CaseInsensitiveSet(active) # Set of active nodes from `pg_stat_replication` + self.sync_wanted = sync_wanted # The desired number of sync nodes + self.leader_wanted = leader_wanted # The desired leader + + def check_invariants(self) -> None: + """Checks invatiant of synchronous_standby_names and /sync key in DCS. + + :raises `QuorumError`: in case of broken state""" + voters = CaseInsensitiveSet(self.voters | CaseInsensitiveSet([self.leader])) + sync = CaseInsensitiveSet(self.sync | CaseInsensitiveSet([self.leader_wanted])) + + # We need to verify that subset of nodes that can acknowledge a commit overlaps + # with any subset of nodes that can achieve quorum to promote a new leader. + if self.voters and not (len(voters | sync) <= self.quorum + self.numsync + 1): + raise QuorumError("Quorum and sync not guaranteed to overlap: nodes %d >= quorum %d + sync %d" % + (len(voters | sync), self.quorum, self.numsync)) + # unstable cases, we are changing synchronous_standby_names and /sync key + # one after another, hence one set is allowed to be a subset of another + if not (voters.issubset(sync) or sync.issubset(voters)): + raise QuorumError("Mismatched sets: quorum only=%s sync only=%s" % + (voters - sync, sync - voters)) + + def quorum_update(self, quorum: int, voters: CaseInsensitiveSet, leader: Optional[str] = None, + adjust_quorum: Optional[bool] = True) -> Iterator[Tuple[str, str, int, CaseInsensitiveSet]]: + """Updates quorum, voters and optionally leader fields. + + :param quorum: the new value for `self.quorum`, could be adjusted depending + on values of `self.numsync_confirmed` and `adjust_quorum` + :param voters: the new value for `self.voters`, could be adjusted if numsync_confirmed == 0 + :param leader: the new value for `self.leader`, optional + :param adjust_quorum: if set to `True` the quorum requirement will be increased by the + difference between `self.numsync` and ``self.numsync_confirmed` + :rtype: Iterator[tuple(type, leader, quorum, voters)] with the new quorum state, + where type could be 'quorum' or 'restart'. The latter means that + quorum could not be updated with the current input data + and the :class:`QuorumStateResolver` should be restarted. + :raises `QuorumError`: in case of invalid data or if invariant after transition could not be satisfied + """ + if quorum < 0: + raise QuorumError("Quorum %d < 0 of (%s)" % (quorum, voters)) + if quorum > 0 and quorum >= len(voters): + raise QuorumError("Quorum %d >= N of (%s)" % (quorum, voters)) + + old_leader = self.leader + if leader is not None: # Change of leader was requested + self.leader = leader + elif self.numsync_confirmed == 0: + # If there are no nodes that known to caught up with the primary we want to reset quorum/votes in /sync key + quorum = 0 + voters = CaseInsensitiveSet() + elif adjust_quorum: + # It could be that the number of nodes that are known to catch up with the primary is below desired numsync. + # We want to increase quorum to guaranty that the sync node will be found during the leader race. + quorum += max(self.numsync - self.numsync_confirmed, 0) + + if (self.leader, quorum, voters) == (old_leader, self.quorum, self.voters): + if self.voters: + return + # If transition produces no change of leader/quorum/voters we want to give a hint to + # the caller to fetch the new state from the database and restart QuorumStateResolver. + yield 'restart', self.leader, self.quorum, self.voters + + self.quorum = quorum + self.voters = voters + self.check_invariants() + logger.debug('quorum %s %s %s', self.leader, self.quorum, self.voters) + yield 'quorum', self.leader, self.quorum, self.voters + + def sync_update(self, numsync: int, sync: CaseInsensitiveSet) -> Iterator[Tuple[str, str, int, CaseInsensitiveSet]]: + """Updates numsync and sync fields. + + :param numsync: the new value for `self.numsync` + :param sync: the new value for `self.sync` + :rtype: Iterator[tuple('sync', leader, numsync, sync)] with the new state of synchronous_standby_names + :raises `QuorumError`: in case of invalid data or if invariant after transition could not be satisfied + """ + if numsync < 0: + raise QuorumError("Sync %d < 0 of (%s)" % (numsync, sync)) + if numsync > len(sync): + raise QuorumError("Sync %s > N of (%s)" % (numsync, sync)) + + self.numsync = numsync + self.sync = sync + self.check_invariants() + logger.debug('sync %s %s %s', self.leader, self.numsync, self.sync) + yield 'sync', self.leader, self.numsync, self.sync + + def __iter__(self) -> Iterator[Tuple[str, str, int, CaseInsensitiveSet]]: + transitions = list(self._generate_transitions()) + # Merge 2 transitions of the same type to a single one. This is always safe because skipping the first + # transition is equivalent to no one observing the intermediate state. + for cur_transition, next_transition in zip(transitions, transitions[1:] + [None]): + if next_transition and cur_transition[0] == next_transition[0]: + continue + yield cur_transition + if cur_transition[0] == 'restart': + break + + def _generate_transitions(self) -> Iterator[Tuple[str, str, int, CaseInsensitiveSet]]: + logger.debug("Quorum state: leader %s quorum %s, voters %s, numsync %s, sync %s, " + "numsync_confirmed %s, active %s, sync_wanted %s leader_wanted %s", + self.leader, self.quorum, self.voters, self.numsync, self.sync, + self.numsync_confirmed, self.active, self.sync_wanted, self.leader_wanted) + try: + if self.leader_wanted != self.leader: + voters = (self.voters - CaseInsensitiveSet([self.leader_wanted])) | CaseInsensitiveSet([self.leader]) + if not self.sync: + # If sync is empty we need to update synchronous_standby_names first + numsync = len(voters) - self.quorum + yield from self.sync_update(numsync, CaseInsensitiveSet(voters)) + # If leader changed we need to add the old leader to quorum (voters) + yield from self.quorum_update(self.quorum, CaseInsensitiveSet(voters), self.leader_wanted) + # right after promote there could be no replication connections yet + if not self.sync & self.active: + return # give another loop_wait seconds for replicas to reconnect before removing them from quorum + else: + self.check_invariants() + except QuorumError as e: + logger.warning('%s', e) + yield from self.quorum_update(len(self.sync) - self.numsync, self.sync) + + assert self.leader == self.leader_wanted + + # numsync_confirmed could be 0 after restart/failover, we will calculate it from quorum + if self.numsync_confirmed == 0 and self.sync & self.active: + self.numsync_confirmed = min(len(self.sync & self.active), len(self.voters) - self.quorum) + logger.debug('numsync_confirmed=0, adjusting it to %d', self.numsync_confirmed) + + # Handle non steady state cases + if self.sync < self.voters: + logger.debug("Case 1: synchronous_standby_names subset of DCS state") + # Case 1: quorum is superset of sync nodes. In the middle of changing quorum. + # Evict from quorum dead nodes that are not being synced. + remove_from_quorum = self.voters - (self.sync | self.active) + if remove_from_quorum: + yield from self.quorum_update( + quorum=len(self.voters) - len(remove_from_quorum) - self.numsync, + voters=CaseInsensitiveSet(self.voters - remove_from_quorum), + adjust_quorum=not (self.sync - self.active)) + # Start syncing to nodes that are in quorum and alive + add_to_sync = (self.voters & self.active) - self.sync + if add_to_sync: + yield from self.sync_update(self.numsync, CaseInsensitiveSet(self.sync | add_to_sync)) + elif self.sync > self.voters: + logger.debug("Case 2: synchronous_standby_names superset of DCS state") + # Case 2: sync is superset of quorum nodes. In the middle of changing replication factor. + # Add to quorum voters nodes that are already synced and active + add_to_quorum = (self.sync - self.voters) & self.active + if add_to_quorum: + voters = CaseInsensitiveSet(self.voters | add_to_quorum) + yield from self.quorum_update(len(voters) - self.numsync, voters) + # Remove from sync nodes that are dead + remove_from_sync = self.sync - self.voters + if remove_from_sync: + yield from self.sync_update( + numsync=min(self.numsync, len(self.sync) - len(remove_from_sync)), + sync=CaseInsensitiveSet(self.sync - remove_from_sync)) + + # After handling these two cases quorum and sync must match. + assert self.voters == self.sync + + safety_margin = self.quorum + min(self.numsync, self.numsync_confirmed) - len(self.voters | self.sync) + if safety_margin > 0: # In the middle of changing replication factor. + if self.numsync > self.sync_wanted: + logger.debug('Case 3: replication factor is bigger than needed') + yield from self.sync_update(max(self.sync_wanted, len(self.voters) - self.quorum), self.sync) + else: + logger.debug('Case 4: quorum is bigger than needed') + yield from self.quorum_update(len(self.sync) - self.numsync, self.voters) + else: + safety_margin = self.quorum + self.numsync - len(self.voters | self.sync) + if self.numsync == self.sync_wanted and safety_margin > 0 and self.numsync > self.numsync_confirmed: + yield from self.quorum_update(len(self.sync) - self.numsync, self.voters) + + # We are in a steady state point. Find if desired state is different and act accordingly. + + # If any nodes have gone away, evict them + to_remove = self.sync - self.active + if to_remove and self.sync == to_remove: + logger.debug("Removing nodes: %s", to_remove) + yield from self.quorum_update(0, CaseInsensitiveSet(), adjust_quorum=False) + yield from self.sync_update(0, CaseInsensitiveSet()) + elif to_remove: + logger.debug("Removing nodes: %s", to_remove) + can_reduce_quorum_by = self.quorum + # If we can reduce quorum size try to do so first + if can_reduce_quorum_by: + # Pick nodes to remove by sorted order to provide deterministic behavior for tests + remove = CaseInsensitiveSet(sorted(to_remove, reverse=True)[:can_reduce_quorum_by]) + sync = CaseInsensitiveSet(self.sync - remove) + # when removing nodes from sync we can safely increase numsync if requested + numsync = min(self.sync_wanted, len(sync)) if self.sync_wanted > self.numsync else self.numsync + yield from self.sync_update(numsync, sync) + voters = CaseInsensitiveSet(self.voters - remove) + to_remove &= self.sync + yield from self.quorum_update(len(voters) - self.numsync, voters, + adjust_quorum=not to_remove) + if to_remove: + assert self.quorum == 0 + numsync = self.numsync - len(to_remove) + sync = CaseInsensitiveSet(self.sync - to_remove) + voters = CaseInsensitiveSet(self.voters - to_remove) + sync_decrease = numsync - min(self.sync_wanted, len(sync)) + quorum = min(sync_decrease, len(voters) - 1) if sync_decrease else 0 + yield from self.quorum_update(quorum, voters, adjust_quorum=False) + yield from self.sync_update(numsync, sync) + + # If any new nodes, join them to quorum + to_add = self.active - self.sync + if to_add: + # First get to requested replication factor + logger.debug("Adding nodes: %s", to_add) + sync_wanted = min(self.sync_wanted, len(self.sync | to_add)) + increase_numsync_by = sync_wanted - self.numsync + if increase_numsync_by > 0: + if self.sync: + add = CaseInsensitiveSet(sorted(to_add)[:increase_numsync_by]) + increase_numsync_by = len(add) + else: # there is only the leader + add = to_add # and it is safe to add all nodes at once if sync is empty + yield from self.sync_update(self.numsync + increase_numsync_by, CaseInsensitiveSet(self.sync | add)) + voters = CaseInsensitiveSet(self.voters | add) + yield from self.quorum_update(len(voters) - sync_wanted, voters) + to_add -= self.sync + if to_add: + voters = CaseInsensitiveSet(self.voters | to_add) + yield from self.quorum_update(len(voters) - sync_wanted, voters, + adjust_quorum=sync_wanted > self.numsync_confirmed) + yield from self.sync_update(sync_wanted, CaseInsensitiveSet(self.sync | to_add)) + + # Apply requested replication factor change + sync_increase = min(self.sync_wanted, len(self.sync)) - self.numsync + if sync_increase > 0: + # Increase replication factor + logger.debug("Increasing replication factor to %s", self.numsync + sync_increase) + yield from self.sync_update(self.numsync + sync_increase, self.sync) + yield from self.quorum_update(len(self.voters) - self.numsync, self.voters) + elif sync_increase < 0: + # Reduce replication factor + logger.debug("Reducing replication factor to %s", self.numsync + sync_increase) + if self.quorum - sync_increase < len(self.voters): + yield from self.quorum_update(len(self.voters) - self.numsync - sync_increase, self.voters, + adjust_quorum=self.sync_wanted > self.numsync_confirmed) + yield from self.sync_update(self.numsync + sync_increase, self.sync) diff --git a/tests/test_ha.py b/tests/test_ha.py index 0a1c01d5..5a99f456 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -1102,7 +1102,7 @@ class TestHa(PostgresInit): self.ha.demote('immediate') follow.assert_called_once_with(None) - def test_process_sync_replication(self): + def test_process__multisync_replication(self): self.ha.has_lock = true mock_set_sync = self.p.sync_handler.set_synchronous_standby_names = Mock() self.p.name = 'leader' @@ -1479,3 +1479,66 @@ class TestHa(PostgresInit): 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)',)) + + def test_process_quorum_replication(self): + self.p._major_version = 150000 + self.ha.has_lock = true + mock_set_sync = self.p.config.set_synchronous_standby_names = Mock() + self.p.name = 'leader' + + self.ha.cluster.config.data.update({'synchronous_mode': 'quorum'}) + self.ha.global_config = self.ha.patroni.config.get_global_config(self.ha.cluster) + + mock_write_sync = self.ha.dcs.write_sync_state = Mock(return_value=None) + # Test /sync key is attempted to set and failed when missing or invalid + self.p.sync_handler.current_state = Mock(return_value=_SyncState('quorum', 1, 1, CaseInsensitiveSet(['other']), + CaseInsensitiveSet(['other']))) + self.ha.run_cycle() + self.assertEqual(mock_write_sync.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': None}) + self.assertEqual(mock_set_sync.call_count, 0) + + mock_write_sync = self.ha.dcs.write_sync_state = Mock(side_effect=[SyncState.empty(), None]) + # Test /sync key is attempted to set and succeed when missing or invalid + with patch.object(SyncState, 'is_empty', Mock(side_effect=[True, False])): + self.ha.run_cycle() + self.assertEqual(mock_write_sync.call_count, 2) + 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': None}) + self.assertEqual(mock_write_sync.call_args_list[1][0], (self.p.name, CaseInsensitiveSet(['other']), 0)) + self.assertEqual(mock_write_sync.call_args_list[1][1], {'index': None}) + self.assertEqual(mock_set_sync.call_count, 0) + + self.p.sync_handler.current_state = Mock(side_effect=[_SyncState('quorum', 1, 0, CaseInsensitiveSet(['foo']), + CaseInsensitiveSet(['other'])), + _SyncState('quorum', 1, 1, CaseInsensitiveSet(['foo']), + CaseInsensitiveSet(['foo']))]) + mock_write_sync = self.ha.dcs.write_sync_state = Mock(return_value=SyncState(1, 'leader', 'foo', 0)) + self.ha.cluster = get_cluster_initialized_with_leader(sync=('leader', 'foo')) + self.ha.cluster.config.data.update({'synchronous_mode': 'quorum'}) + self.ha.global_config = self.ha.patroni.config.get_global_config(self.ha.cluster) + # Test the sync node is removed from voters, added to ssn + with patch.object(Postgresql, 'synchronous_standby_names', Mock(return_value='other')),\ + patch('time.sleep', Mock()): + self.ha.run_cycle() + self.assertEqual(mock_write_sync.call_count, 1) + self.assertEqual(mock_write_sync.call_args_list[0][0], (self.p.name, CaseInsensitiveSet(), 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], ('ANY 1 (other)',)) + + # Test ANY 1 (*) when synchronous_mode_strict and no nodes available + self.ha.cluster.config.data.update({'synchronous_mode_strict': True}) + self.ha.global_config = self.ha.patroni.config.get_global_config(self.ha.cluster) + self.p.sync_handler.current_state = Mock(return_value=_SyncState('quorum', 1, 0, + CaseInsensitiveSet(['other', 'foo']), + CaseInsensitiveSet())) + mock_write_sync.reset_mock() + mock_set_sync.reset_mock() + self.ha.run_cycle() + self.assertEqual(mock_write_sync.call_count, 1) + self.assertEqual(mock_write_sync.call_args_list[0][0], (self.p.name, CaseInsensitiveSet(), 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], ('ANY 1 (*)',)) diff --git a/tests/test_quorum.py b/tests/test_quorum.py new file mode 100644 index 00000000..8dddb4f1 --- /dev/null +++ b/tests/test_quorum.py @@ -0,0 +1,473 @@ +import unittest + +from typing import List, Set, Tuple + +from patroni.quorum import QuorumStateResolver, QuorumError + + +class QuorumTest(unittest.TestCase): + + def check_state_transitions(self, leader: str, quorum: int, voters: Set[str], numsync: int, sync: Set[str], + numsync_confirmed: int, active: Set[str], sync_wanted: int, leader_wanted: str, + expected: List[Tuple[str, str, int, Set[str]]]) -> None: + kwargs = { + 'leader': leader, 'quorum': quorum, 'voters': voters, + 'numsync': numsync, 'sync': sync, 'numsync_confirmed': numsync_confirmed, + 'active': active, 'sync_wanted': sync_wanted, 'leader_wanted': leader_wanted + } + result = list(QuorumStateResolver(**kwargs)) + self.assertEqual(result, expected) + + # also check interrupted transitions + if len(result) > 0 and result[0][0] != 'restart' and kwargs['leader'] == result[0][1]: + if result[0][0] == 'sync': + kwargs.update(numsync=result[0][2], sync=result[0][3]) + else: + kwargs.update(leader=result[0][1], quorum=result[0][2], voters=result[0][3]) + kwargs['expected'] = expected[1:] + self.check_state_transitions(**kwargs) + + def test_1111(self): + leader = 'a' + + # Add node + self.check_state_transitions(leader=leader, quorum=0, voters=set(), + numsync=0, sync=set(), numsync_confirmed=0, active=set('b'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('sync', leader, 1, set('b')), + ('restart', leader, 0, set()), + ]) + self.check_state_transitions(leader=leader, quorum=0, voters=set(), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('b'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set('b')) + ]) + + self.check_state_transitions(leader=leader, quorum=0, voters=set(), + numsync=0, sync=set(), numsync_confirmed=0, active=set('bcde'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('sync', leader, 2, set('bcde')), + ('restart', leader, 0, set()), + ]) + self.check_state_transitions(leader=leader, quorum=0, voters=set(), + numsync=2, sync=set('bcde'), numsync_confirmed=1, active=set('bcde'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 3, set('bcde')), + ]) + + def test_1222(self): + """2 node cluster""" + leader = 'a' + + # Active set matches state + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('b'), + sync_wanted=2, leader_wanted=leader, expected=[]) + + # Add node by increasing quorum + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('BC'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('quorum', leader, 1, set('bC')), + ('sync', leader, 1, set('bC')), + ]) + + # Add node by increasing sync + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('bc'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('sync', leader, 2, set('bc')), + ('quorum', leader, 1, set('bc')), + ]) + # Reduce quorum after added node caught up + self.check_state_transitions(leader=leader, quorum=1, voters=set('bc'), + numsync=2, sync=set('bc'), numsync_confirmed=2, active=set('bc'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set('bc')), + ]) + + # Add multiple nodes by increasing both sync and quorum + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('BCdE'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('sync', leader, 2, set('bC')), + ('quorum', leader, 3, set('bCdE')), + ('sync', leader, 2, set('bCdE')), + ]) + # Reduce quorum after added nodes caught up + self.check_state_transitions(leader=leader, quorum=3, voters=set('bcde'), + numsync=2, sync=set('bcde'), numsync_confirmed=3, active=set('bcde'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 2, set('bcde')), + ]) + + # Primary is alone + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=1, sync=set('b'), numsync_confirmed=0, active=set(), + sync_wanted=1, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set()), + ('sync', leader, 0, set()), + ]) + + # Swap out sync replica + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=1, sync=set('b'), numsync_confirmed=0, active=set('c'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set()), + ('sync', leader, 1, set('c')), + ('restart', leader, 0, set()), + ]) + # Update quorum when added node caught up + self.check_state_transitions(leader=leader, quorum=0, voters=set(), + numsync=1, sync=set('c'), numsync_confirmed=1, active=set('c'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set('c')), + ]) + + def test_1233(self): + """Interrupted transition from 2 node cluster to 3 node fully sync cluster""" + leader = 'a' + + # Node c went away, transition back to 2 node cluster + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=2, sync=set('bc'), numsync_confirmed=1, active=set('b'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('sync', leader, 1, set('b')), + ]) + + # Node c is available transition to larger quorum set, but not yet caught up. + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=2, sync=set('bc'), numsync_confirmed=1, active=set('bc'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 1, set('bc')), + ]) + + # Add in a new node at the same time, but node c didn't caught up yet + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=2, sync=set('bc'), numsync_confirmed=1, active=set('bcd'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 2, set('bcd')), + ('sync', leader, 2, set('bcd')), + ]) + # All sync nodes caught up, reduce quorum + self.check_state_transitions(leader=leader, quorum=2, voters=set('bcd'), + numsync=2, sync=set('bcd'), numsync_confirmed=3, active=set('bcd'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 1, set('bcd')), + ]) + + # Change replication factor at the same time + self.check_state_transitions(leader=leader, quorum=0, voters=set('b'), + numsync=2, sync=set('bc'), numsync_confirmed=1, active=set('bc'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('quorum', leader, 1, set('bc')), + ('sync', leader, 1, set('bc')), + ]) + + def test_2322(self): + """Interrupted transition from 2 node cluster to 3 node cluster with replication factor 2""" + leader = 'a' + + # Node c went away, transition back to 2 node cluster + self.check_state_transitions(leader=leader, quorum=1, voters=set('bc'), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('b'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set('b')), + ]) + + # Node c is available transition to larger quorum set. + self.check_state_transitions(leader=leader, quorum=1, voters=set('bc'), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('bc'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('sync', leader, 1, set('bc')), + ]) + + # Add in a new node at the same time + self.check_state_transitions(leader=leader, quorum=1, voters=set('bc'), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('bcd'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('sync', leader, 1, set('bc')), + ('quorum', leader, 2, set('bcd')), + ('sync', leader, 1, set('bcd')), + ]) + + # Convert to a fully synced cluster + self.check_state_transitions(leader=leader, quorum=1, voters=set('bc'), + numsync=1, sync=set('b'), numsync_confirmed=1, active=set('bc'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('sync', leader, 2, set('bc')), + ]) + # Reduce quorum after all nodes caught up + self.check_state_transitions(leader=leader, quorum=1, voters=set('bc'), + numsync=2, sync=set('bc'), numsync_confirmed=2, active=set('bc'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set('bc')), + ]) + + def test_3535(self): + leader = 'a' + + # remove nodes + self.check_state_transitions(leader=leader, quorum=2, voters=set('bcde'), + numsync=2, sync=set('bcde'), numsync_confirmed=2, active=set('bc'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('sync', leader, 2, set('bc')), + ('quorum', leader, 0, set('bc')), + ]) + self.check_state_transitions(leader=leader, quorum=2, voters=set('bcde'), + numsync=2, sync=set('bcde'), numsync_confirmed=3, active=set('bcd'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('sync', leader, 2, set('bcd')), + ('quorum', leader, 1, set('bcd')), + ]) + + # remove nodes and decrease sync + self.check_state_transitions(leader=leader, quorum=2, voters=set('bcde'), + numsync=2, sync=set('bcde'), numsync_confirmed=2, active=set('bc'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('sync', leader, 2, set('bc')), + ('quorum', leader, 1, set('bc')), + ('sync', leader, 1, set('bc')), + ]) + self.check_state_transitions(leader=leader, quorum=1, voters=set('bcde'), + numsync=3, sync=set('bcde'), numsync_confirmed=2, active=set('bc'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('sync', leader, 3, set('bcd')), + ('quorum', leader, 1, set('bc')), + ('sync', leader, 1, set('bc')), + ]) + + # Increase replication factor and decrease quorum + self.check_state_transitions(leader=leader, quorum=2, voters=set('bcde'), + numsync=2, sync=set('bcde'), numsync_confirmed=2, active=set('bcde'), + sync_wanted=3, leader_wanted=leader, expected=[ + ('sync', leader, 3, set('bcde')), + ]) + # decrease quorum after more nodes caught up + self.check_state_transitions(leader=leader, quorum=2, voters=set('bcde'), + numsync=3, sync=set('bcde'), numsync_confirmed=3, active=set('bcde'), + sync_wanted=3, leader_wanted=leader, expected=[ + ('quorum', leader, 1, set('bcde')), + ]) + + # Add node with decreasing sync and increasing quorum + self.check_state_transitions(leader=leader, quorum=2, voters=set('bcde'), + numsync=2, sync=set('bcde'), numsync_confirmed=2, active=set('bcdef'), + sync_wanted=1, leader_wanted=leader, expected=[ + # increase quorum by 2, 1 for added node and another for reduced sync + ('quorum', leader, 4, set('bcdef')), + # now reduce replication factor to requested value + ('sync', leader, 1, set('bcdef')), + ]) + + # Remove node with increasing sync and decreasing quorum + self.check_state_transitions(leader=leader, quorum=2, voters=set('bcde'), + numsync=2, sync=set('bcde'), numsync_confirmed=2, active=set('bcd'), + sync_wanted=3, leader_wanted=leader, expected=[ + # node e removed from sync wth replication factor increase + ('sync', leader, 3, set('bcd')), + # node e removed from voters with quorum decrease + ('quorum', leader, 1, set('bcd')), + ]) + + def test_remove_nosync_node(self): + leader = 'a' + self.check_state_transitions(leader=leader, quorum=0, voters=set('bc'), + numsync=2, sync=set('bc'), numsync_confirmed=1, active=set('b'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set('b')), + ('sync', leader, 1, set('b')) + ]) + + def test_swap_sync_node(self): + leader = 'a' + self.check_state_transitions(leader=leader, quorum=0, voters=set('bc'), + numsync=2, sync=set('bc'), numsync_confirmed=1, active=set('bd'), + sync_wanted=2, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set('b')), + ('sync', leader, 2, set('bd')), + ('quorum', leader, 1, set('bd')) + ]) + + def test_promotion(self): + # Beginning stat: 'a' in the primary, 1 of bcd in sync + # a fails, c gets quorum votes and promotes + self.check_state_transitions(leader='a', quorum=2, voters=set('bcd'), + numsync=0, sync=set(), numsync_confirmed=0, active=set(), + sync_wanted=1, leader_wanted='c', expected=[ + ('sync', 'a', 1, set('abd')), # set a and b to sync + ('quorum', 'c', 2, set('abd')), # set c as a leader and move a to voters + # and stop because there are no active nodes + ]) + + # next loop, b managed to reconnect + self.check_state_transitions(leader='c', quorum=2, voters=set('abd'), + numsync=1, sync=set('abd'), numsync_confirmed=0, active=set('b'), + sync_wanted=1, leader_wanted='c', expected=[ + ('sync', 'c', 1, set('b')), # remove a from sync as inactive + ('quorum', 'c', 0, set('b')), # remove a from voters and reduce quorum + ]) + + # alternative reality: next loop, no one reconnected + self.check_state_transitions(leader='c', quorum=2, voters=set('abd'), + numsync=1, sync=set('abd'), numsync_confirmed=0, active=set(), + sync_wanted=1, leader_wanted='c', expected=[ + ('quorum', 'c', 0, set()), + ('sync', 'c', 0, set()), + ]) + + def test_nonsync_promotion(self): + # Beginning state: 1 of bc in sync. e.g. (a primary, ssn = ANY 1 (b c)) + # a fails, d sees b and c, knows that it is in sync and decides to promote. + # We include in sync state former primary increasing replication factor + # and let situation resolve. Node d ssn=ANY 1 (b c) + leader = 'd' + self.check_state_transitions(leader='a', quorum=1, voters=set('bc'), + numsync=0, sync=set(), numsync_confirmed=0, active=set(), + sync_wanted=1, leader_wanted=leader, expected=[ + # Set a, b, and c to sync and increase replication factor + ('sync', 'a', 2, set('abc')), + # Set ourselves as the leader and move the old leader to voters + ('quorum', leader, 1, set('abc')), + # and stop because there are no active nodes + ]) + # next loop, b and c managed to reconnect + self.check_state_transitions(leader=leader, quorum=1, voters=set('abc'), + numsync=2, sync=set('abc'), numsync_confirmed=0, active=set('bc'), + sync_wanted=1, leader_wanted=leader, expected=[ + ('sync', leader, 2, set('bc')), # Remove a from being synced to. + ('quorum', leader, 1, set('bc')), # Remove a from quorum + ('sync', leader, 1, set('bc')), # Can now reduce replication factor back + ]) + # alternative reality: next loop, no one reconnected + self.check_state_transitions(leader=leader, quorum=1, voters=set('abc'), + numsync=2, sync=set('abc'), numsync_confirmed=0, active=set(), + sync_wanted=1, leader_wanted=leader, expected=[ + ('quorum', leader, 0, set()), + ('sync', leader, 0, set()), + ]) + + def test_invalid_states(self): + leader = 'a' + + # Main invariant is not satisfied, system is in an unsafe state + resolver = QuorumStateResolver(leader=leader, quorum=0, voters=set('bc'), + numsync=1, sync=set('bc'), numsync_confirmed=1, + active=set('bc'), sync_wanted=1, leader_wanted=leader) + self.assertRaises(QuorumError, resolver.check_invariants) + self.assertEqual(list(resolver), [ + ('quorum', leader, 1, set('bc')) + ]) + + # Quorum and sync states mismatched, somebody other than Patroni modified system state + resolver = QuorumStateResolver(leader=leader, quorum=1, voters=set('bc'), + numsync=2, sync=set('bd'), numsync_confirmed=1, + active=set('bd'), sync_wanted=1, leader_wanted=leader) + self.assertRaises(QuorumError, resolver.check_invariants) + self.assertEqual(list(resolver), [ + ('quorum', leader, 1, set('bd')), + ('sync', leader, 1, set('bd')), + ]) + self.assertTrue(repr(resolver.sync).startswith('