From 9732106e5d1a4dfb7695df12b367ca34e557fa07 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 25 Jun 2015 14:01:06 +0200 Subject: [PATCH] BUGFIX: keep list of replication slots on slaves empty Otherwise slaves accumulate wal files... --- governor.py | 8 +++++++- helpers/ha.py | 14 ++++++-------- helpers/postgresql.py | 9 +++++++-- tests/test_governor.py | 8 +++++++- 4 files changed, 27 insertions(+), 12 deletions(-) diff --git a/governor.py b/governor.py index d8a1e620..eb029b22 100755 --- a/governor.py +++ b/governor.py @@ -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() diff --git a/helpers/ha.py b/helpers/ha.py index 0453686e..4f0e61bc 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -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(): diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 861074e8..8d615a0c 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -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() diff --git a/tests/test_governor.py b/tests/test_governor.py index b84637a7..b850deda 100644 --- a/tests/test_governor.py +++ b/tests/test_governor.py @@ -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