Introduce Status class (#2853)

It represents the `/status` key in DCS and makes it easier to introduce new values stored in the `/status` key without need to refactor all DCS implementations.
This commit is contained in:
Alexander Kukushkin
2023-09-14 14:40:44 +02:00
committed by GitHub
parent b31a4d55c9
commit 238b8db91e
10 changed files with 121 additions and 140 deletions
+75 -11
View File
@@ -778,16 +778,71 @@ class TimelineHistory(NamedTuple):
return TimelineHistory(version, value, lines)
class Status(NamedTuple):
"""Immutable object (namedtuple) which represents `/status` key.
Consists of the following fields:
:ivar last_lsn: :class:`int` object containing position of last known leader LSN.
:ivar slots: state of permanent replication slots on the primary in the format: ``{"slot_name": int}``.
"""
last_lsn: int
slots: Optional[Dict[str, int]]
@staticmethod
def empty() -> 'Status':
"""Construct an empty :class:`Status` instance.
:returns: empty :class:`Status` object.
"""
return Status(0, None)
@staticmethod
def from_node(value: Union[str, Dict[str, Any], None]) -> 'Status':
"""Factory method to parse *value* as :class:`Status` object.
:param value: JSON serialized string
:returns: constructed :class:`Status` object.
"""
try:
if isinstance(value, str):
value = json.loads(value)
except Exception:
return Status.empty()
if isinstance(value, int): # legacy
return Status(value, None)
if not isinstance(value, dict):
return Status.empty()
try:
last_lsn = int(value.get('optime', ''))
except Exception:
last_lsn = 0
slots: Union[str, Dict[str, int], None] = value.get('slots')
if isinstance(slots, str):
try:
slots = json.loads(slots)
except Exception:
slots = None
if not isinstance(slots, dict):
slots = None
return Status(last_lsn, slots)
class Cluster(NamedTuple('Cluster',
[('initialize', Optional[str]),
('config', Optional[ClusterConfig]),
('leader', Optional[Leader]),
('last_lsn', int),
('status', Status),
('members', List[Member]),
('failover', Optional[Failover]),
('sync', SyncState),
('history', Optional[TimelineHistory]),
('slots', Optional[Dict[str, int]]),
('failsafe', Optional[Dict[str, str]]),
('workers', Dict[int, 'Cluster'])])):
"""Immutable object (namedtuple) which represents PostgreSQL or Citus cluster.
@@ -801,13 +856,11 @@ class Cluster(NamedTuple('Cluster',
:ivar initialize: shows whether this cluster has initialization key stored in DC or not.
:ivar config: global dynamic configuration, reference to `ClusterConfig` object.
:ivar leader: :class:`Leader` object which represents current leader of the cluster.
:ivar last_lsn: :class:int object containing position of last known leader LSN.
This value is stored in the `/status` key or `/optime/leader` (legacy) key.
:ivar status: :class:`Status` object which represents the `/status` key.
:ivar members: list of:class:` Member` objects, all PostgreSQL cluster members including leader
:ivar failover: reference to :class:`Failover` object.
:ivar sync: reference to :class:`SyncState` object, last observed synchronous replication state.
:ivar history: reference to `TimelineHistory` object.
:ivar slots: state of permanent logical replication slots on the primary in the format: {"slot_name": int}.
:ivar failsafe: failsafe topology. Node is allowed to become the leader only if its name is found in this list.
:ivar workers: dictionary of workers of the Citus cluster, optional. Each key is an :class:`int` representing
the group, and the corresponding value is a :class:`Cluster` instance.
@@ -819,10 +872,20 @@ class Cluster(NamedTuple('Cluster',
kwargs['workers'] = {}
return super(Cluster, cls).__new__(cls, *args, **kwargs)
@property
def last_lsn(self) -> int:
"""Last known leader LSN."""
return self.status.last_lsn
@property
def slots(self) -> Optional[Dict[str, int]]:
"""State of permanent replication slots on the primary in the format: ``{"slot_name": int}``."""
return self.status.slots
@staticmethod
def empty() -> 'Cluster':
"""Produce an empty :class:`Cluster` instance."""
return Cluster(None, None, None, 0, [], None, SyncState.empty(), None, None, None, {})
return Cluster(None, None, None, Status.empty(), [], None, SyncState.empty(), None, None, {})
def is_empty(self):
"""Validate definition of all attributes of this :class:`Cluster` instance.
@@ -845,7 +908,7 @@ class Cluster(NamedTuple('Cluster',
>>> assert bool(cluster) is False
>>> cluster = Cluster(None, None, None, 0, [1, 2, 3], None, SyncState.empty(), None, None, None, {})
>>> cluster = Cluster(None, None, None, Status(0, None), [1, 2, 3], None, SyncState.empty(), None, None, {})
>>> len(cluster)
1
@@ -1147,19 +1210,20 @@ class Cluster(NamedTuple('Cluster',
:Example:
No history provided:
>>> Cluster(0, 0, 0, 0, 0, 0, 0, 0, 0, None, {}).timeline
>>> Cluster(0, 0, 0, Status.empty(), 0, 0, 0, 0, None, {}).timeline
0
Empty history assume timeline is ``1``:
>>> Cluster(0, 0, 0, 0, 0, 0, 0, TimelineHistory.from_node(1, '[]'), 0, None, {}).timeline
>>> Cluster(0, 0, 0, Status.empty(), 0, 0, 0, TimelineHistory.from_node(1, '[]'), None, {}).timeline
1
Invalid history format, a string of ``a``, returns ``0``:
>>> Cluster(0, 0, 0, 0, 0, 0, 0, TimelineHistory.from_node(1, '[["a"]]'), 0, None, {}).timeline
>>> Cluster(0, 0, 0, Status.empty(), 0, 0, 0, TimelineHistory.from_node(1, '[["a"]]'), None, {}).timeline
0
History as a list of strings:
>>> Cluster(0, 0, 0, 0, 0, 0, 0, TimelineHistory.from_node(1, '[["3", "2", "1"]]'), 0, None, {}).timeline
>>> history = TimelineHistory.from_node(1, '[["3", "2", "1"]]')
>>> Cluster(0, 0, 0, Status.empty(), 0, 0, 0, history, None, {}).timeline
4
"""
if self.history:
+4 -19
View File
@@ -15,7 +15,7 @@ from urllib3.exceptions import HTTPError
from urllib.parse import urlencode, urlparse, quote
from typing import Any, Callable, Dict, List, Mapping, NamedTuple, Optional, Union, Tuple, TYPE_CHECKING
from . import AbstractDCS, Cluster, ClusterConfig, Failover, Leader, Member, SyncState, \
from . import AbstractDCS, Cluster, ClusterConfig, Failover, Leader, Member, Status, SyncState, \
TimelineHistory, ReturnFalseException, catch_return_false_exception, citus_group_re
from ..exceptions import DCSError
from ..utils import deep_compare, parse_bool, Retry, RetryFailedError, split_host_port, uri, USER_AGENT
@@ -383,23 +383,8 @@ class Consul(AbstractDCS):
history = history and TimelineHistory.from_node(history['ModifyIndex'], history['Value'])
# get last known leader lsn and slots
status = nodes.get(self._STATUS)
if status:
try:
status = json.loads(status['Value'])
last_lsn = status.get(self._OPTIME)
slots = status.get('slots')
except Exception:
slots = last_lsn = None
else:
last_lsn = nodes.get(self._LEADER_OPTIME)
last_lsn = last_lsn and last_lsn['Value']
slots = None
try:
last_lsn = int(last_lsn or '')
except Exception:
last_lsn = 0
status = nodes.get(self._STATUS) or nodes.get(self._LEADER_OPTIME)
status = Status.from_node(status and status['Value'])
# get list of members
members = [self.member(n) for k, n in nodes.items() if k.startswith(self._MEMBERS) and k.count('/') == 1]
@@ -428,7 +413,7 @@ class Consul(AbstractDCS):
except Exception:
failsafe = None
return Cluster(initialize, config, leader, last_lsn, members, failover, sync, history, slots, failsafe)
return Cluster(initialize, config, leader, status, members, failover, sync, history, failsafe)
@property
def _consistency(self) -> str:
+4 -19
View File
@@ -21,7 +21,7 @@ from urllib.parse import urlparse
from urllib3 import Timeout
from urllib3.exceptions import HTTPError, ReadTimeoutError, ProtocolError
from . import AbstractDCS, Cluster, ClusterConfig, Failover, Leader, Member, SyncState, \
from . import AbstractDCS, Cluster, ClusterConfig, Failover, Leader, Member, Status, SyncState, \
TimelineHistory, ReturnFalseException, catch_return_false_exception, citus_group_re
from ..exceptions import DCSError
from ..request import get as requests_get
@@ -677,23 +677,8 @@ class Etcd(AbstractEtcd):
history = history and TimelineHistory.from_node(history.modifiedIndex, history.value)
# get last know leader lsn and slots
status = nodes.get(self._STATUS)
if status:
try:
status = json.loads(status.value)
last_lsn = status.get(self._OPTIME)
slots = status.get('slots')
except Exception:
slots = last_lsn = None
else:
last_lsn = nodes.get(self._LEADER_OPTIME)
last_lsn = last_lsn and last_lsn.value
slots = None
try:
last_lsn = int(last_lsn or '')
except Exception:
last_lsn = 0
status = nodes.get(self._STATUS) or nodes.get(self._LEADER_OPTIME)
status = Status.from_node(status and status.value)
# get list of members
members = [self.member(n) for k, n in nodes.items() if k.startswith(self._MEMBERS) and k.count('/') == 1]
@@ -722,7 +707,7 @@ class Etcd(AbstractEtcd):
except Exception:
failsafe = None
return Cluster(initialize, config, leader, last_lsn, members, failover, sync, history, slots, failsafe)
return Cluster(initialize, config, leader, status, members, failover, sync, history, failsafe)
def _cluster_loader(self, path: str) -> Cluster:
result = self.retry(self._client.read, path, recursive=True, quorum=self._ctl)
+4 -19
View File
@@ -15,7 +15,7 @@ from urllib3.exceptions import ReadTimeoutError, ProtocolError
from threading import Condition, Lock, Thread
from typing import Any, Callable, Collection, Dict, Iterator, List, Optional, Tuple, Type, TYPE_CHECKING, Union
from . import ClusterConfig, Cluster, Failover, Leader, Member, SyncState, \
from . import ClusterConfig, Cluster, Failover, Leader, Member, Status, SyncState, \
TimelineHistory, catch_return_false_exception, citus_group_re
from .etcd import AbstractEtcdClientWithFailover, AbstractEtcd, catch_etcd_errors, DnsCachingResolver, Retry
from ..exceptions import DCSError, PatroniException
@@ -723,23 +723,8 @@ class Etcd3(AbstractEtcd):
history = history and TimelineHistory.from_node(history['mod_revision'], history['value'])
# get last know leader lsn and slots
status = nodes.get(self._STATUS)
if status:
try:
status = json.loads(status['value'])
last_lsn = status.get(self._OPTIME)
slots = status.get('slots')
except Exception:
slots = last_lsn = None
else:
last_lsn = nodes.get(self._LEADER_OPTIME)
last_lsn = last_lsn and last_lsn['value']
slots = None
try:
last_lsn = int(last_lsn or '')
except Exception:
last_lsn = 0
status = nodes.get(self._STATUS) or nodes.get(self._LEADER_OPTIME)
status = Status.from_node(status and status['value'])
# get list of members
members = [self.member(n) for k, n in nodes.items() if k.startswith(self._MEMBERS) and k.count('/') == 1]
@@ -770,7 +755,7 @@ class Etcd3(AbstractEtcd):
except Exception:
failsafe = None
return Cluster(initialize, config, leader, last_lsn, members, failover, sync, history, slots, failsafe)
return Cluster(initialize, config, leader, status, members, failover, sync, history, failsafe)
def _cluster_loader(self, path: str) -> Cluster:
nodes = {node['key'][len(path):]: node
+4 -14
View File
@@ -19,7 +19,7 @@ from urllib3.exceptions import HTTPError
from threading import Condition, Lock, Thread
from typing import Any, Callable, Collection, Dict, List, Optional, Tuple, Type, Union, TYPE_CHECKING
from . import AbstractDCS, Cluster, ClusterConfig, Failover, Leader, Member, SyncState, \
from . import AbstractDCS, Cluster, ClusterConfig, Failover, Leader, Member, Status, SyncState, \
TimelineHistory, CITUS_COORDINATOR_GROUP_ID, citus_group_re
from ..exceptions import DCSError
from ..utils import deep_compare, iter_response_objects, keepalive_socket_options, \
@@ -888,18 +888,8 @@ class Kubernetes(AbstractDCS):
self._leader_resource_version = metadata.resource_version if metadata else None
annotations: Dict[str, str] = metadata and metadata.annotations or {}
# get last known leader lsn
try:
last_lsn = int(annotations.get(self._OPTIME, ''))
except Exception:
last_lsn = 0
# get permanent slots state (confirmed_flush_lsn)
slots = annotations.get('slots')
try:
slots = json.loads(annotations.get('slots', ''))
except Exception:
slots = None
# get last known leader lsn and slots
status = Status.from_node(annotations)
# get failsafe topology
try:
@@ -945,7 +935,7 @@ class Kubernetes(AbstractDCS):
metadata = sync and sync.metadata
sync = SyncState.from_node(metadata and metadata.resource_version, metadata and metadata.annotations)
return Cluster(initialize, config, leader, last_lsn, members, failover, sync, history, slots, failsafe)
return Cluster(initialize, config, leader, status, members, failover, sync, history, failsafe)
def _cluster_loader(self, path: Dict[str, Any]) -> Cluster:
return self._cluster_from_nodes(path['group'], path['nodes'], path['pods'].values())
+5 -19
View File
@@ -12,7 +12,8 @@ from pysyncobj.transport import TCPTransport, CONNECTION_STATE
from pysyncobj.utility import TcpUtility
from typing import Any, Callable, Collection, Dict, List, Optional, Set, Union, TYPE_CHECKING
from . import AbstractDCS, ClusterConfig, Cluster, Failover, Leader, Member, SyncState, TimelineHistory, citus_group_re
from . import AbstractDCS, ClusterConfig, Cluster, Failover, Leader, Member, Status, SyncState, \
TimelineHistory, citus_group_re
from ..exceptions import DCSError
from ..utils import validate_directory
if TYPE_CHECKING: # pragma: no cover
@@ -343,23 +344,8 @@ class Raft(AbstractDCS):
history = history and TimelineHistory.from_node(history['index'], history['value'])
# get last know leader lsn and slots
status = nodes.get(self._STATUS)
if status:
try:
status = json.loads(status['value'])
last_lsn = status.get(self._OPTIME)
slots = status.get('slots')
except Exception:
slots = last_lsn = None
else:
last_lsn = nodes.get(self._LEADER_OPTIME)
last_lsn = last_lsn and last_lsn['value']
slots = None
try:
last_lsn = int(last_lsn or '')
except Exception:
last_lsn = 0
status = nodes.get(self._STATUS) or nodes.get(self._LEADER_OPTIME)
status = Status.from_node(status and status['value'])
# get list of members
members = [self.member(k, n) for k, n in nodes.items() if k.startswith(self._MEMBERS) and k.count('/') == 1]
@@ -387,7 +373,7 @@ class Raft(AbstractDCS):
except Exception:
failsafe = None
return Cluster(initialize, config, leader, last_lsn, members, failover, sync, history, slots, failsafe)
return Cluster(initialize, config, leader, status, members, failover, sync, history, failsafe)
def _cluster_loader(self, path: str) -> Cluster:
response = self._sync_obj.get(path, recursive=True)
+11 -25
View File
@@ -12,7 +12,8 @@ from kazoo.retry import RetryFailedError
from kazoo.security import ACL, make_acl
from typing import Any, Callable, Dict, List, Optional, Union, Tuple, TYPE_CHECKING
from . import AbstractDCS, ClusterConfig, Cluster, Failover, Leader, Member, SyncState, TimelineHistory, citus_group_re
from . import AbstractDCS, ClusterConfig, Cluster, Failover, Leader, Member, Status, SyncState, \
TimelineHistory, citus_group_re
from ..exceptions import DCSError
from ..utils import deep_compare
if TYPE_CHECKING: # pragma: no cover
@@ -200,29 +201,15 @@ class ZooKeeper(AbstractDCS):
except NoNodeError:
return None
def get_status(self, path: str, leader: Optional[Leader]) -> Tuple[int, Optional[Dict[str, int]]]:
def get_status(self, path: str, leader: Optional[Leader]) -> Status:
watch = self.status_watcher if not leader or leader.name != self._name else None
status = self.get_node(path + self._STATUS, watch)
if not status:
status = self.get_node(path + self._LEADER_OPTIME, watch)
if status:
try:
status = json.loads(status[0])
last_lsn = status.get(self._OPTIME)
slots = status.get('slots')
except Exception:
slots = last_lsn = None
else:
last_lsn = self.get_node(path + self._LEADER_OPTIME, watch)
last_lsn = last_lsn and last_lsn[0]
slots = None
try:
last_lsn = int(last_lsn or '')
except Exception:
last_lsn = 0
self._fetch_status = False
return last_lsn, slots
self._fetch_status = False
return Status.from_node(status and status[0])
@staticmethod
def member(name: str, value: str, znode: ZnodeStat) -> Member:
@@ -276,7 +263,7 @@ class ZooKeeper(AbstractDCS):
self._fetch_cluster = member.version == -1
# get last known leader lsn and slots
last_lsn, slots = self.get_status(path, leader)
status = self.get_status(path, leader)
# failover key
failover = self.get_node(path + self._FAILOVER, watch=self.cluster_watcher) if self._FAILOVER in nodes else None
@@ -289,7 +276,7 @@ class ZooKeeper(AbstractDCS):
except Exception:
failsafe = None
return Cluster(initialize, config, leader, last_lsn, members, failover, sync, history, slots, failsafe)
return Cluster(initialize, config, leader, status, members, failover, sync, history, failsafe)
def _citus_cluster_loader(self, path: str) -> Dict[int, Cluster]:
fetch_cluster = False
@@ -320,11 +307,10 @@ class ZooKeeper(AbstractDCS):
self.event.clear()
else:
try:
last_lsn, slots = self.get_status(self.client_path(''), cluster.leader)
status = self.get_status(self.client_path(''), cluster.leader)
self.event.clear()
new_cluster: List[Any] = list(cluster)
new_cluster[3] = last_lsn
new_cluster[8] = slots
new_cluster[3] = status
cluster = Cluster(*new_cluster)
except Exception:
pass
+3 -2
View File
@@ -14,7 +14,7 @@ from . import psycopg
from .__main__ import Patroni
from .async_executor import AsyncExecutor, CriticalTask
from .collections import CaseInsensitiveSet
from .dcs import AbstractDCS, Cluster, Leader, Member, RemoteMember
from .dcs import AbstractDCS, Cluster, Leader, Member, RemoteMember, Status
from .exceptions import DCSError, PostgresConnectionException, PatroniFatalException
from .postgresql.callback_executor import CallbackAction
from .postgresql.misc import postgres_version_to_int
@@ -123,7 +123,8 @@ class Failsafe(object):
leader = self.leader
if leader:
# We rely on the strict order of fields in the namedtuple
cluster = Cluster(*cluster[0:2], leader, *cluster[3:8], leader.member.data['slots'], *cluster[9:])
status = Status(cluster.status.last_lsn, leader.member.data['slots'])
cluster = Cluster(*cluster[0:2], leader, status, *cluster[4:])
return cluster
def is_active(self) -> bool:
+2 -2
View File
@@ -6,7 +6,7 @@ import sys
from mock import Mock, MagicMock, PropertyMock, patch, mock_open
from patroni.collections import CaseInsensitiveSet
from patroni.config import Config
from patroni.dcs import Cluster, ClusterConfig, Failover, Leader, Member, get_dcs, SyncState, TimelineHistory
from patroni.dcs import Cluster, ClusterConfig, Failover, Leader, Member, get_dcs, Status, SyncState, TimelineHistory
from patroni.dcs.etcd import AbstractEtcdClientWithFailover
from patroni.exceptions import DCSError, PostgresConnectionException, PatroniFatalException
from patroni.ha import Ha, _MemberStatus
@@ -39,7 +39,7 @@ def get_cluster(initialize, leader, members, failover, sync, cluster_config=None
history = TimelineHistory(1, '[[1,67197376,"no recovery target specified","' + t + '","foo"]]',
[(1, 67197376, 'no recovery target specified', t, 'foo')])
cluster_config = cluster_config or ClusterConfig(1, {'check_timeline': True}, 1)
return Cluster(initialize, cluster_config, leader, 10, members, failover, sync, history, None, failsafe)
return Cluster(initialize, cluster_config, leader, Status(10, None), members, failover, sync, history, failsafe)
def get_cluster_not_initialized_without_leader(cluster_config=None):
+9 -10
View File
@@ -8,7 +8,7 @@ from threading import Thread
from patroni import psycopg
from patroni.config import GlobalConfig
from patroni.dcs import Cluster, ClusterConfig, Member, SyncState
from patroni.dcs import Cluster, ClusterConfig, Member, Status, SyncState
from patroni.postgresql import Postgresql
from patroni.postgresql.misc import fsync_dir
from patroni.postgresql.slots import SlotsAdvanceThread, SlotsHandler
@@ -33,15 +33,15 @@ class TestSlotsHandler(BaseTestPostgresql):
self.s = self.p.slots_handler
self.p.start()
config = ClusterConfig(1, {'slots': {'ls': {'database': 'a', 'plugin': 'b'}, 'ls2': None}}, 1)
self.cluster = Cluster(True, config, self.leader, 0, [self.me, self.other, self.leadermem],
None, SyncState.empty(), None, {'ls': 12345, 'ls2': 12345}, None)
self.cluster = Cluster(True, config, self.leader, Status(0, {'ls': 12345, 'ls2': 12345}),
[self.me, self.other, self.leadermem], None, SyncState.empty(), None, None)
def test_sync_replication_slots(self):
config = ClusterConfig(1, {'slots': {'test_3': {'database': 'a', 'plugin': 'b'},
'A': 0, 'ls': 0, 'b': {'type': 'logical', 'plugin': '1'}},
'ignore_slots': [{'name': 'blabla'}]}, 1)
cluster = Cluster(True, config, self.leader, 0, [self.me, self.other, self.leadermem],
None, SyncState.empty(), None, {'test_3': 10}, None)
cluster = Cluster(True, config, self.leader, Status(0, {'test_3': 10}),
[self.me, self.other, self.leadermem], None, SyncState.empty(), None, None)
with mock.patch('patroni.postgresql.Postgresql._query', Mock(side_effect=psycopg.OperationalError)):
self.s.sync_replication_slots(cluster, False)
self.p.set_role('standby_leader')
@@ -77,9 +77,8 @@ class TestSlotsHandler(BaseTestPostgresql):
'state': 'running', 'conn_url': 'postgres://replicator:[email protected]:5436/postgres',
'tags': {'replicatefrom': 'postgresql0'}
})
cluster = Cluster(True, config, self.leader, 0,
[self.me, self.other, self.leadermem, cascading_replica],
None, SyncState.empty(), None, {'ls': 10}, None)
cluster = Cluster(True, config, self.leader, Status(0, {'ls': 10}),
[self.me, self.other, self.leadermem, cascading_replica], None, SyncState.empty(), None, None)
self.p.set_role('replica')
with patch.object(Postgresql, '_query') as mock_query, \
patch.object(Postgresql, 'is_primary', Mock(return_value=False)):
@@ -90,8 +89,8 @@ class TestSlotsHandler(BaseTestPostgresql):
def test_process_permanent_slots(self):
config = ClusterConfig(1, {'slots': {'ls': {'database': 'a', 'plugin': 'b'}},
'ignore_slots': [{'name': 'blabla'}]}, 1)
cluster = Cluster(True, config, self.leader, 0, [self.me, self.other, self.leadermem],
None, SyncState.empty(), None, None, None)
cluster = Cluster(True, config, self.leader, Status.empty(), [self.me, self.other, self.leadermem],
None, SyncState.empty(), None, None)
self.s.sync_replication_slots(cluster, False)
with patch.object(Postgresql, '_query') as mock_query: