From dd7c3c349f64391f7310904d061c178376687470 Mon Sep 17 00:00:00 2001 From: Dmitry Dolgov <9erthalion6@gmail.com> Date: Fri, 7 Sep 2018 10:10:56 +0200 Subject: [PATCH] [WIP] Standby cluster implementation (#679) Implementation of "standby cluster" described in #657. Standby cluster consists of a "standby leader", that replicates from a "remote master" (which is not a part of current patroni cluster and can be anywhere), and cascade replicas, that replicate from the corresponding standby leader. "Standby leader" behaves pretty much like a regular leader, which means that it holds a leader lock in DSC, in case if disappears there will be an election of a new "standby leader". One can define such a cluster using the section "standby_cluster" in patroni config file. This section provides parameters for standby cluster, that will be applied only once during bootstrap and can be changed only through DSC. --- docs/SETTINGS.rst | 8 ++ docs/replica_bootstrap.rst | 32 +++++++ features/basic_replication.feature | 2 +- features/standby_cluster.feature | 16 ++++ features/steps/standby_cluster.py | 78 ++++++++++++++++ patroni/config.py | 34 ++++++- patroni/dcs/__init__.py | 56 ++++++++++- patroni/ha.py | 107 ++++++++++++++++++--- patroni/postgresql.py | 36 ++++++-- postgres0.yml | 4 + tests/test_config.py | 2 +- tests/test_ctl.py | 15 +-- tests/test_ha.py | 144 +++++++++++++++++++++++++++-- tests/test_postgresql.py | 4 +- 14 files changed, 497 insertions(+), 41 deletions(-) create mode 100644 features/standby_cluster.feature create mode 100644 features/steps/standby_cluster.py diff --git a/docs/SETTINGS.rst b/docs/SETTINGS.rst index a7c68604..0ce1ba5b 100644 --- a/docs/SETTINGS.rst +++ b/docs/SETTINGS.rst @@ -25,6 +25,14 @@ Bootstrap configuration - **use\_slots**: whether or not to use replication_slots. Must be False for PostgreSQL 9.3. You should comment out max_replication_slots before it becomes ineligible for leader status. - **recovery\_conf**: additional configuration settings written to recovery.conf when configuring follower. - **parameters**: list of configuration settings for Postgres. Many of these are required for replication to work. + - **standby\_cluster**: if this section is defined, we want to bootstrap a standby cluster. + - **host**: an address of remote master + - **port**: a port of remote master + - **primary\_slot\_name**: which slot on the remote master to use for replication. This parameter is optional, the default value is derived from the instance name (see function `slot_name_from_member_name`). + - **create\_replica\_methods**: an ordered list of methods that can be used to bootstrap standby leader from the remote master, can be different from the list defined in :ref:`postgresql_settings` + - **restore\_command**: command to restore WAL records from the remote master to standby leader, can be different from the list defined in :ref:`postgresql_settings` + - **archive\_cleanup\_command**: cleanup command for standby leader + - **recovery\_min\_apply\_delay**: how long to wait before actually apply WAL records on a standby leader - **method**: custom script to use for bootstrapping this cluster. See :ref:`custom bootstrap methods documentation ` for details. When ``initdb`` is specified revert to the default ``initdb`` command. ``initdb`` is also triggered when no ``method`` diff --git a/docs/replica_bootstrap.rst b/docs/replica_bootstrap.rst index d5d79155..3f832a89 100644 --- a/docs/replica_bootstrap.rst +++ b/docs/replica_bootstrap.rst @@ -129,3 +129,35 @@ and - max-rate: '100M' If all replica creation methods fail, Patroni will try again all methods in order during the next event loop cycle. + +Standby cluster +--------------- + +Another available option is to run a "standby cluster", that contains only of +standby nodes replicating from some remote master. This type of clusters has: + +* "standby leader", that behaves pretty much like a regular cluster leader, + except it replicates from a remote master. + +* cascade replicas, that are replicating from standby leader. + +Standby leader holds and updates a leader lock in DCS. If the leader lock +expires, cascade replicas will perform an election to choose another leader +from the standbys. For the sake of flexibility, you can specify different +methods of creating a replica and recovery WAL records when a cluster is in the +"standby mode", and after it was detached to function as a normal cluster. + +To configure such cluster you need to specify the section ``standby_cluster`` +in a patroni configuration: + +.. code:: YAML + + bootstrap: + dcs: + standby_cluster: + host: 1.2.3.4 + port: 5432 + primary_slot_name: patroni + +Note, that these options will be applied only once during cluster bootstrap, +and the only way to change them afterwards is through DCS. diff --git a/features/basic_replication.feature b/features/basic_replication.feature index 8aa112fa..e0eafc2a 100644 --- a/features/basic_replication.feature +++ b/features/basic_replication.feature @@ -32,7 +32,7 @@ Feature: basic replication Then I receive a response returncode 0 When I sleep for 2 seconds And I shut down postgres0 - And I run patronictl.py resume batman + And I run patronictl.py resume batman Then I receive a response returncode 0 And postgres2 role is the primary after 24 seconds When I issue a PATCH request to http://127.0.0.1:8010/config with {"synchronous_mode": null, "master_start_timeout": 0} diff --git a/features/standby_cluster.feature b/features/standby_cluster.feature new file mode 100644 index 00000000..51127efe --- /dev/null +++ b/features/standby_cluster.feature @@ -0,0 +1,16 @@ +Feature: standby cluster + + Scenario: check replication of a single table in a standby cluster + Given I start postgres0 without slots sync + And I create a replication slot postgres1 on postgres0 + And I start postgres1 in a standby cluster batman1 as a clone of postgres0 + Then postgres1 is a leader of batman1 after 10 seconds + When I add the table foo to postgres0 + Then table foo is present on postgres1 after 20 seconds + When I start postgres2 in a cluster batman1 + Then postgres2 role is the replica after 24 seconds + And table foo is present on postgres2 after 20 seconds + + Scenario: check failover + When I kill postgres1 + Then postgres2 is replicating from postgres0 after 20 seconds diff --git a/features/steps/standby_cluster.py b/features/steps/standby_cluster.py new file mode 100644 index 00000000..6eb8975d --- /dev/null +++ b/features/steps/standby_cluster.py @@ -0,0 +1,78 @@ +import time + +from behave import step + + +select_replication_query = """ +SELECT * FROM pg_catalog.pg_stat_replication +WHERE application_name = '{0}' +""" + +create_replication_slot_query = """ +SELECT pg_create_physical_replication_slot('{0}') +""" + + +@step('I start {name:w} without slots sync') +def start_patroni_without_slots_sync(context, name): + return context.pctl.start(name, custom_config={ + "bootstrap": { + "dcs": { + "postgresql": { + "use_slots": False + } + } + } + }) + + +@step('I start {name:w} in a cluster {cluster_name:w}') +def start_patroni(context, name, cluster_name): + return context.pctl.start(name, custom_config={ + "scope": cluster_name + }) + + +@step('I start {name:w} in a standby cluster {cluster_name:w} as a clone of {name2:w}') +def start_patroni_stanby_cluster(context, name, cluster_name, name2): + port = context.pctl._processes[name2]._connkwargs.get('port') + return context.pctl.start(name, custom_config={ + "scope": cluster_name, + "bootstrap": { + "dcs": { + "standby_cluster": { + "host": "localhost", + "port": port, + "primary_slot_name": "postgres1", + } + } + } + }) + + +@step('{pg_name1:w} is replicating from {pg_name2:w} after {timeout:d} seconds') +def check_replication_status(context, pg_name1, pg_name2, timeout): + bound_time = time.time() + timeout + + while time.time() < bound_time: + cur = context.pctl.query( + pg_name2, + select_replication_query.format(pg_name1), + fail_ok=True + ) + + if cur and len(cur.fetchall()) != 0: + return True + + time.sleep(1) + + return False + + +@step('I create a replication slot {slot_name:w} on {pg_name:w}') +def create_replication_slot(context, slot_name, pg_name): + return context.pctl.query( + pg_name, + create_replication_slot_query.format(slot_name), + fail_ok=True + ) diff --git a/patroni/config.py b/patroni/config.py index ee809f5b..dc8a6045 100644 --- a/patroni/config.py +++ b/patroni/config.py @@ -1,13 +1,14 @@ import json import logging import os +import six import sys import tempfile import yaml from collections import defaultdict from copy import deepcopy -from patroni.dcs import ClusterConfig +from patroni.dcs import ClusterConfig, is_standby_cluster from patroni.postgresql import Postgresql from patroni.utils import deep_compare, parse_bool, parse_int, patch_config from requests.structures import CaseInsensitiveDict @@ -45,6 +46,15 @@ class Config(object): 'master_start_timeout': 300, 'synchronous_mode': False, 'synchronous_mode_strict': False, + 'standby_cluster': { + 'create_replica_methods': '', + 'host': '', + 'port': '', + 'primary_slot_name': '', + 'restore_command': '', + 'archive_cleanup_command': '', + 'recovery_min_apply_delay': '' + }, 'postgresql': { 'bin_dir': '', 'use_slots': True, @@ -88,6 +98,10 @@ class Config(object): def dynamic_configuration(self): return deepcopy(self._dynamic_configuration) + @property + def is_standby_cluster(self): + return is_standby_cluster(self._dynamic_configuration.get('standby_cluster')) + def check_mode(self, mode): return bool(parse_bool(self._dynamic_configuration.get(mode))) @@ -183,6 +197,13 @@ class Config(object): config['postgresql'][name].update(self._process_postgresql_parameters(value)) elif name not in ('connect_address', 'listen', 'data_dir', 'pgpass', 'authentication'): config['postgresql'][name] = deepcopy(value) + elif name == 'standby_cluster': + allowed_keys = self.__DEFAULT_CONFIG['standby_cluster'].keys() + expected = { + k: v for k, v in (value or {}).items() + if (k in allowed_keys and isinstance(v, six.string_types)) + } + config['standby_cluster'].update(expected) elif name in config: # only variables present in __DEFAULT_CONFIG allowed to be overriden from DCS if name in ('synchronous_mode', 'synchronous_mode_strict'): config[name] = value @@ -317,8 +338,15 @@ class Config(object): if 'name' not in config and 'name' in pg_config: config['name'] = pg_config['name'] - pg_config.update({p: config[p] for p in ('name', 'scope', 'retry_timeout', - 'synchronous_mode', 'maximum_lag_on_failover') if p in config}) + updated_fields = ( + 'name', + 'scope', + 'retry_timeout', + 'synchronous_mode', + 'maximum_lag_on_failover' + ) + + pg_config.update({p: config[p] for p in updated_fields if p in config}) return config diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index 346a22b0..8d050fd5 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -108,12 +108,29 @@ class Member(namedtuple('Member', 'index,name,session,data')): @property def conn_url(self): - return self.data.get('conn_url') + conn_url = self.data.get('conn_url') + conn_kwargs = self.data.get('conn_kwargs') + if conn_url: + return conn_url + + if conn_kwargs: + conn_url = 'postgresql://{host}:{port}'.format( + host=conn_kwargs.get('host'), + port=conn_kwargs.get('port'), + ) + self.data['conn_url'] = conn_url + return conn_url def conn_kwargs(self, auth=None): + defaults = { + "host": "", + "port": "", + "database": "" + } ret = self.data.get('conn_kwargs') if ret: - ret = ret.copy() + defaults.update(ret) + ret = defaults else: r = urlparse(self.conn_url) ret = { @@ -159,6 +176,27 @@ class Member(namedtuple('Member', 'index,name,session,data')): return self.state == 'running' +class RemoteMember(Member): + """ Represents a remote master for a standby cluster + """ + def __new__(cls, name, data): + return super(RemoteMember, cls).__new__(cls, None, name, None, data) + + @staticmethod + def allowed_keys(): + return ('primary_slot_name', + 'create_replica_methods', + 'restore_command', + 'archive_cleanup_command', + 'recovery_min_apply_delay') + + def __getattr__(self, name): + if name not in RemoteMember.allowed_keys(): + return + + return self.data.get(name) + + class Leader(namedtuple('Leader', 'index,session,member')): """Immutable object (namedtuple) which represents leader key. @@ -359,6 +397,9 @@ class Cluster(namedtuple('Cluster', 'initialize,config,leader,last_leader_operat def is_synchronous_mode(self): return self.check_mode('synchronous_mode') + def is_standby_cluster(self): + return is_standby_cluster(self.config and self.config.data.get('standby_cluster')) + @six.add_metaclass(abc.ABCMeta) class AbstractDCS(object): @@ -612,3 +653,14 @@ class AbstractDCS(object): self.event.wait(timeout) return self.event.isSet() + + +def is_standby_cluster(config): + """ Check whether or not provided configuration describes a standby cluster. + Config can be both patroni config or cluster.config.data + """ + return isinstance(config, dict) and ( + config.get('host') or + config.get('port') or + config.get('restore_command') + ) diff --git a/patroni/ha.py b/patroni/ha.py index f0707b56..0a712dd7 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -6,6 +6,7 @@ import psycopg2 import requests import sys import time +import uuid from collections import namedtuple from multiprocessing.pool import ThreadPool @@ -13,6 +14,7 @@ from patroni.async_executor import AsyncExecutor, CriticalTask from patroni.exceptions import DCSError, PostgresConnectionException, PatroniException from patroni.postgresql import ACTION_ON_START from patroni.utils import polling_loop, tzutc +from patroni.dcs import RemoteMember from threading import RLock logger = logging.getLogger(__name__) @@ -193,14 +195,24 @@ class Ha(object): self._async_executor.schedule('bootstrap {0}'.format(msg)) self._async_executor.run_async(self.clone, args=(clone_member, msg)) return 'trying to bootstrap {0}'.format(msg) + # no initialize key and node is allowed to be master and has 'bootstrap' section in a configuration file elif self.cluster.initialize is None and not self.patroni.nofailover and 'bootstrap' in self.patroni.config: if self.dcs.initialize(create_new=True): # race for initialization self.state_handler.bootstrapping = True self._post_bootstrap_task = CriticalTask() - self._async_executor.schedule('bootstrap') - self._async_executor.run_async(self.state_handler.bootstrap, args=(self.patroni.config['bootstrap'],)) - return 'trying to bootstrap a new cluster' + + if self.patroni.config.is_standby_cluster: + self._async_executor.schedule('bootstrap_standby_leader') + self._async_executor.run_async(self.bootstrap_standby_leader) + return 'trying to bootstrap a new standby leader' + else: + self._async_executor.schedule('bootstrap') + self._async_executor.run_async( + self.state_handler.bootstrap, + args=(self.patroni.config['bootstrap'],) + ) + return 'trying to bootstrap a new cluster' else: return 'failed to acquire initialize lock' else: @@ -211,6 +223,21 @@ class Ha(object): return 'trying to ' + msg return 'waiting for leader to bootstrap' + def bootstrap_standby_leader(self): + """ If we found 'standby' key in the configuration, we need to bootstrap + not a real master, but a 'standby leader', that will take base backup + from a remote master and start follow it. + """ + patroni_config = self.patroni.config.dynamic_configuration + clone_source = self.get_remote_master(patroni_config) + msg = 'clone from remote master {0}'.format(clone_source.conn_url) + result = self.clone(clone_source, msg) + self._post_bootstrap_task.complete(result) + if result: + self.state_handler.set_role('standby_leader') + + return result + def _handle_rewind(self): if self.state_handler.rewind_needed_and_possible(self.cluster.leader): self._async_executor.schedule('running pg_rewind from ' + self.cluster.leader.name) @@ -266,12 +293,18 @@ class Ha(object): def _get_node_to_follow(self, cluster): # determine the node to follow. If replicatefrom tag is set, # try to follow the node mentioned there, otherwise, follow the leader. - if not self.patroni.replicatefrom or self.patroni.replicatefrom == self.state_handler.name: - node_to_follow = cluster.leader - else: - node_to_follow = cluster.get_member(self.patroni.replicatefrom) + is_leader = self.cluster.leader and self.state_handler.name == self.cluster.leader.name - return node_to_follow if node_to_follow and node_to_follow.name != self.state_handler.name else None + if self.cluster.is_standby_cluster() and is_leader: + node_to_follow = self.get_remote_master(cluster.config.data) + elif self.patroni.replicatefrom and self.patroni.replicatefrom != self.state_handler.name: + node_to_follow = cluster.get_member(self.patroni.replicatefrom) + else: + node_to_follow = cluster.leader + + return (node_to_follow if + node_to_follow and + node_to_follow.name != self.state_handler.name else None) def follow(self, demote_reason, follow_reason, refresh=True): if refresh: @@ -411,6 +444,12 @@ class Ha(object): line.append(cluster_history[line[0]][3]) self.dcs.set_history_value(json.dumps(history, separators=(',', ':'))) + def enforce_follow_remote_master(self, message): + self.state_handler.set_role('standby_leader') + demote_reason = 'cannot be a real master in standby cluster' + + return self.follow(demote_reason, message) + def enforce_master_role(self, message, promote_message): if not self.is_paused() and not self.watchdog.is_running and not self.watchdog.activate(): if self.state_handler.is_leader(): @@ -740,8 +779,18 @@ class Ha(object): logger.info('Cleaning up failover key after acquiring leader lock...') self.dcs.manual_failover('', '') self.load_cluster_from_dcs() - return self.enforce_master_role('acquired session lock as a leader', - 'promoted self to leader by acquiring session lock') + + if self.cluster.is_standby_cluster(): + # standby leader disappeared, and this is a healthiest + # replica, so it should become a new standby leader. + # This imply that we need to start following a remote master + msg = 'promoted self to a standby leader because i had the session lock' + return self.enforce_follow_remote_master(msg) + else: + return self.enforce_master_role( + 'acquired session lock as a leader', + 'promoted self to leader by acquiring session lock' + ) else: return self.follow('demoted self after trying and failing to obtain lock', 'following new leader after trying and failing to obtain lock') @@ -773,8 +822,17 @@ class Ha(object): if msg is not None: return msg - return self.enforce_master_role('no action. i am the leader with the lock', - 'promoted self to leader because i had the session lock') + if self.cluster.is_standby_cluster(): + # in case of standby cluster we don't really need to + # enforce anything, since the leader is not a master. + # So just remind the role. + msg = 'no action. i am the standby leader with the lock' + return self.enforce_follow_remote_master(msg) + else: + return self.enforce_master_role( + 'no action. i am the leader with the lock', + 'promoted self to leader because i had the session lock' + ) else: # Either there is no connection to DCS or someone else acquired the lock logger.error('failed to update leader lock') @@ -1001,6 +1059,7 @@ class Ha(object): self.set_is_leader(True) self.state_handler.call_nowait(ACTION_ON_START) self.load_cluster_from_dcs() + return 'initialized a new cluster' def handle_starting_instance(self): @@ -1215,3 +1274,27 @@ class Ha(object): no "active" leader watch request in progress. This usually happens on the master or if the node is running async action""" self.dcs.event.set() + + def get_remote_master(self, config): + """ In case of standby cluster this will tel us from which remote + master to stream. Config can be both patroni config or + cluster.config.data + """ + config = config or (self.config is not None and self.config.data) + + if config and config.get('standby_cluster'): + cluster_params = config.get('standby_cluster') + unique_name = 'remote_master:{}'.format(uuid.uuid1()) + data = { + 'conn_kwargs': { + "host": cluster_params.get('host'), + "port": cluster_params.get('port'), + }, + 'no_replication_slot': 'primary_slot_name' not in cluster_params, + } + data.update({ + k: v for k, v in cluster_params.items() + if k in RemoteMember.allowed_keys() + }) + + return RemoteMember(unique_name, data) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index b5c305b3..f69a52bb 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -15,11 +15,13 @@ from patroni.callback_executor import CallbackExecutor from patroni.exceptions import PostgresConnectionException, PostgresException from patroni.utils import compare_values, parse_bool, parse_int, Retry, RetryFailedError, polling_loop, split_host_port from patroni.postmaster import PostmasterProcess +from patroni.dcs import RemoteMember from requests.structures import CaseInsensitiveDict from six import string_types from six.moves.urllib.parse import quote_plus from threading import current_thread, Lock + logger = logging.getLogger(__name__) ACTION_ON_START = "on_start" @@ -665,9 +667,17 @@ class Postgresql(object): self.set_state('creating replica') self._sysid = None - # get list of replica methods from config. - # If there is no configuration key, or no value is specified, use basebackup - replica_methods = self._create_replica_methods or ['basebackup'] + is_remote_master = isinstance(clone_member, RemoteMember) + create_replica_methods = is_remote_master and clone_member.create_replica_methods + + # get list of replica methods either from clone member or from + # the config. If there is no configuration key, or no value is + # specified, use basebackup + replica_methods = ( + create_replica_methods + or self._create_replica_methods + or ['basebackup'] + ) if clone_member and clone_member.conn_url: r = clone_member.conn_kwargs(self._replication) @@ -1413,6 +1423,12 @@ class Postgresql(object): return self._rewind_state == REWIND_STATUS.FAILED def follow(self, member, 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.primary_conninfo(member) change_role = self.role in ('master', 'demoted') @@ -1420,8 +1436,16 @@ class Postgresql(object): recovery_params.update({'standby_mode': 'on', 'recovery_target_timeline': 'latest'}) if primary_conninfo: recovery_params['primary_conninfo'] = primary_conninfo - if self.use_slots: - recovery_params['primary_slot_name'] = slot_name_from_member_name(self.name) + if self.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 self.write_recovery_conf(recovery_params) @@ -1533,7 +1557,7 @@ $$""".format(name, ' '.join(options)), name, password, password) # the current master, because that member would replicate from elsewhere. We still create the slot if # the replicatefrom destination member is currently not a member of the cluster (fallback to the # master), or if replicatefrom destination member happens to be the current master - if self.role == 'master': + if self.role in ('master', 'standby_leader'): slot_members = [m.name for m in cluster.members if m.name != self.name and (m.replicatefrom is None or m.replicatefrom == self.name or not cluster.has_member(m.replicatefrom))] diff --git a/postgres0.yml b/postgres0.yml index 0f3a1935..5b4cfe56 100644 --- a/postgres0.yml +++ b/postgres0.yml @@ -29,6 +29,10 @@ bootstrap: maximum_lag_on_failover: 1048576 # master_start_timeout: 300 # synchronous_mode: false + #standby_cluster: + #host: 127.0.0.1 + #port: 1111 + #primary_slot_name: patroni postgresql: use_pg_rewind: true # use_slots: true diff --git a/tests/test_config.py b/tests/test_config.py index 3d899e38..a2f661f9 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -23,7 +23,7 @@ class TestConfig(unittest.TestCase): def test_set_dynamic_configuration(self): with patch.object(Config, '_build_effective_configuration', Mock(side_effect=Exception)): self.assertIsNone(self.config.set_dynamic_configuration({'foo': 'bar'})) - self.assertTrue(self.config.set_dynamic_configuration({'synchronous_mode': True})) + self.assertTrue(self.config.set_dynamic_configuration({'synchronous_mode': True, 'standby_cluster': {}})) def test_reload_local_configuration(self): os.environ.update({ diff --git a/tests/test_ctl.py b/tests/test_ctl.py index e388d910..0afea71d 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -344,12 +344,15 @@ class TestCtl(unittest.TestCase): assert result.exit_code == 0 @patch('requests.post', Mock(side_effect=requests.exceptions.ConnectionError('foo'))) - def test_request_patroni(self): - context = {'restapi': {'keyfile': '/etc/patroni/key.pem', 'certfile': 'cert.pem'}} - with patch('click.get_current_context') as mock_context: - mock_context.return_value.obj = context - member = get_cluster_initialized_with_leader().leader.member - self.assertRaises(requests.exceptions.ConnectionError, request_patroni, member, 'post', 'dummy', {}) + @patch('click.get_current_context') + def test_request_patroni(self, mock_context): + member = get_cluster_initialized_with_leader().leader.member + + mock_context.return_value.obj = {'ctl': {'cacert': 'cert.pem'}} + self.assertRaises(requests.exceptions.ConnectionError, request_patroni, member, 'post', 'dummy', {}) + + mock_context.return_value.obj = {'ctl': {'insecure': True}} + self.assertRaises(requests.exceptions.ConnectionError, request_patroni, member, 'post', 'dummy', {}) def test_ctl(self): self.runner.invoke(ctl, ['list']) diff --git a/tests/test_ha.py b/tests/test_ha.py index b5364669..be9b94ae 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -2,8 +2,10 @@ import datetime import etcd import os import unittest +import sys from mock import Mock, MagicMock, PropertyMock, patch +from patroni.async_executor import CriticalTask from patroni.config import Config from patroni.dcs import Cluster, ClusterConfig, Failover, Leader, Member, get_dcs, SyncState, TimelineHistory from patroni.dcs.etcd import Client @@ -26,16 +28,17 @@ def false(*args, **kwargs): return False -def get_cluster(initialize, leader, members, failover, sync): +def get_cluster(initialize, leader, members, failover, sync, cluster_config=None): history = TimelineHistory(1, [(1, 67197376, 'no recovery target specified', datetime.datetime.now().isoformat())]) - return Cluster(initialize, ClusterConfig(1, {1: 2}, 1), leader, 10, members, failover, sync, history) + cluster_config = cluster_config or ClusterConfig(1, {1: 2}, 1) + return Cluster(initialize, cluster_config, leader, 10, members, failover, sync, history) -def get_cluster_not_initialized_without_leader(): - return get_cluster(None, None, [], None, SyncState(None, None, None)) +def get_cluster_not_initialized_without_leader(cluster_config=None): + return get_cluster(None, None, [], None, SyncState(None, None, None), cluster_config) -def get_cluster_initialized_without_leader(leader=False, failover=None, sync=None): +def get_cluster_initialized_without_leader(leader=False, failover=None, sync=None, cluster_config=None): m1 = Member(0, 'leader', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', 'api_url': 'http://127.0.0.1:8008/patroni', 'xlog_location': 4}) leader = Leader(0, 0, m1) if leader else None @@ -47,16 +50,38 @@ def get_cluster_initialized_without_leader(leader=False, failover=None, sync=Non 'scheduled_restart': {'schedule': "2100-01-01 10:53:07.560445+00:00", 'postgres_version': '99.0.0'}}) syncstate = SyncState(0 if sync else None, sync and sync[0], sync and sync[1]) - return get_cluster(SYSID, leader, [m1, m2], failover, syncstate) + return get_cluster(SYSID, leader, [m1, m2], failover, syncstate, cluster_config) def get_cluster_initialized_with_leader(failover=None, sync=None): return get_cluster_initialized_without_leader(leader=True, failover=failover, sync=sync) -def get_cluster_initialized_with_only_leader(failover=None): +def get_cluster_initialized_with_only_leader(failover=None, cluster_config=None): leader = get_cluster_initialized_without_leader(leader=True, failover=failover).leader - return get_cluster(True, leader, [leader], failover, None) + return get_cluster(True, leader, [leader], failover, None, cluster_config) + + +def get_cluster_not_initialized_standby(failover=None, sync=None): + return get_cluster_not_initialized_without_leader( + cluster_config=ClusterConfig(1, { + "standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }}, 1) + ) + + +def get_standby_cluster_initialized_with_only_leader(failover=None, sync=None): + return get_cluster_initialized_with_only_leader( + cluster_config=ClusterConfig(1, { + "standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }}, 1) + ) def get_node_status(reachable=True, in_recovery=True, wal_position=10, nofailover=False, watchdog_failed=False): @@ -97,6 +122,10 @@ zookeeper: hosts: [localhost] port: 8181 """ + # We rely on sys.argv in Config, so it's necessary to reset + # all the extra values that are coming from py.test + sys.argv = sys.argv[:1] + self.config = Config() self.postgresql = p self.dcs = d @@ -182,6 +211,50 @@ class TestHa(unittest.TestCase): self.p.is_healthy = false self.assertEquals(self.ha.run_cycle(), 'starting as a secondary') + @patch('patroni.dcs.etcd.Etcd.initialize', return_value=True) + def test_start_as_standby_leader(self, initialize): + self.p.data_directory_empty = true + self.ha.cluster = get_cluster_not_initialized_standby() + self.ha.cluster.is_unlocked = true + self.ha.patroni.config._dynamic_configuration = {"standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }} + self.assertEquals( + self.ha.run_cycle(), + 'trying to bootstrap a new standby leader' + ) + + @patch.object(Cluster, 'get_clone_member', + Mock(return_value=Member(0, 'test', 1, {'api_url': 'http://127.0.0.1:8011/patroni'}))) + @patch.object(Postgresql, 'create_replica', Mock(return_value=0)) + def test_start_as_cascade_replica_in_standby_cluster(self): + self.p.data_directory_empty = true + self.ha.cluster = get_standby_cluster_initialized_with_only_leader() + self.ha.cluster.is_unlocked = false + self.ha.patroni.config._dynamic_configuration = {"standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }} + self.assertEquals( + self.ha.run_cycle(), + "trying to bootstrap from replica 'test'" + ) + + @patch.object(Postgresql, 'create_replica', Mock(return_value=0)) + def test_bootstrap_standby_leader(self): + self.ha.cluster = get_cluster_not_initialized_standby() + self.ha.cluster.is_unlocked = true + self.ha.patroni.config._dynamic_configuration = {"standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }} + self.ha._post_bootstrap_task = CriticalTask() + self.assertEquals(self.ha.bootstrap_standby_leader(), True) + def test_recover_replica_failed(self): self.p.controldata = lambda: {'Database cluster state': 'in recovery', 'Database system identifier': SYSID} self.p.is_running = false @@ -616,6 +689,61 @@ class TestHa(unittest.TestCase): self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', self.p.name, None)) self.assertEquals(self.ha.run_cycle(), 'PAUSE: waiting to become master after promote...') + def test_process_healthy_standby_cluster_as_standby_leader(self): + self.p.is_leader = false + self.p.name = 'leader' + self.ha.patroni.config._dynamic_configuration = {"standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }} + self.ha.cluster = get_standby_cluster_initialized_with_only_leader() + msg = 'no action. i am the standby leader with the lock' + self.assertEquals(self.ha.run_cycle(), msg) + + def test_process_healthy_standby_cluster_as_cascade_replica(self): + self.p.is_leader = false + self.p.name = 'replica' + self.ha.patroni.config._dynamic_configuration = {"standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }} + self.ha.cluster = get_standby_cluster_initialized_with_only_leader() + msg = 'no action. i am a secondary and i am following a leader' + self.assertEquals(self.ha.run_cycle(), msg) + + @patch('patroni.dcs.etcd.Etcd.initialize', return_value=True) + def test_process_unhealthy_standby_cluster_as_standby_leader(self, initialize): + self.p.is_leader = false + self.p.name = 'leader' + self.ha.patroni.config._dynamic_configuration = {"standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }} + self.ha.cluster = get_standby_cluster_initialized_with_only_leader() + self.ha.cluster.is_unlocked = true + self.ha.sysid_valid = true + self.p._sysid = True + msg = 'promoted self to a standby leader because i had the session lock' + self.assertEquals(self.ha.run_cycle(), msg) + + @patch.object(Postgresql, 'rewind_needed_and_possible', Mock(return_value=True)) + @patch('patroni.dcs.etcd.Etcd.initialize', return_value=True) + def test_process_unhealthy_standby_cluster_as_cascade_replica(self, initialize): + self.p.is_leader = false + self.p.name = 'replica' + self.ha.patroni.config._dynamic_configuration = {"standby_cluster": { + "host": "localhost", + "port": 5432, + "primary_slot_name": "", + }} + self.ha.cluster = get_standby_cluster_initialized_with_only_leader() + self.ha.is_unlocked = true + msg = 'running pg_rewind from leader' + self.assertEquals(self.ha.run_cycle(), msg) + def test_failed_to_update_lock_in_pause(self): self.ha.update_lock = false self.ha.is_paused = true diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 2c0f2874..4019979c 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -8,7 +8,7 @@ import unittest from mock import Mock, MagicMock, PropertyMock, patch, mock_open from patroni.async_executor import CriticalTask -from patroni.dcs import Cluster, Leader, Member, SyncState +from patroni.dcs import Cluster, Leader, Member, RemoteMember, SyncState from patroni.exceptions import PostgresConnectionException, PostgresException from patroni.postgresql import Postgresql, STATE_REJECT, STATE_NO_RESPONSE from patroni.postmaster import PostmasterProcess @@ -401,7 +401,7 @@ class TestPostgresql(unittest.TestCase): @patch.object(Postgresql, 'is_running', Mock(return_value=False)) @patch.object(Postgresql, 'start', Mock()) def test_follow(self): - self.p.follow(None) + self.p.follow(RemoteMember('123', {'recovery_command': 'foo'})) @patch('subprocess.check_output', Mock(return_value=0, side_effect=pg_controldata_string)) def test_can_rewind(self):