mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Call touch_member at the end of HA loop (#321)
To make sure that we have up-to-date state of member in DCS after HA loop has changed something.
This commit is contained in:
committed by
GitHub
parent
298357c099
commit
e38dfaf1ba
@@ -200,7 +200,7 @@ class Consul(AbstractDCS):
|
||||
|
||||
def touch_member(self, data, **kwargs):
|
||||
cluster = self.cluster
|
||||
member = cluster and ([m for m in cluster.members if m.name == self._name] or [None])[0]
|
||||
member = cluster and cluster.get_member(self._name, fallback_to_leader=False)
|
||||
create_member = self.refresh_session()
|
||||
if member and (create_member or member.session != self._session):
|
||||
try:
|
||||
@@ -223,6 +223,9 @@ class Consul(AbstractDCS):
|
||||
|
||||
@catch_consul_errors
|
||||
def attempt_to_acquire_leader(self, permanent=False):
|
||||
if not self._session and not permanent:
|
||||
self.refresh_session()
|
||||
|
||||
args = {} if permanent else {'acquire': self._session}
|
||||
ret = self.retry(self._client.kv.put, self.leader_path, self._name, **args)
|
||||
if not ret:
|
||||
|
||||
@@ -232,7 +232,7 @@ class ZooKeeper(AbstractDCS):
|
||||
|
||||
def touch_member(self, data, ttl=None, permanent=False):
|
||||
cluster = self.cluster
|
||||
member = cluster and ([m for m in cluster.members if m.name == self._name] or [None])[0]
|
||||
member = cluster and cluster.get_member(self._name, fallback_to_leader=False)
|
||||
data = data.encode('utf-8')
|
||||
if member and self._client.client_id is not None and member.session != self._client.client_id[0]:
|
||||
try:
|
||||
|
||||
+7
-3
@@ -165,7 +165,6 @@ class Ha(object):
|
||||
return message
|
||||
else:
|
||||
self.state_handler.promote()
|
||||
self.touch_member()
|
||||
return promote_message
|
||||
|
||||
@staticmethod
|
||||
@@ -309,7 +308,6 @@ class Ha(object):
|
||||
self.state_handler.stop()
|
||||
self.state_handler.set_role('demoted')
|
||||
self.dcs.delete_leader()
|
||||
self.touch_member()
|
||||
self.dcs.reset_cluster()
|
||||
sleep(2) # Give a time to somebody to take the leader lock
|
||||
cluster = self.dcs.get_cluster()
|
||||
@@ -581,10 +579,12 @@ class Ha(object):
|
||||
return None
|
||||
|
||||
def _run_cycle(self):
|
||||
dcs_failed = False
|
||||
try:
|
||||
self.load_cluster_from_dcs()
|
||||
|
||||
self.touch_member()
|
||||
if not self.cluster.has_member(self.state_handler.name):
|
||||
self.touch_member()
|
||||
|
||||
# cluster has leader key but not initialize key
|
||||
if not (self.cluster.is_unlocked() or self.sysid_valid(self.cluster.initialize)) and self.has_lock():
|
||||
@@ -645,6 +645,7 @@ class Ha(object):
|
||||
self.state_handler.call_nowait(ACTION_ON_START)
|
||||
self.state_handler.sync_replication_slots(self.cluster)
|
||||
except DCSError:
|
||||
dcs_failed = True
|
||||
logger.error('Error communicating with DCS')
|
||||
if not self.is_paused() and self.state_handler.is_running() and self.state_handler.is_leader():
|
||||
self.demote(delete_leader=False)
|
||||
@@ -652,6 +653,9 @@ class Ha(object):
|
||||
return 'DCS is not accessible'
|
||||
except (psycopg2.Error, PostgresConnectionException):
|
||||
return 'Error communicating with PostgreSQL. Will try again later'
|
||||
finally:
|
||||
if not dcs_failed:
|
||||
self.touch_member()
|
||||
|
||||
def run_cycle(self):
|
||||
with self._async_executor:
|
||||
|
||||
@@ -99,6 +99,8 @@ class TestConsul(unittest.TestCase):
|
||||
|
||||
@patch.object(consul.Consul.KV, 'put', Mock(return_value=False))
|
||||
def test_take_leader(self):
|
||||
self.c.set_ttl(20)
|
||||
self.c.refresh_session = Mock()
|
||||
self.c.take_leader()
|
||||
|
||||
@patch.object(consul.Consul.KV, 'put', Mock(return_value=True))
|
||||
|
||||
Reference in New Issue
Block a user