mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
BUGFIX: keep list of replication slots on slaves empty
Otherwise slaves accumulate wal files...
This commit is contained in:
+7
-1
@@ -74,7 +74,13 @@ class Governor:
|
||||
while True:
|
||||
self.touch_member()
|
||||
logging.info(self.ha.run_cycle())
|
||||
|
||||
try:
|
||||
if self.ha.state_handler.is_leader():
|
||||
self.ha.cluster and self.ha.state_handler.create_replication_slots(self.ha.cluster)
|
||||
else:
|
||||
self.ha.state_handler.drop_replication_slots()
|
||||
except:
|
||||
logging.exception('Exception when changing replication slots')
|
||||
self.schedule_next_run()
|
||||
|
||||
|
||||
|
||||
+6
-8
@@ -50,8 +50,9 @@ class Ha:
|
||||
if self.acquire_lock():
|
||||
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.state_handler.promote()
|
||||
return 'promoted self to leader by acquiring session lock'
|
||||
else:
|
||||
self.load_cluster_from_etcd()
|
||||
if self.state_handler.is_leader():
|
||||
@@ -70,14 +71,11 @@ class Ha:
|
||||
return 'following a different leader because i am not the healthiest node'
|
||||
else:
|
||||
if self.has_lock() and self.update_lock():
|
||||
try:
|
||||
if self.state_handler.is_leader() or self.state_handler.is_promoted:
|
||||
return 'no action. i am the leader with the lock'
|
||||
if self.state_handler.is_leader() or self.state_handler.is_promoted:
|
||||
return 'no action. i am the leader with the lock'
|
||||
else:
|
||||
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)
|
||||
else:
|
||||
logger.info('does not have lock')
|
||||
if self.state_handler.is_leader():
|
||||
|
||||
@@ -268,8 +268,7 @@ primary_conninfo = '{}'
|
||||
cursor = self.query("SELECT slot_name FROM pg_replication_slots WHERE slot_type='physical'")
|
||||
self.members = [r[0] for r in cursor]
|
||||
|
||||
def create_replication_slots(self, cluster):
|
||||
members = [m.name for m in cluster.members if m.name != self.name]
|
||||
def sync_replication_slots(self, members):
|
||||
# drop unused slots
|
||||
for slot in set(self.members) - set(members):
|
||||
self.query("""SELECT pg_drop_replication_slot(%s)
|
||||
@@ -283,5 +282,11 @@ primary_conninfo = '{}'
|
||||
WHERE slot_name = %s)""", slot, slot)
|
||||
self.members = members
|
||||
|
||||
def create_replication_slots(self, cluster):
|
||||
self.sync_replication_slots([m.name for m in cluster.members if m.name != self.name])
|
||||
|
||||
def drop_replication_slots(self):
|
||||
self.sync_replication_slots([])
|
||||
|
||||
def last_operation(self):
|
||||
return self.xlog_position()
|
||||
|
||||
@@ -23,7 +23,7 @@ def nop(*args, **kwargs):
|
||||
pass
|
||||
|
||||
|
||||
def time_sleep(_):
|
||||
def time_sleep(*args):
|
||||
raise Exception()
|
||||
|
||||
|
||||
@@ -63,6 +63,12 @@ class TestGovernor(unittest.TestCase):
|
||||
time.sleep = time_sleep
|
||||
self.assertRaises(Exception, main)
|
||||
|
||||
def test_governor_run(self):
|
||||
time.sleep = time_sleep
|
||||
self.g.postgresql.is_leader = lambda: False
|
||||
self.g.ha.state_handler.sync_replication_slots = time_sleep
|
||||
self.assertRaises(Exception, self.g.run)
|
||||
|
||||
def touch_member(self):
|
||||
if not self.touched:
|
||||
self.touched = True
|
||||
|
||||
Reference in New Issue
Block a user