mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-09-01 09:09:21 +00:00
[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.
This commit is contained in:
@@ -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 <custom_bootstrap>` for details.
|
||||
When ``initdb`` is specified revert to the default ``initdb`` command. ``initdb`` is also triggered when no ``method``
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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}
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
)
|
||||
+31
-3
@@ -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
|
||||
|
||||
|
||||
+54
-2
@@ -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')
|
||||
)
|
||||
|
||||
+95
-12
@@ -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)
|
||||
|
||||
+30
-6
@@ -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))]
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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({
|
||||
|
||||
+9
-6
@@ -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'])
|
||||
|
||||
+136
-8
@@ -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:[email protected]: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
|
||||
|
||||
@@ -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):
|
||||
|
||||
Reference in New Issue
Block a user