diff --git a/patroni/postgresql/slots.py b/patroni/postgresql/slots.py index d08a09c3..fde81a7a 100644 --- a/patroni/postgresql/slots.py +++ b/patroni/postgresql/slots.py @@ -130,7 +130,7 @@ class SlotsHandler(object): self._postgresql = postgresql self._advance = None self._replication_slots: Dict[str, Dict[str, Any]] = {} # already existing replication slots - self._unready_logical_slots: Dict[str, Optional[int]] = {} + self._logical_slots_processing_queue: Dict[str, Optional[int]] = {} self.pg_replslot_dir = os.path.join(self._postgresql.data_dir, 'pg_replslot') self.schedule() @@ -190,7 +190,8 @@ class SlotsHandler(object): self._replication_slots = replication_slots self._schedule_load_slots = False if self._force_readiness_check: - self._unready_logical_slots = {n: None for n, v in replication_slots.items() if v['type'] == 'logical'} + self._logical_slots_processing_queue = {n: None for n, v in replication_slots.items() + if v['type'] == 'logical'} self._force_readiness_check = False def ignore_replication_slot(self, cluster: Cluster, name: str) -> bool: @@ -327,10 +328,10 @@ class SlotsHandler(object): self._ensure_physical_slots(slots) if self._postgresql.is_leader(): - self._unready_logical_slots.clear() + self._logical_slots_processing_queue.clear() self._ensure_logical_slots_primary(slots) elif cluster.slots and slots: - self.check_logical_slots_readiness(cluster, nofailover, replicatefrom) + self.check_logical_slots_readiness(cluster, replicatefrom) ret = self._ensure_logical_slots_replica(cluster, slots) @@ -347,48 +348,100 @@ class SlotsHandler(object): with get_connection_cursor(connect_timeout=3, options="-c statement_timeout=2000", **conn_kwargs) as cur: yield cur - def check_logical_slots_readiness(self, cluster: Cluster, nofailover: bool, replicatefrom: Optional[str]) -> None: + def check_logical_slots_readiness(self, cluster: Cluster, replicatefrom: Optional[str]) -> bool: + """Determine whether all known logical slots are synchronised from the leader. + + 1) Retrieve the current ``catalog_xmin`` value for the physical slot from the cluster leader, and + 2) using previously stored list of "unready" logical slots, those which have yet to be checked hence have no + stored slot attributes, + 3) store logical slot ``catalog_xmin`` when the physical slot ``catalog_xmin`` becomes valid. + + :param cluster: object containing stateful information for the cluster. + :param replicatefrom: name of the member that should be used to replicate from. + + :returns: ``False`` if any issue while checking logical slots readiness, ``True`` otherwise. + """ catalog_xmin = None - if self._unready_logical_slots and cluster.leader: + if self._logical_slots_processing_queue and cluster.leader: slot_name = cluster.get_my_slot_name_on_primary(self._postgresql.name, replicatefrom) try: with self._get_leader_connection_cursor(cluster.leader) as cur: cur.execute("SELECT slot_name, catalog_xmin FROM pg_catalog.pg_get_replication_slots()" " WHERE NOT pg_catalog.pg_is_in_recovery() AND slot_name = ANY(%s)", - ([n for n, v in self._unready_logical_slots.items() if v is None] + [slot_name],)) + ([n for n, v in self._logical_slots_processing_queue.items() + if v is None] + [slot_name],)) slots = {row[0]: row[1] for row in cur} if slot_name not in slots: - return logger.warning('Physical slot %s does not exist on the primary', slot_name) + logger.warning('Physical slot %s does not exist on the primary', slot_name) + return False catalog_xmin = slots.pop(slot_name) except Exception as e: - return logger.error("Failed to check %s physical slot on the primary: %r", slot_name, e) - # Remember catalog_xmin of logical slots on the primary when catalog_xmin of - # the physical slot became valid. Logical slots on replica will be safe to use after - # promote when catalog_xmin of the physical slot overtakes these values. - if catalog_xmin is not None: - for name, value in slots.items(): - self._unready_logical_slots[name] = value - else: # Replica isn't streaming or the hot_standby_feedback isn't enabled - try: - cur = self._query("SELECT pg_catalog.current_setting('hot_standby_feedback')::boolean") - row = cur.fetchone() - if row and not row[0]: - logger.error('Logical slot failover requires "hot_standby_feedback".' - ' Please check postgresql.auto.conf') - except Exception as e: - logger.error('Failed to check the hot_standby_feedback setting: %r', e) - return # since `catalog_xmin` isn't valid further checks don't make any sense + logger.error("Failed to check %s physical slot on the primary: %r", slot_name, e) + return False - for name in list(self._unready_logical_slots): - value = self._replication_slots.get(name) - # The logical slot on a replica is safe to use when the physical replica slot on the primary: - # 1. has a nonzero/non-null catalog_xmin - # 2. has a catalog_xmin that is not newer (greater) than the catalog_xmin of any slot on the standby - # 3. overtook the catalog_xmin of remembered values of logical slots on the primary. - if not value or catalog_xmin is not None and\ - self._unready_logical_slots[name] <= catalog_xmin <= value['catalog_xmin']: - del self._unready_logical_slots[name] - if value: + if not self._update_pending_logical_slot_primary(slots, catalog_xmin): + return False # since `catalog_xmin` isn't valid further checks don't make any sense + + self._ready_logical_slots(catalog_xmin) + return True + + def _update_pending_logical_slot_primary(self, slots: Dict[str, Any], catalog_xmin: Optional[int] = None) -> bool: + """Store pending logical slot information for ``catalog_xmin`` on the primary. + + Remember ``catalog_xmin`` of logical slots on the primary when ``catalog_xmin`` of the physical slot became + valid. Logical slots on replica will be safe to use after promote when ``catalog_xmin`` of the physical slot + overtakes these values. + + :param slots: dictionary of slot information from the primary + :param catalog_xmin: ``catalog_xmin`` of the physical slot used by this replica to stream changes from primary. + + :returns: ``False`` if any issue was faced while processing, ``True`` otherwise. + """ + if catalog_xmin is not None: + for name, value in slots.items(): + self._logical_slots_processing_queue[name] = value + return True + + # Replica isn't streaming or the hot_standby_feedback isn't enabled + try: + cur = self._query("SELECT pg_catalog.current_setting('hot_standby_feedback')::boolean") + row = cur.fetchone() + if row and not row[0]: + logger.error('Logical slot failover requires "hot_standby_feedback".' + ' Please check postgresql.auto.conf') + except Exception as e: + logger.error('Failed to check the hot_standby_feedback setting: %r', e) + return False + + def _ready_logical_slots(self, primary_physical_catalog_xmin: Optional[int] = None) -> None: + """Ready logical slots by comparing primary physical slot ``catalog_xmin`` to logical ``catalog_xmin``. + + The logical slot on a replica is safe to use when the physical replica slot on the primary: + + 1. has a nonzero/non-null ``catalog_xmin`` represented by ``primary_physical_xmin``. + 2. has a ``catalog_xmin`` that is not newer (greater) than the ``catalog_xmin`` of any slot on the standby + 3. overtook the ``catalog_xmin`` of remembered values of logical slots on the primary. + + :param primary_physical_catalog_xmin: is the value retrieved from ``pg_catalog.pg_get_replication_slots()`` for + the physical replication slot on the primary. + """ + # Make a copy of processing queue keys as a list as the queue dictionary is modified inside the loop. + for name in list(self._logical_slots_processing_queue): + primary_logical_catalog_xmin = self._logical_slots_processing_queue[name] + standby_logical_slot = self._replication_slots.get(name, {}) + standby_logical_catalog_xmin = standby_logical_slot.get('catalog_xmin', 0) + if TYPE_CHECKING: # pragma: no cover + assert primary_logical_catalog_xmin is not None + + if ( + not standby_logical_slot + or primary_physical_catalog_xmin is not None + and primary_logical_catalog_xmin <= primary_physical_catalog_xmin <= standby_logical_catalog_xmin + ): + + del self._logical_slots_processing_queue[name] + + if standby_logical_slot: logger.info('Logical slot %s is safe to be used after a failover', name) def copy_logical_slots(self, cluster: Cluster, create_slots: List[str]) -> None: @@ -433,7 +486,7 @@ class SlotsHandler(object): shutil.rmtree(slot_dir) os.rename(slot_tmp_dir, slot_dir) fsync_dir(slot_dir) - self._unready_logical_slots[name] = None + self._logical_slots_processing_queue[name] = None fsync_dir(self._postgresql.slots_handler.pg_replslot_dir) self._postgresql.start() @@ -446,6 +499,6 @@ class SlotsHandler(object): if self._advance: self._advance.on_promote() - if self._unready_logical_slots: + if self._logical_slots_processing_queue: logger.warning('Logical replication slots that might be unsafe to use after promote: %s', - set(self._unready_logical_slots)) + set(self._logical_slots_processing_queue)) diff --git a/tests/test_slots.py b/tests/test_slots.py index a962dbe3..3f21998f 100644 --- a/tests/test_slots.py +++ b/tests/test_slots.py @@ -93,7 +93,7 @@ class TestSlotsHandler(BaseTestPostgresql): def test__ensure_logical_slots_replica(self): self.p.set_role('replica') self.cluster.slots['ls'] = 12346 - with patch.object(SlotsHandler, 'check_logical_slots_readiness', Mock()): + with patch.object(SlotsHandler, 'check_logical_slots_readiness', Mock(return_value=False)): self.assertEqual(self.s.sync_replication_slots(self.cluster, False), []) self.s._schedule_load_slots = False with patch.object(MockCursor, 'execute', Mock(side_effect=psycopg.OperationalError)),\ @@ -121,12 +121,12 @@ class TestSlotsHandler(BaseTestPostgresql): self.s.copy_logical_slots(self.cluster, ['ls']) with patch.object(MockCursor, '__iter__', Mock(return_value=iter([('postgresql0', None)]))),\ patch.object(MockCursor, 'fetchone', Mock(side_effect=Exception)): - self.assertIsNone(self.s.check_logical_slots_readiness(self.cluster, False, None)) + self.assertFalse(self.s.check_logical_slots_readiness(self.cluster, None)) with patch.object(MockCursor, '__iter__', Mock(return_value=iter([('postgresql0', None)]))),\ patch.object(MockCursor, 'fetchone', Mock(return_value=(False,))): - self.assertIsNone(self.s.check_logical_slots_readiness(self.cluster, False, None)) + self.assertFalse(self.s.check_logical_slots_readiness(self.cluster, None)) with patch.object(MockCursor, '__iter__', Mock(return_value=iter([('ls', 100)]))): - self.s.check_logical_slots_readiness(self.cluster, False, None) + self.s.check_logical_slots_readiness(self.cluster, None) @patch.object(Postgresql, 'stop', Mock(return_value=True)) @patch.object(Postgresql, 'start', Mock(return_value=True))