Make work with dcs.cluster thread-safe

This commit is contained in:
Alexander Kukushkin
2015-10-05 14:30:47 +02:00
parent 4c444c943e
commit 601ba7db8d
6 changed files with 104 additions and 81 deletions
+3 -2
View File
@@ -48,8 +48,9 @@ class RestApiHandler(BaseHTTPRequestHandler):
response = self.get_postgresql_status()
patroni = self.server.patroni
if patroni.dcs.cluster: # dcs available
if patroni.dcs.cluster.leader and patroni.dcs.cluster.leader.name == patroni.postgresql.name: # is_leader
cluster = patroni.dcs.cluster
if cluster: # dcs available
if cluster.leader and cluster.leader.name == patroni.postgresql.name: # is_leader
status_code = 200 if 'master' in path else 503
elif 'role' not in response:
status_code = 503
+28 -5
View File
@@ -4,7 +4,7 @@ import json
from collections import namedtuple
from patroni.exceptions import DCSError
from six.moves.urllib_parse import urlparse, urlunparse, parse_qsl
from threading import Event
from threading import Event, Lock
def parse_connection_string(value):
@@ -125,7 +125,8 @@ class AbstractDCS:
self._scope = config['scope']
self._base_path = '/service/' + self._scope
self.cluster = None
self._cluster = None
self._cluster_thread_lock = Lock()
self.event = Event()
def client_path(self, path):
@@ -156,10 +157,32 @@ class AbstractDCS:
return self.client_path(self._LEADER_OPTIME)
@abc.abstractmethod
def _load_cluster(self):
"""Internally this method should build `Cluster` object which
represents current state and topology of the cluster in DCS.
this method supposed to be called only by `get_cluster` method.
raise `~DCSError` in case of communication or other problems with DCS.
If the current node was running as a master and exception raised,
instance would be demoted."""
def get_cluster(self):
""":returns: `Cluster` object which represent current state and topology of the cluster
raise `~DCSError` in case of communication or other problems with DCS. If current instance was
running as a master and exception raised instance would be demoted."""
with self._cluster_thread_lock:
try:
self._load_cluster()
except:
self._cluster = None
raise
return self._cluster
@property
def cluster(self):
with self._cluster_thread_lock:
return self._cluster
def reset_cluster(self):
with self._cluster_thread_lock:
self._cluster = None
@abc.abstractmethod
def write_leader_optime(self, last_operation):
+6 -7
View File
@@ -171,7 +171,7 @@ class Etcd(AbstractDCS):
def member(node):
return Member.from_node(node.modifiedIndex, os.path.basename(node.key), node.ttl, node.value)
def get_cluster(self):
def _load_cluster(self):
try:
result = self.retry(self.client.read, self.client_path(''), recursive=True)
nodes = {os.path.relpath(node.key, result.key): node for node in result.leaves}
@@ -198,14 +198,12 @@ class Etcd(AbstractDCS):
if failover:
failover = Failover.from_node(failover.modifiedIndex, failover.value)
self.cluster = Cluster(initialize, leader, last_leader_operation, members, failover)
self._cluster = Cluster(initialize, leader, last_leader_operation, members, failover)
except etcd.EtcdKeyNotFound:
self.cluster = Cluster(False, None, None, [], None)
self._cluster = Cluster(False, None, None, [], None)
except:
self.cluster = None
logger.exception('get_cluster')
raise EtcdError('Etcd is not responding properly')
return self.cluster
@catch_etcd_errors
def touch_member(self, connection_string, ttl=None):
@@ -249,10 +247,11 @@ class Etcd(AbstractDCS):
return self.retry(self.client.delete, self.initialize_path, prevValue=self._name)
def watch(self, timeout):
cluster = self.cluster
# watch on leader key changes if it is defined and current node is not lock owner
if self.cluster and self.cluster.leader and self.cluster.leader.name != self._name:
if cluster and cluster.leader and cluster.leader.name != self._name:
end_time = time.time() + timeout
index = self.cluster.leader.index
index = cluster.leader.index
while index and timeout >= 1: # when timeout is too small urllib3 doesn't have enough time to connect
try:
+42 -39
View File
@@ -16,13 +16,17 @@ class Ha:
self.patroni = patroni
self.state_handler = patroni.postgresql
self.dcs = patroni.dcs
self.cluster = None
self.old_cluster = None
self._async_executor = AsyncExecutor()
def load_cluster_from_dcs(self):
cluster = self.dcs.get_cluster()
# We want to keep the state of cluster when it was healhy
if not self.dcs.get_cluster().is_unlocked() or not self.old_cluster:
self.old_cluster = self.dcs.cluster
if not cluster.is_unlocked() or not self.old_cluster:
self.old_cluster = cluster
self.cluster = cluster
def acquire_lock(self):
return self.dcs.attempt_to_acquire_leader()
@@ -37,7 +41,7 @@ class Ha:
return ret
def has_lock(self):
lock_owner = self.dcs.cluster.leader and self.dcs.cluster.leader.name
lock_owner = self.cluster.leader and self.cluster.leader.name
logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name)
return lock_owner == self.state_handler.name
@@ -55,8 +59,8 @@ class Ha:
pass
self.dcs.touch_member(json.dumps(data, separators=(',', ':')))
def copy_backup_from_leader(self):
if self.state_handler.bootstrap(self.dcs.cluster.leader):
def copy_backup_from_leader(self, leader):
if self.state_handler.bootstrap(leader):
logger.info('bootstrapped from leader')
else:
self.state_handler.stop('immediate')
@@ -64,14 +68,11 @@ class Ha:
logger.error('failed to bootstrap from leader')
def bootstrap(self):
if not self.dcs.cluster.is_unlocked(): # cluster already has leader
if self._async_executor.busy:
self.copy_backup_from_leader()
else:
self._async_executor.schedule('bootstrap from leader')
self._async_executor.run_async(self.copy_backup_from_leader)
if not self.cluster.is_unlocked(): # cluster already has leader
self._async_executor.schedule('bootstrap from leader')
self._async_executor.run_async(self.copy_backup_from_leader, args=(self.cluster.leader, ))
return 'trying to bootstrap from leader'
elif not self.dcs.cluster.initialize: # no initialize key
elif not self.cluster.initialize: # no initialize key
if self.dcs.initialize(): # race for initialization
try:
self.state_handler.bootstrap()
@@ -91,11 +92,12 @@ class Ha:
def recover(self):
has_lock = self.has_lock()
self.state_handler.write_recovery_conf(None if has_lock else self.dcs.cluster.leader)
self.state_handler.write_recovery_conf(None if has_lock else self.cluster.leader)
if not self.state_handler.start():
if not has_lock:
return 'failed to start postgres'
self.dcs.delete_leader()
self.dcs.reset_cluster()
return 'removed leader key after trying and failing to start postgres'
if not has_lock:
return 'started as a secondary'
@@ -105,9 +107,11 @@ class Ha:
def follow_the_leader(self, demote_reason, follow_reason, refresh=True):
refresh and self.load_cluster_from_dcs()
ret = demote_reason if self.state_handler.is_leader() else follow_reason
if not self.state_handler.check_recovery_conf(self.dcs.cluster.leader):
if not self.state_handler.check_recovery_conf(self.cluster.leader):
self._async_executor.schedule('changing primary_conninfo and restarting')
self._async_executor.run_async(self.state_handler.follow_the_leader, (self.dcs.cluster.leader, ))
leader = self.cluster.leader
leader = None if (leader and leader.name) == self.state_handler.name else leader
self._async_executor.run_async(self.state_handler.follow_the_leader, (leader, ))
return ret
def enforce_master_role(self, message, promote_message):
@@ -150,8 +154,7 @@ class Ha:
if self.state_handler.is_leader():
return True
if check_replication_lag and \
not self.state_handler.check_replication_lag(self.dcs.cluster.last_leader_operation):
if check_replication_lag and not self.state_handler.check_replication_lag(self.cluster.last_leader_operation):
return False # Too far behind last reported xlog location on master
# Prepare list of nodes to run check against
@@ -182,13 +185,13 @@ class Ha:
return ret
def manual_failover_process_no_leader(self):
failover = self.dcs.cluster.failover
failover = self.cluster.failover
if failover.member: # manual failover to specific member
if failover.member == self.state_handler.name: # manual failover to me
return True
# find specific node and check that it is healthy
members = [m for m in self.dcs.cluster.members if m.name == failover.member]
members = [m for m in self.cluster.members if m.name == failover.member]
if members:
member, reachable, in_recovery, xlog_location = self.fetch_node_status(members[0])
if reachable: # node is healthy
@@ -204,7 +207,7 @@ class Ha:
if failover.leader:
if self.state_handler.name == failover.leader: # I was the leader
# exclude me and desired member which is unhealthy (failover.member can be None)
members = [m for m in self.dcs.cluster.members if m.name != failover.member]
members = [m for m in self.cluster.members if m.name != failover.member]
if self.is_failover_possible(members): # check that there are healthy members
return False
else: # I was the leader and it looks like currently I am the only healthy member
@@ -213,28 +216,29 @@ class Ha:
# at this point we assume that our node is a candidate for a failover among all nodes except former leader
# exclude former leader from the list (failover.leader can be None)
members = [m for m in self.dcs.cluster.members if m.name != failover.leader]
members = [m for m in self.cluster.members if m.name != failover.leader]
return self._is_healthiest_node(members, check_replication_lag=False)
def is_healthiest_node(self):
if self.dcs.cluster.failover:
if self.cluster.failover:
return self.manual_failover_process_no_leader()
# run usual health check
members = {m.name: m for m in self.dcs.cluster.members + self.old_cluster.members}
members = {m.name: m for m in self.cluster.members + self.old_cluster.members}
return self._is_healthiest_node(members.values())
def demote(self, delete_leader=True):
if delete_leader:
self.state_handler.stop()
self.dcs.delete_leader()
self.dcs.reset_cluster()
self.state_handler.follow_the_leader(None)
def process_manual_failover_from_leader(self):
failover = self.dcs.cluster.failover
failover = self.cluster.failover
if not failover.leader or failover.leader == self.state_handler.name:
if not failover.member or failover.member != self.state_handler.name:
members = [m for m in self.dcs.cluster.members if not failover.member or m.name == failover.member]
members = [m for m in self.cluster.members if not failover.member or m.name == failover.member]
if self.is_failover_possible(members): # check that there are healthy members
self._async_executor.schedule('manual failover: demote')
self._async_executor.run_async(self.demote)
@@ -245,15 +249,15 @@ class Ha:
logger.warning('manual failover: I am already the leader, no need to failover')
else:
logger.warning('manual failover: leader name does not match: %s != %s',
self.dcs.cluster.failover.leader, self.state_handler.name)
self.cluster.failover.leader, self.state_handler.name)
logger.info('Trying to clean up failover key')
self.dcs.manual_failover('', '', self.dcs.cluster.failover.index)
self.dcs.manual_failover('', '', self.cluster.failover.index)
def process_unhealthy_cluster(self):
if self.is_healthiest_node():
if self.acquire_lock():
if self.dcs.cluster.failover:
if self.cluster.failover:
logger.info('Cleanning up failover key after acquiring leader lock...')
self.dcs.manual_failover('', '')
return self.enforce_master_role('acquired session lock as a leader',
@@ -267,7 +271,7 @@ class Ha:
def process_healthy_cluster(self):
if self.has_lock():
if self.dcs.cluster.failover:
if self.cluster.failover:
msg = self.process_manual_failover_from_leader()
if msg is not None:
return msg
@@ -307,22 +311,21 @@ class Ha:
else:
return (False, 'restart failed')
def reinitialize(self):
def reinitialize(self, cluster):
self.state_handler.stop('immediate')
self.state_handler.remove_data_directory()
self.load_cluster_from_dcs()
self.bootstrap()
self.copy_backup_from_leader(cluster.leader)
def process_scheduled_action(self):
if self.reinitialize_scheduled():
if self.dcs.cluster.is_unlocked():
if self.cluster.is_unlocked():
logger.error('Cluster has no leader, can not reinitialize')
self._async_executor.reset_scheduled_action()
elif self.has_lock():
logger.error('I am the leader, can not reinitialize')
self._async_executor.reset_scheduled_action()
else:
self._async_executor.run_async(self.reinitialize)
self._async_executor.run_async(self.reinitialize, args=(self.cluster, ))
return 'reinitialize started'
def handle_long_action_in_progress(self):
@@ -331,7 +334,7 @@ class Ha:
return 'updated leader lock during ' + self._async_executor.scheduled_action
else:
return 'failed to update leader lock during ' + self._async_executor.scheduled_action
elif self.dcs.cluster.is_unlocked():
elif self.cluster.is_unlocked():
return 'not healthy enough for leader race'
else:
return self._async_executor.scheduled_action + ' in progress'
@@ -343,7 +346,7 @@ class Ha:
self.touch_member()
# cluster has leader key but not initialize key
if not self.dcs.cluster.is_unlocked() and not self.dcs.cluster.initialize:
if not self.cluster.is_unlocked() and not self.cluster.initialize:
self.dcs.initialize() # fix it
if self._async_executor.busy:
@@ -358,7 +361,7 @@ class Ha:
if self.state_handler.data_directory_empty():
return self.bootstrap() # new node
# "bootstrap", but data directory is not empty
elif not self.dcs.cluster.initialize and self.dcs.cluster.is_unlocked():
elif not self.cluster.initialize and self.cluster.is_unlocked():
self.dcs.initialize()
# try to start dead postgres
@@ -368,12 +371,12 @@ class Ha:
return msg
try:
if self.dcs.cluster.is_unlocked():
if self.cluster.is_unlocked():
return self.process_unhealthy_cluster()
else:
return self.process_healthy_cluster()
finally:
self.state_handler.sync_replication_slots(self.dcs.cluster)
self.state_handler.sync_replication_slots(self.cluster)
except DCSError:
logger.error('Error communicating with DCS')
if self.state_handler.is_running() and self.state_handler.is_leader():
+4 -5
View File
@@ -164,9 +164,9 @@ class ZooKeeper(AbstractDCS):
# get last leader operation
self.last_leader_operation = self.get_node(self.leader_optime_path) if self.fetch_cluster else None
self.last_leader_operation = 0 if self.last_leader_operation is None else int(self.last_leader_operation[0])
self.cluster = Cluster(initialize, leader, self.last_leader_operation, members, failover)
self._cluster = Cluster(initialize, leader, self.last_leader_operation, members, failover)
def get_cluster(self):
def _load_cluster(self):
if self.exhibitor and self.exhibitor.poll():
self.client.set_hosts(self.exhibitor.zookeeper_hosts)
@@ -174,11 +174,9 @@ class ZooKeeper(AbstractDCS):
try:
self.client.retry(self._inner_load_cluster)
except:
self.cluster = None
logger.exception('get_cluster')
self.session_listener(KazooState.LOST)
raise ZooKeeperError('ZooKeeper in not responding properly')
return self.cluster
def _create(self, path, value, **kwargs):
try:
@@ -206,7 +204,8 @@ class ZooKeeper(AbstractDCS):
return self._create(self.initialize_path, self._name, makepath=True)
def touch_member(self, data, ttl=None):
me = self.cluster and ([m for m in self.cluster.members if m.name == self._name] or [None])[0]
cluster = self.cluster
me = cluster and ([m for m in cluster.members if m.name == self._name] or [None])[0]
path = self.member_path
data = data.encode('utf-8')
create = not me
+21 -23
View File
@@ -85,7 +85,7 @@ class MockPatroni:
def run_async(func, args=()):
func(args) if args else func()
func(*args) if args else func()
class TestHa(unittest.TestCase):
@@ -101,7 +101,7 @@ class TestHa(unittest.TestCase):
self.ha = Ha(MockPatroni(self.p, self.e))
self.ha._async_executor.run_async = run_async
self.ha.old_cluster = self.e.get_cluster()
self.e.cluster = get_cluster_not_initialized_without_leader()
self.ha.cluster = get_cluster_not_initialized_without_leader()
self.ha.load_cluster_from_dcs = Mock()
def test_update_lock(self):
@@ -161,28 +161,28 @@ class TestHa(unittest.TestCase):
self.assertEquals(self.ha.run_cycle(), 'following a different leader because i am not the healthiest node')
def test_promote_because_have_lock(self):
self.e.cluster.is_unlocked = false
self.ha.cluster.is_unlocked = false
self.ha.has_lock = true
self.p.is_leader = false
self.assertEquals(self.ha.run_cycle(), 'promoted self to leader because i had the session lock')
def test_leader_with_lock(self):
self.e.cluster.is_unlocked = false
self.ha.cluster.is_unlocked = false
self.ha.has_lock = true
self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock')
def test_demote_because_not_having_lock(self):
self.e.cluster.is_unlocked = false
self.ha.cluster.is_unlocked = false
self.assertEquals(self.ha.run_cycle(), 'demoting self because i do not have the lock and i was a leader')
def test_demote_because_update_lock_failed(self):
self.e.cluster.is_unlocked = false
self.ha.cluster.is_unlocked = false
self.ha.has_lock = true
self.ha.update_lock = false
self.assertEquals(self.ha.run_cycle(), 'demoting self because i do not have the lock and i was a leader')
def test_follow_the_leader(self):
self.e.cluster.is_unlocked = false
self.ha.cluster.is_unlocked = false
self.p.is_leader = false
self.assertEquals(self.ha.run_cycle(), 'no action. i am a secondary and i am following a leader')
@@ -191,27 +191,25 @@ class TestHa(unittest.TestCase):
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.e.cluster = get_cluster_initialized_with_leader()
self.ha.cluster = get_cluster_initialized_with_leader()
self.p.bootstrap = false
self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap from leader')
self.ha._async_executor._busy = True
self.assertEquals(self.ha.bootstrap(), 'trying to bootstrap from leader')
def test_bootstrap_waiting_for_leader(self):
self.e.cluster = get_cluster_initialized_without_leader()
self.ha.cluster = get_cluster_initialized_without_leader()
self.assertEquals(self.ha.bootstrap(), 'waiting for leader to bootstrap')
def test_bootstrap_initialize_lock_failed(self):
self.e.cluster = get_cluster_not_initialized_without_leader()
self.ha.cluster = get_cluster_not_initialized_without_leader()
self.assertEquals(self.ha.bootstrap(), 'failed to acquire initialize lock')
def test_bootstrap_initialized_new_cluster(self):
self.e.cluster = get_cluster_not_initialized_without_leader()
self.ha.cluster = get_cluster_not_initialized_without_leader()
self.e.initialize = true
self.assertEquals(self.ha.bootstrap(), 'initialized a new cluster')
def test_bootstrap_release_initialize_key_on_failure(self):
self.e.cluster = get_cluster_not_initialized_without_leader()
self.ha.cluster = get_cluster_not_initialized_without_leader()
self.e.initialize = true
self.p.bootstrap = Mock(side_effect=PostgresException("Could not bootstrap master PostgreSQL"))
self.assertRaises(PostgresException, self.ha.bootstrap)
@@ -222,7 +220,7 @@ class TestHa(unittest.TestCase):
self.ha.run_cycle()
self.assertIsNone(self.ha._async_executor.scheduled_action)
self.e.cluster = get_cluster_initialized_with_leader()
self.ha.cluster = get_cluster_initialized_with_leader()
self.ha.has_lock = true
self.ha.schedule_reinitialize()
self.ha.run_cycle()
@@ -244,7 +242,7 @@ class TestHa(unittest.TestCase):
self.assertTrue(self.ha.restart_scheduled())
self.assertEquals(self.ha.run_cycle(), 'not healthy enough for leader race')
self.e.cluster = get_cluster_initialized_with_leader()
self.ha.cluster = get_cluster_initialized_with_leader()
self.assertEquals(self.ha.run_cycle(), 'restart in progress')
self.ha.has_lock = true
@@ -256,26 +254,26 @@ class TestHa(unittest.TestCase):
@patch('requests.get', requests_get)
def test_manual_failover_from_leader(self):
self.ha.has_lock = true
self.e.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', ''))
self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, 'blabla', ''))
self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock')
self.e.cluster = get_cluster_initialized_with_leader(Failover(0, '', MockPostgresql.name))
self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', MockPostgresql.name))
self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock')
self.e.cluster = get_cluster_initialized_with_leader(Failover(0, '', 'blabla'))
self.ha.cluster = get_cluster_initialized_with_leader(Failover(0, '', 'blabla'))
self.assertEquals(self.ha.run_cycle(), 'no action. i am the leader with the lock')
f = Failover(0, MockPostgresql.name, '')
self.e.cluster = get_cluster_initialized_with_leader(f)
self.ha.cluster = get_cluster_initialized_with_leader(f)
self.assertEquals(self.ha.run_cycle(), 'manual failover: demoting myself')
@patch('requests.get', requests_get)
def test_manual_failover_process_no_leader(self):
self.p.is_leader = false
self.e.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', MockPostgresql.name))
self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', MockPostgresql.name))
self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock')
self.e.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'leader'))
self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, '', 'leader'))
self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock')
self.ha.fetch_node_status = lambda e: (e, True, True, 0) # accessible, in_recovery
self.assertEquals(self.ha.run_cycle(), 'following a different leader because i am not the healthiest node')
self.e.cluster = get_cluster_initialized_without_leader(failover=Failover(0, MockPostgresql.name, ''))
self.ha.cluster = get_cluster_initialized_without_leader(failover=Failover(0, MockPostgresql.name, ''))
self.assertEquals(self.ha.run_cycle(), 'following a different leader because i am not the healthiest node')
self.ha.fetch_node_status = lambda e: (e, False, True, 0) # accessible, in_recovery
self.assertEquals(self.ha.run_cycle(), 'promoted self to leader by acquiring session lock')