Implement synchronous_mode=quorum

This commit is contained in:
Alexander Kukushkin
2023-05-11 12:18:08 +02:00
parent f5f0adba14
commit e97d2f0999
5 changed files with 898 additions and 3 deletions
+1 -1
View File
@@ -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',
+55 -1
View File
@@ -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.
+305
View File
@@ -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)
+64 -1
View File
@@ -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 (*)',))
+473
View File
@@ -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('<CaseInsensitiveSet'))
def test_sync_high_quorum_low_safety_margin_high(self):
leader = 'a'
self.check_state_transitions(leader=leader, quorum=2, voters=set('bcdef'),
numsync=4, sync=set('bcdef'), numsync_confirmed=3, active=set('bcdef'),
sync_wanted=2, leader_wanted=leader, expected=[
('quorum', leader, 3, set('bcdef')), # Adjust quorum requirements
('sync', leader, 2, set('bcdef')), # Reduce synchronization
])
def test_quorum_update(self):
resolver = QuorumStateResolver(leader='a', quorum=1, voters=set('bc'), numsync=1, sync=set('bc'),
numsync_confirmed=1, active=set('bc'), sync_wanted=1, leader_wanted='a')
self.assertRaises(QuorumError, list, resolver.quorum_update(-1, set()))
self.assertRaises(QuorumError, list, resolver.quorum_update(1, set()))
def test_sync_update(self):
resolver = QuorumStateResolver(leader='a', quorum=1, voters=set('bc'), numsync=1, sync=set('bc'),
numsync_confirmed=1, active=set('bc'), sync_wanted=1, leader_wanted='a')
self.assertRaises(QuorumError, list, resolver.sync_update(-1, set()))
self.assertRaises(QuorumError, list, resolver.sync_update(1, set()))
def test_remove_nodes_with_decreasing_sync(self):
leader = 'a'
# Remove node with decreasing sync
self.check_state_transitions(leader=leader, quorum=1, voters=set('bcdef'),
numsync=4, sync=set('bcdef'), numsync_confirmed=2, active=set('bcd'),
sync_wanted=2, leader_wanted=leader, expected=[
# node f removed from sync
('sync', leader, 4, set('bcde')),
# nodes e and f removed from voters with quorum decrease
('quorum', leader, 1, set('bcd')),
# node e removed from sync with replication factor decrease
('sync', leader, 2, set('bcd')),
])
# Interrupted state, and node g joined
self.check_state_transitions(leader=leader, quorum=1, voters=set('bcdef'),
numsync=4, sync=set('bcde'), numsync_confirmed=2, active=set('bcdg'),
sync_wanted=2, leader_wanted=leader, expected=[
# remove nodes e and f from voters
('quorum', leader, 1, set('bcd')),
# remove node e from sync and reduce replication factor
('sync', leader, 3, set('bcd')),
# add node g to voters with quorum increase
('quorum', leader, 2, set('bcdg')),
# add node g to sync and reduce replication factor
('sync', leader, 2, set('bcdg')),
])
# node f returned
self.check_state_transitions(leader=leader, quorum=1, voters=set('bcdef'),
numsync=4, sync=set('bcde'), numsync_confirmed=2, active=set('bcdf'),
sync_wanted=2, leader_wanted=leader, expected=[
# replace node e with f in sync
('sync', leader, 4, set('bcdf')),
# remove nodes e from voters with quorum decrease
('quorum', leader, 2, set('bcdf')),
# reduce replication factor as it was requested
('sync', leader, 2, set('bcdf')),
])
# node e returned
self.check_state_transitions(leader=leader, quorum=1, voters=set('bcdef'),
numsync=4, sync=set('bcde'), numsync_confirmed=2, active=set('bcde'),
sync_wanted=2, leader_wanted=leader, expected=[
# remove nodes f from voters with quorum decrease
('quorum', leader, 2, set('bcde')),
# reduce replication factor as it was requested
('sync', leader, 2, set('bcde')),
])
# node b is also lost
self.check_state_transitions(leader=leader, quorum=1, voters=set('bcdef'),
numsync=4, sync=set('bcde'), numsync_confirmed=2, active=set('cd'),
sync_wanted=1, leader_wanted=leader, expected=[
# remove nodes b, e, and f from voters
('quorum', leader, 1, set('cd')),
# remove nodes b and e from sync with replication factor decrease
('sync', leader, 1, set('cd')),
])
def test_empty_ssn(self):
# Beginning stat: 'a' in the primary, 1 of bc in sync
# a fails, c gets quorum votes and promotes
self.check_state_transitions(leader='a', quorum=1, voters=set('bc'),
numsync=1, sync=set(), numsync_confirmed=0, active=set(),
sync_wanted=1, leader_wanted='c', expected=[
('sync', 'a', 1, set('ab')), # remove a from sync as inactive
('quorum', 'c', 1, set('ab')), # 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=1, voters=set('ab'),
numsync=1, sync=set('ab'), 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
])