diff --git a/patroni/__init__.py b/patroni/__init__.py index cc008fda..491edfee 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -81,11 +81,9 @@ class Patroni: logger.info('waiting on DCS') sleep(5) elif self.postgresql.is_running(): - self.postgresql.load_replication_slots() + self.postgresql.schedule_load_slots = True def schedule_next_run(self): - if self.postgresql.is_promoted: - self.next_run = time.time() self.next_run += self.nap_time current_time = time.time() nap_time = self.next_run - current_time @@ -102,10 +100,7 @@ class Patroni: self.touch_member() logger.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() + self.ha.cluster and self.ha.state_handler.sync_replication_slots(self.ha.cluster) except: logger.exception('Exception when changing replication slots') reap_children() diff --git a/patroni/ha.py b/patroni/ha.py index f31b2b26..7725d3c2 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -1,7 +1,7 @@ import logging +import psycopg2 from patroni.dcs import DCSError -from psycopg2 import InterfaceError, OperationalError logger = logging.getLogger(__name__) @@ -56,7 +56,7 @@ class Ha: if self.cluster.is_unlocked(): if self.state_handler.is_healthiest_node(self.old_cluster): if self.acquire_lock(): - if self.state_handler.is_leader() or self.state_handler.is_promoted: + if self.state_handler.is_leader() or self.state_handler.role == 'master': return 'acquired session lock as a leader' else: self.state_handler.promote() @@ -79,7 +79,7 @@ class Ha: return 'following a different leader because i am not the healthiest node' else: if self.has_lock() and self.update_lock(): - if self.state_handler.is_leader() or self.state_handler.is_promoted: + if self.state_handler.is_leader() or self.state_handler.role == 'master': return 'no action. i am the leader with the lock' else: self.state_handler.promote() @@ -97,5 +97,5 @@ class Ha: if self.state_handler.is_leader(): self.state_handler.demote(None) return 'demoted self because DCS is not accessible and i was a leader' - except (InterfaceError, OperationalError): - logger.error('Error communicating with Postgresql. Will try again') + except psycopg2.Error: + logger.exception('Error communicating with Postgresql. Will try again') diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 0706dd2e..c77471d6 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -49,6 +49,7 @@ class Postgresql: self.admin = config['admin'] self.callback = config.get('callbacks', {}) self.use_slots = config.get('use_slots', True) + self.schedule_load_slots = self.use_slots self.recovery_conf = os.path.join(self.data_dir, 'recovery.conf') self.configuration_to_save = (os.path.join(self.data_dir, 'pg_hba.conf'), os.path.join(self.data_dir, 'postgresql.conf')) @@ -56,7 +57,6 @@ class Postgresql: 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._role = 'replica' - self.is_promoted = False self._pg_ctl = ['pg_ctl', '-w', '-D', self.data_dir] @@ -103,9 +103,7 @@ class Postgresql: cursor = self._cursor() cursor.execute(sql, params) return cursor - except psycopg2.InterfaceError as e: - ex = e - except psycopg2.OperationalError as e: + except psycopg2.Error as e: if self._connection and self._connection.closed == 0: raise e ex = e @@ -155,11 +153,7 @@ class Postgresql: return ret def is_leader(self, check_only=False): - ret = not self.query('SELECT pg_is_in_recovery()').fetchone()[0] - if ret and self.is_promoted and not check_only: - self.delete_trigger_file() - self.is_promoted = False - return ret + return not self.query('SELECT pg_is_in_recovery()').fetchone()[0] def is_running(self): return subprocess.call(' '.join(self._pg_ctl) + ' status > /dev/null 2>&1', shell=True) == 0 @@ -170,26 +164,30 @@ class Postgresql: return False cmd = self.callback[cb_name] try: - subprocess.Popen(shlex.split(cmd) + [cb_name, self._role, self.scope]) + subprocess.Popen(shlex.split(cmd) + [cb_name, self.role, self.scope]) except: - logger.exception('callback %s %s %s %s failed', cmd, cb_name, self._role, self.scope) + logger.exception('callback %s %s %s %s failed', cmd, cb_name, self.role, self.scope) return False return True - def start(self, role='replica', block_callbacks=False): + @property + def role(self): + return self._role + + def start(self, block_callbacks=False): if self.is_running(): self._role = 'master' if self.is_leader(check_only=True) else 'replica' - self.load_replication_slots() + self.schedule_load_slots = self.use_slots logger.error('Cannot start PostgreSQL because one is already running.') return False - self._role = role + self._role = 'replica' if os.path.exists(self.recovery_conf) else 'master' if os.path.exists(self.postmaster_pid): os.remove(self.postmaster_pid) logger.info('Removed %s', self.postmaster_pid) ret = subprocess.call(self._pg_ctl + ['start', '-o', self.server_options()]) == 0 - ret and self.load_replication_slots() + self.schedule_load_slots = ret and self.use_slots self.save_configuration_files() # block_callbacks is used during restart to avoid # running start/stop callbacks in addition to restart ones @@ -306,7 +304,7 @@ recovery_target_timeline = 'latest' def follow_the_leader(self, leader): if not self.check_recovery_conf(leader): self.write_recovery_conf(leader) - run_callback = self._role == 'master' + run_callback = self.role == 'master' self.restart() run_callback and self.call_nowait(ACTION_ON_ROLE_CHANGE) @@ -327,11 +325,13 @@ recovery_target_timeline = 'latest' logger.exception('unable to restore configuration from WAL-E backup') def promote(self): - self.is_promoted = subprocess.call(self._pg_ctl + ['promote']) == 0 - if self.is_promoted: + if self.role == 'master': + return True + ret = subprocess.call(self._pg_ctl + ['promote']) == 0 + if ret: self._role = 'master' self.call_nowait(ACTION_ON_ROLE_CHANGE) - return self.is_promoted + return ret def demote(self, leader): self.follow_the_leader(leader) @@ -359,12 +359,15 @@ recovery_target_timeline = 'latest' END, '0/0')""").fetchone()[0] def load_replication_slots(self): - if self.use_slots: + if self.use_slots and self.schedule_load_slots: cursor = self.query("SELECT slot_name FROM pg_replication_slots WHERE slot_type='physical'") self.members = [r[0] for r in cursor] + self.schedule_load_slots = False - def sync_replication_slots(self, members): + def sync_replication_slots(self, cluster): if self.use_slots: + self.load_replication_slots() + members = [m.name for m in cluster.members if m.name != self.name] if self.role == 'master' else [] # drop unused slots for slot in set(self.members) - set(members): self.query("""SELECT pg_drop_replication_slot(%s) @@ -377,13 +380,7 @@ recovery_target_timeline = 'latest' WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots 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([]) + self.members = members def last_operation(self): return str(self.xlog_position()) @@ -400,7 +397,7 @@ recovery_target_timeline = 'latest' """ ret = False if not current_leader: - ret = self.initialize() and self.start(role='master') + ret = self.initialize() and self.start() if ret: self.create_replication_user() self.create_connection_users() diff --git a/tests/test_ha.py b/tests/test_ha.py index 46aa9a0a..d82668e6 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -19,7 +19,7 @@ class MockPostgresql: def __init__(self): self.name = 'postgresql0' - self.is_promoted = False + self.role = 'replica' def is_healthy(self): return True diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 110543da..7e8e768f 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -164,10 +164,10 @@ class TestPostgresql(unittest.TestCase): p = Postgresql(cfg) p.create_connection_users() - def test_create_replication_slots(self): + def test_sync_replication_slots(self): self.p.start() cluster = Cluster(True, self.leader, 0, [self.me, self.other, self.leadermem]) - self.p.create_replication_slots(cluster) + self.p.sync_replication_slots(cluster) def test_query(self): self.p.query('select 1') @@ -191,11 +191,6 @@ class TestPostgresql(unittest.TestCase): self.p.config['maximum_lag_on_failover'] = -3 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()) @@ -206,6 +201,7 @@ class TestPostgresql(unittest.TestCase): def test_promote(self): self.assertTrue(self.p.promote()) + self.assertTrue(self.p.promote()) def test_last_operation(self): self.assertEquals(self.p.last_operation(), '0')