Try to avoid "double" promotion.

Also check presence of trigger_file on master after promotion when
pg_is_in_recovery() = false and if it is there - remove it.
Plus check presence of trigger_file on a new slave (after running
pg_basebackup) and if it is there - also remove it.
This commit is contained in:
Alexander Kukushkin
2015-06-01 13:48:44 +02:00
parent 8c4e547d7a
commit adcc7ac256
4 changed files with 28 additions and 12 deletions
+7 -8
View File
@@ -48,10 +48,10 @@ class Ha:
if self.cluster.is_unlocked():
if self.state_handler.is_healthiest_node(self.cluster):
if self.acquire_lock():
if not self.state_handler.is_leader():
self.state_handler.promote()
return 'promoted self to leader by acquiring session lock'
return 'acquired session lock as a leader'
if self.state_handler.is_leader() or self.state_handler.is_promoted:
return 'acquired session lock as a leader'
self.state_handler.promote()
return 'promoted self to leader by acquiring session lock'
else:
self.load_cluster_from_etcd()
if self.state_handler.is_leader():
@@ -71,11 +71,10 @@ class Ha:
else:
if self.has_lock() and self.update_lock():
try:
if not self.state_handler.is_leader():
self.state_handler.promote()
return 'promoted self to leader because i had the session lock'
else:
if self.state_handler.is_leader() or self.state_handler.is_promoted:
return 'no action. i am the leader with the lock'
self.state_handler.promote()
return 'promoted self to leader because i had the session lock'
finally:
# create replication slots
self.state_handler.create_replication_slots(self.cluster)
+18 -4
View File
@@ -41,6 +41,10 @@ class Postgresql:
self.replication = config['replication']
self.recovery_conf = os.path.join(self.data_dir, 'recovery.conf')
self.pid_path = os.path.join(self.data_dir, 'postmaster.pid')
self.trigger_file = config.get('recovery_conf', {}).get('trigger_file', None) or 'promote'
self.trigger_file = os.path.abspath(os.path.join(self.data_dir, self.trigger_file))
self.is_promoted = False
self._pg_ctl = ['pg_ctl', '-w', '-D', self.data_dir]
self.local_address = self.get_local_address()
@@ -101,6 +105,9 @@ class Postgresql:
ret and self.write_pg_hba()
return ret
def delete_trigger_file(self):
os.path.exists(self.trigger_file) and os.unlink(self.trigger_file)
def sync_from_leader(self, leader):
r = parseurl(leader.address)
@@ -111,11 +118,17 @@ class Postgresql:
env = os.environ.copy()
env['PGPASSFILE'] = pgpass
return subprocess.call(['pg_basebackup', '-R', '-D', self.data_dir, '--host=' + r['host'],
'--port=' + str(r['port']), '-U', r['user']], env=env) == 0
ret = subprocess.call(['pg_basebackup', '-R', '-D', self.data_dir, '--host=' + r['host'],
'--port=' + str(r['port']), '-U', r['user']], env=env) == 0
self.delete_trigger_file()
return ret
def is_leader(self):
return not self.query('SELECT pg_is_in_recovery()').fetchone()[0]
ret = not self.query('SELECT pg_is_in_recovery()').fetchone()[0]
if ret and self.is_promoted:
self.delete_trigger_file()
self.is_promoted = False
return ret
def is_running(self):
return subprocess.call(' '.join(self._pg_ctl) + ' status > /dev/null', shell=True) == 0
@@ -226,7 +239,8 @@ primary_conninfo = '{}'
self.restart()
def promote(self):
return subprocess.call(self._pg_ctl + ['promote']) == 0
self.is_promoted = subprocess.call(self._pg_ctl + ['promote']) == 0
return self.is_promoted
def demote(self, leader):
self.follow_the_leader(leader)
+1
View File
@@ -19,6 +19,7 @@ class MockPostgresql:
def __init__(self):
self.name = 'postgresql0'
self.is_promoted = False
def is_healthy(self):
return True
+2
View File
@@ -160,7 +160,9 @@ class TestPostgresql(unittest.TestCase):
self.assertFalse(self.p.is_healthiest_node(cluster))
def test_is_leader(self):
self.p.is_promoted = True
self.assertTrue(self.p.is_leader())
self.assertFalse(self.p.is_promoted)
def test_reload(self):
self.assertTrue(self.p.reload())