Build recovery params in a separate method (#1219)

In addition to that try to protect from the case when some recovery parameters are set in one of included files by explicitly setting their value to an empty string on postgres 12.

Simplifies https://github.com/zalando/patroni/pull/1208
This commit is contained in:
Alexander Kukushkin
2019-10-11 20:18:06 +02:00
committed by GitHub
parent 863aed314b
commit f4623c4e8e
3 changed files with 29 additions and 27 deletions
+3 -26
View File
@@ -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
+25
View File
@@ -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)
+1 -1
View File
@@ -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))