mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Merge pull request #166 from zalando/feature/clonefrom
Correct implementation of 'clonefrom' feature
This commit is contained in:
@@ -1,13 +1,13 @@
|
||||
Feature: cascading replication
|
||||
We should check that patroni can do base backup and streaming from the replica
|
||||
|
||||
Scenario: check a base backup from the replica
|
||||
Scenario: check a base backup and streaming replication from a replica
|
||||
Given I start postgres0
|
||||
And postgres0 is a leader after 10 seconds
|
||||
And I start postgres1
|
||||
And I configure and start postgres1 with a tag clonefrom true
|
||||
And replication works from postgres0 to postgres1 after 15 seconds
|
||||
And I create label with "postgres0" in postgres0 data directory
|
||||
And I create label with "postgres1" in postgres1 data directory
|
||||
And I configure and start postgres2 with a tag clonefrom postgres1
|
||||
And I configure and start postgres2 with a tag replicatefrom postgres1
|
||||
Then replication works from postgres0 to postgres2 after 30 seconds
|
||||
And there is a label with "postgres1" in postgres2 data directory
|
||||
|
||||
+2
-5
@@ -20,7 +20,8 @@ class Patroni(object):
|
||||
|
||||
def __init__(self, config):
|
||||
self.nap_time = config['loop_wait']
|
||||
self.tags = config.get('tags', dict())
|
||||
self.tags = {tag: value for tag, value in config.get('tags', {}).items()
|
||||
if tag not in ('clonefrom', 'nofailover', 'noloadbalance') or value}
|
||||
self.postgresql = Postgresql(config['postgresql'])
|
||||
self.dcs = self.get_dcs(self.postgresql.name, config)
|
||||
self.version = __version__
|
||||
@@ -36,10 +37,6 @@ class Patroni(object):
|
||||
def replicatefrom(self):
|
||||
return self.tags.get('replicatefrom')
|
||||
|
||||
@property
|
||||
def clonefrom(self):
|
||||
return self.tags.get('clonefrom')
|
||||
|
||||
@staticmethod
|
||||
def get_dcs(name, config):
|
||||
if 'etcd' in config:
|
||||
|
||||
+1
-1
@@ -290,7 +290,7 @@ class RestApiHandler(BaseHTTPRequestHandler):
|
||||
return {'state': state}
|
||||
|
||||
def get_tags(self):
|
||||
return {'tags': self.server.patroni.tags}
|
||||
return {'tags': self.server.patroni.tags} if self.server.patroni.tags else {}
|
||||
|
||||
def log_message(self, fmt, *args):
|
||||
logger.debug("API thread: %s - - [%s] %s", self.client_address[0], self.log_date_time_string(), fmt % args)
|
||||
|
||||
+17
-4
@@ -4,6 +4,7 @@ import json
|
||||
import six
|
||||
|
||||
from collections import namedtuple
|
||||
from random import randint
|
||||
from six.moves.urllib_parse import urlparse, urlunparse, parse_qsl
|
||||
from threading import Event, Lock
|
||||
|
||||
@@ -64,13 +65,21 @@ class Member(namedtuple('Member', 'index,name,session,data')):
|
||||
def api_url(self):
|
||||
return self.data.get('api_url')
|
||||
|
||||
@property
|
||||
def tags(self):
|
||||
return self.data.get('tags', {})
|
||||
|
||||
@property
|
||||
def nofailover(self):
|
||||
return self.data.get('tags', {}).get('nofailover', False)
|
||||
return self.tags.get('nofailover', False)
|
||||
|
||||
@property
|
||||
def replicatefrom(self):
|
||||
return self.data.get('tags', {}).get('replicatefrom')
|
||||
return self.tags.get('replicatefrom')
|
||||
|
||||
@property
|
||||
def clonefrom(self):
|
||||
return self.tags.get('clonefrom', False)
|
||||
|
||||
|
||||
class Leader(namedtuple('Leader', 'index,session,member')):
|
||||
@@ -147,8 +156,12 @@ class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,mem
|
||||
def has_member(self, member_name):
|
||||
return any(m for m in self.members if m.name == member_name)
|
||||
|
||||
def get_member(self, member_name):
|
||||
return ([m for m in self.members if m.name == member_name] or [None])[0]
|
||||
def get_member(self, member_name, fallback_to_leader=True):
|
||||
return ([m for m in self.members if m.name == member_name] or [self.leader if fallback_to_leader else None])[0]
|
||||
|
||||
def get_clone_member(self):
|
||||
candidates = [m for m in self.members if m.clonefrom and (not self.leader or m.name != self.leader.name)]
|
||||
return candidates[randint(0, len(candidates) - 1)] if candidates else self.leader
|
||||
|
||||
|
||||
@six.add_metaclass(abc.ABCMeta)
|
||||
|
||||
+22
-22
@@ -55,9 +55,10 @@ class Ha(object):
|
||||
'conn_url': self.state_handler.connection_string,
|
||||
'api_url': self.patroni.api.connection_string,
|
||||
'state': self.state_handler.state,
|
||||
'role': self.state_handler.role,
|
||||
'tags': self.patroni.tags
|
||||
'role': self.state_handler.role
|
||||
}
|
||||
if self.patroni.tags:
|
||||
data['tags'] = self.patroni.tags
|
||||
if data['state'] in ['running', 'restarting', 'starting']:
|
||||
try:
|
||||
data['xlog_location'] = self.state_handler.xlog_position()
|
||||
@@ -65,25 +66,22 @@ class Ha(object):
|
||||
pass
|
||||
self.dcs.touch_member(json.dumps(data, separators=(',', ':')))
|
||||
|
||||
def clone(self, clone_member, clone_member_name="leader"):
|
||||
def clone(self, clone_member=None, msg='(without leader)'):
|
||||
if self.state_handler.bootstrap(cluster_initialized=True, clone_member=clone_member):
|
||||
logger.info('bootstrapped from {0}'.format(clone_member_name)
|
||||
if clone_member else 'bootstrapped without leader')
|
||||
logger.info('bootstrapped %s', msg)
|
||||
else:
|
||||
logger.error('failed to bootstrap %s', msg)
|
||||
self.state_handler.stop('immediate')
|
||||
self.state_handler.remove_data_directory()
|
||||
logger.error('failed to bootstrap from {0}'.format(clone_member_name)
|
||||
if clone_member else 'failed to bootstrap (without leader)')
|
||||
|
||||
def bootstrap(self):
|
||||
if not self.cluster.is_unlocked(): # cluster already has leader
|
||||
clonefrom = self.patroni.clonefrom
|
||||
clone_member = self.cluster.get_member(clonefrom)\
|
||||
if self.cluster.has_member(clonefrom) else self.cluster.leader
|
||||
clone_member_name = 'leader' if clone_member == self.cluster.leader else 'replica \'{0}\''.format(clonefrom)
|
||||
self._async_executor.schedule('bootstrap from {0}'.format(clone_member_name))
|
||||
self._async_executor.run_async(self.clone, args=(clone_member, clone_member_name))
|
||||
return 'trying to bootstrap from {0}'.format(clone_member_name)
|
||||
clone_member = self.cluster.get_clone_member()
|
||||
member_role = 'leader' if clone_member == self.cluster.leader else 'replica'
|
||||
msg = "from {0} '{1}'".format(member_role, clone_member.name)
|
||||
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)
|
||||
elif not self.cluster.initialize and not self.patroni.nofailover: # no initialize key
|
||||
if self.dcs.initialize(create_new=True): # race for initialization
|
||||
try:
|
||||
@@ -103,8 +101,8 @@ class Ha(object):
|
||||
return 'failed to acquire initialize lock'
|
||||
else:
|
||||
if self.state_handler.can_create_replica_without_replication_connection():
|
||||
self._async_executor.run_async(self.clone, args=(None, ))
|
||||
return "trying to bootstrap without leader"
|
||||
self._async_executor.run_async(self.clone)
|
||||
return "trying to bootstrap (without leader)"
|
||||
return 'waiting for leader to bootstrap'
|
||||
|
||||
def recover(self):
|
||||
@@ -127,8 +125,7 @@ class Ha(object):
|
||||
# try to follow the node mentioned there, otherwise, follow the leader.
|
||||
|
||||
if self.patroni.replicatefrom:
|
||||
node_to_follow = [m for m in self.cluster.members if m.name == self.patroni.replicatefrom]
|
||||
node_to_follow = node_to_follow[0] if node_to_follow else self.cluster.leader
|
||||
node_to_follow = self.cluster.get_member(self.patroni.replicatefrom, fallback_to_leader=True)
|
||||
else:
|
||||
node_to_follow = self.cluster.leader
|
||||
if node_to_follow and node_to_follow.name == self.state_handler.name:
|
||||
@@ -225,9 +222,9 @@ class Ha(object):
|
||||
return True
|
||||
|
||||
# find specific node and check that it is healthy
|
||||
members = [m for m in self.cluster.members if m.name == failover.candidate]
|
||||
if members:
|
||||
member, reachable, _, _, tags = self.fetch_node_status(members[0])
|
||||
member = self.cluster.get_member(failover.candidate, fallback_to_leader=False)
|
||||
if member:
|
||||
member, reachable, _, _, tags = self.fetch_node_status(member)
|
||||
if reachable and not tags.get('nofailover', False): # node is healthy
|
||||
logger.info('manual failover: to %s, i am %s', member.name, self.state_handler.name)
|
||||
return False
|
||||
@@ -390,7 +387,10 @@ class Ha(object):
|
||||
def reinitialize(self, cluster):
|
||||
self.state_handler.stop('immediate')
|
||||
self.state_handler.remove_data_directory()
|
||||
self.clone(cluster.leader)
|
||||
|
||||
clone_member = cluster.get_clone_member()
|
||||
member_role = 'leader' if clone_member == cluster.leader else 'replica'
|
||||
self.clone(clone_member, "from {0} '{1}'".format(member_role, clone_member.name))
|
||||
|
||||
def process_scheduled_action(self):
|
||||
if self.reinitialize_scheduled():
|
||||
|
||||
@@ -269,9 +269,8 @@ class Postgresql(object):
|
||||
# If there is no configuration key, or no value is specified, use basebackup
|
||||
replica_methods = self.config.get('create_replica_method') or ['basebackup']
|
||||
# if we don't have any source, leave only replica methods that work without it
|
||||
replica_methods = \
|
||||
[r for r in replica_methods if self.replica_method_can_work_without_replication_connection(r)]\
|
||||
if not clone_member else replica_methods
|
||||
replica_methods = replica_methods if clone_member else \
|
||||
[r for r in replica_methods if self.replica_method_can_work_without_replication_connection(r)]
|
||||
# go through them in priority order
|
||||
ret = 1
|
||||
for replica_method in replica_methods:
|
||||
|
||||
+3
-9
@@ -33,7 +33,7 @@ def get_cluster_initialized_without_leader(leader=False, failover=None):
|
||||
'api_url': 'http://127.0.0.1:8008/patroni', 'xlog_location': 4})
|
||||
l = Leader(0, 0, m1) if leader else None
|
||||
m2 = Member(0, 'other', 28, {'conn_url': 'postgres://replicator:[email protected]:5436/postgres',
|
||||
'api_url': 'http://127.0.0.1:8011/patroni'})
|
||||
'api_url': 'http://127.0.0.1:8011/patroni', 'tags': {'clonefrom': True}})
|
||||
return get_cluster(True, l, [m1, m2], failover)
|
||||
|
||||
|
||||
@@ -52,7 +52,7 @@ class MockPatroni(object):
|
||||
self.postgresql = p
|
||||
self.dcs = d
|
||||
self.api = Mock()
|
||||
self.tags = {}
|
||||
self.tags = {'foo': 'bar'}
|
||||
self.nofailover = None
|
||||
self.nap_time = 10
|
||||
self.replicatefrom = None
|
||||
@@ -202,14 +202,8 @@ class TestHa(unittest.TestCase):
|
||||
self.ha.load_cluster_from_dcs = Mock(side_effect=DCSError('Etcd is not responding properly'))
|
||||
self.assertEquals(self.ha.run_cycle(), 'demoted self because DCS is not accessible and i was a leader')
|
||||
|
||||
def test_bootstrap_from_leader(self):
|
||||
self.ha.cluster = get_cluster_initialized_with_leader()
|
||||
self.p.bootstrap = false
|
||||
self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap from leader')
|
||||
|
||||
def test_bootstrap_from_another_member(self):
|
||||
self.ha.cluster = get_cluster_initialized_with_leader()
|
||||
self.ha.patroni.clonefrom = 'other'
|
||||
self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap from replica \'other\'')
|
||||
|
||||
def test_bootstrap_waiting_for_leader(self):
|
||||
@@ -219,7 +213,7 @@ class TestHa(unittest.TestCase):
|
||||
def test_bootstrap_without_leader(self):
|
||||
self.ha.cluster = get_cluster_initialized_without_leader()
|
||||
self.p.can_create_replica_without_replication_connection = MagicMock(return_value=True)
|
||||
self.assertEquals(self.ha.bootstrap(), "trying to bootstrap without leader")
|
||||
self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap (without leader)')
|
||||
|
||||
def test_bootstrap_initialize_lock_failed(self):
|
||||
self.ha.cluster = get_cluster_not_initialized_without_leader()
|
||||
|
||||
Reference in New Issue
Block a user