diff --git a/patroni/postgresql/__init__.py b/patroni/postgresql/__init__.py index f3654a02..0deb7a16 100644 --- a/patroni/postgresql/__init__.py +++ b/patroni/postgresql/__init__.py @@ -18,7 +18,6 @@ from patroni.postgresql.postmaster import PostmasterProcess from patroni.postgresql.slots import SlotsHandler from patroni.exceptions import PostgresConnectionException from patroni.utils import Retry, RetryFailedError, polling_loop -from patroni.dcs import slot_name_from_member_name, RemoteMember from threading import current_thread, Lock @@ -674,37 +673,15 @@ class Postgresql(object): logger.exception('Failed to read and parse %s', (history_path,)) def follow(self, member, role='replica', timeout=None): - is_remote_master = isinstance(member, RemoteMember) - no_replication_slot = is_remote_master and member.no_replication_slot - restore_command = is_remote_master and member.restore_command - min_apply_delay = is_remote_master and member.recovery_min_apply_delay - archive_cleanup = is_remote_master and member.archive_cleanup_command - - primary_conninfo = self.config.primary_conninfo_params(member) - change_role = self.cb_called and (self.role in ('master', 'demoted') or - not {'standby_leader', 'replica'} - {self.role, role}) - - recovery_params = self.config.get('recovery_conf', {}).copy() - recovery_params.update({'standby_mode': 'on', 'recovery_target_timeline': 'latest'}) - if primary_conninfo: - recovery_params['primary_conninfo'] = primary_conninfo - if self.slots_handler.use_slots and not no_replication_slot: - required_name = is_remote_master and member.data.get('primary_slot_name') - name = required_name or slot_name_from_member_name(self.name) - recovery_params['primary_slot_name'] = name - if restore_command: - recovery_params['restore_command'] = restore_command - if min_apply_delay: - recovery_params['recovery_min_apply_delay'] = min_apply_delay - if archive_cleanup: - recovery_params['archive_cleanup_command'] = archive_cleanup - + recovery_params = self.config.build_recovery_params(member) self.config.write_recovery_conf(recovery_params) # When we demoting the master or standby_leader to replica or promoting replica to a standby_leader # and we know for sure that postgres was already running before, we will only execute on_role_change # callback and prevent execution of on_restart/on_start callback. # If the role remains the same (replica or standby_leader), we will execute on_start or on_restart + change_role = self.cb_called and (self.role in ('master', 'demoted') or + not {'standby_leader', 'replica'} - {self.role, role}) if change_role: self.__cb_pending = ACTION_NOOP diff --git a/patroni/postgresql/config.py b/patroni/postgresql/config.py index 83ecbc48..9bcea92f 100644 --- a/patroni/postgresql/config.py +++ b/patroni/postgresql/config.py @@ -9,6 +9,7 @@ import time from requests.structures import CaseInsensitiveDict from six.moves.urllib_parse import urlparse, parse_qsl, unquote +from ..dcs import slot_name_from_member_name, RemoteMember from ..utils import compare_values, parse_bool, parse_int, split_host_port, uri logger = logging.getLogger(__name__) @@ -485,6 +486,30 @@ class ConfigHandler(object): value = self.format_dsn(value) fd.write_param(name, value) + def build_recovery_params(self, member): + recovery_params = CaseInsensitiveDict({p: v for p, v in self.get('recovery_conf', {}).items() + if not p.lower().startswith('recovery_target') and + p.lower() not in ('primary_conninfo', 'primary_slot_name')}) + recovery_params.update({'standby_mode': 'on', 'recovery_target_timeline': 'latest'}) + if self._postgresql.major_version >= 120000: + # on pg12 we want to protect from following params being set in one of included files + # not doing so might result in a standby being paused, promoted or shutted down. + recovery_params.update({'recovery_target': '', 'recovery_target_name': '', 'recovery_target_time': '', + 'recovery_target_xid': '', 'recovery_target_lsn': ''}) + + is_remote_master = isinstance(member, RemoteMember) + primary_conninfo = self.primary_conninfo_params(member) + if primary_conninfo: + recovery_params['primary_conninfo'] = primary_conninfo + if self._postgresql.slots_handler.use_slots and not (is_remote_master and member.no_replication_slot): + recovery_params['primary_slot_name'] = member.primary_slot_name if is_remote_master \ + else slot_name_from_member_name(self._postgresql.name) + + if is_remote_master: # standby_cluster config might have different parameters, we want to override them + recovery_params.update({p: member.data.get(p) for p in ('restore_command', 'recovery_min_apply_delay', + 'archive_cleanup_command') if member.data.get(p)}) + return recovery_params + def recovery_conf_exists(self): if self._postgresql.major_version >= 120000: return os.path.exists(self._standby_signal) or os.path.exists(self._recovery_signal) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index b4116684..b25532c4 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -258,7 +258,7 @@ class TestPostgresql(BaseTestPostgresql): @patch.object(Postgresql, 'start', Mock()) def test_follow(self): self.p.call_nowait('on_start') - m = RemoteMember('1', {'restore_command': '2', 'recovery_min_apply_delay': 3, 'archive_cleanup_command': '4'}) + m = RemoteMember('1', {'restore_command': '2', 'primary_slot_name': 'foo', 'conn_kwargs': {'host': 'bar'}}) self.p.follow(m) @patch.object(Postgresql, 'is_running', Mock(return_value=True))