From 80f92b1dee24a5246fbaf22bfb6b7a1b56e82368 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Wed, 2 Sep 2015 14:14:16 +0200 Subject: [PATCH 1/4] Run CHECKPOINT before calling shutdown. In addition, restart is now performed as stop/start, which would allow it to benefit from the shutdown speedup. The hooks in start/stop are modified in order not to run when called as a part of restart. --- helpers/postgresql.py | 24 ++++++++++++++++++------ 1 file changed, 18 insertions(+), 6 deletions(-) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 695e45c1..11029a26 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -183,7 +183,7 @@ class Postgresql: return False return True - def start(self): + def start(self, block_callbacks=False): if self.is_running(): self.load_replication_slots() logger.error('Cannot start PostgreSQL because one is already running.') @@ -196,18 +196,28 @@ class Postgresql: ret = subprocess.call(self._pg_ctl + ['start', '-o', self.server_options()]) == 0 ret and self.load_replication_slots() self.save_configuration_files() - if ret and ACTION_ON_START in self.callback: + # block_callbacks is used during restart to avoid + # running start/stop callbacks in addition to restart ones + if not block_callbacks and ret and ACTION_ON_START in self.callback: self.call_nowait(ACTION_ON_START) return ret - def stop(self): + def stop(self, block_callbacks=False): try: is_leader = self.is_leader(check_only=True) except: is_leader = None pass - ret = subprocess.call(self._pg_ctl + ['stop', '-m', 'fast']) - if ret == 0 and ACTION_ON_STOP in self.callback: + try: + self.query("CHECKPOINT") + except psycopg2.OperationalError: + # likely PostgreSQL is already stopped. + ret = 0 + else: + ret = subprocess.call(self._pg_ctl + ['stop', '-m', 'fast']) + # block_callbacks is used during restart to avoid + # running start/stop callbacks in addition to restart ones + if not block_callbacks and ret == 0 and ACTION_ON_STOP in self.callback: self.call_nowait(ACTION_ON_STOP, is_leader=is_leader) return ret == 0 @@ -223,7 +233,9 @@ class Postgresql: except: is_leader = None pass - ret = subprocess.call(self._pg_ctl + ['restart', '-m', 'fast']) + ret = self.stop(block_callbacks=True) + if ret == 0: + ret = self.start(block_callbacks=True) if ret == 0 and ACTION_ON_RESTART in self.callback: self.call_nowait(ACTION_ON_RESTART, is_leader=is_leader) return ret == 0 From dbcc5aff9bbab560b8538e7aee3d91dd3d4ae1dd Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 16 Sep 2015 16:22:10 +0200 Subject: [PATCH 2/4] Track postgresql role in a Postgresql class --- patroni/postgresql.py | 76 +++++++++++++++------------------------- tests/test_patroni.py | 4 +-- tests/test_postgresql.py | 29 +++++++++------ 3 files changed, 49 insertions(+), 60 deletions(-) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 22286b19..0706dd2e 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -55,6 +55,7 @@ class Postgresql: self.postmaster_pid = 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._role = 'replica' self.is_promoted = False self._pg_ctl = ['pg_ctl', '-w', '-D', self.data_dir] @@ -163,31 +164,26 @@ class Postgresql: def is_running(self): return subprocess.call(' '.join(self._pg_ctl) + ' status > /dev/null 2>&1', shell=True) == 0 - def call_nowait(self, cb_name, is_leader=None): + def call_nowait(self, cb_name): """ pick a callback command and call it without waiting for it to finish """ if not self.callback or cb_name not in self.callback: return False cmd = self.callback[cb_name] - if is_leader is None: - try: - is_leader = self.is_leader(check_only=True) - except psycopg2.OperationalError as e: - logger.warning("unable to perform {0} action, cannot obtain the cluster role: {1}".format(cb_name, e)) - return False try: - role = "master" if is_leader else "replica" - subprocess.Popen(shlex.split(cmd) + [cb_name, 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, 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, block_callbacks=False): + def start(self, role='replica', block_callbacks=False): if self.is_running(): + self._role = 'master' if self.is_leader(check_only=True) else 'replica' self.load_replication_slots() logger.error('Cannot start PostgreSQL because one is already running.') return False + self._role = role if os.path.exists(self.postmaster_pid): os.remove(self.postmaster_pid) logger.info('Removed %s', self.postmaster_pid) @@ -197,47 +193,32 @@ class Postgresql: self.save_configuration_files() # block_callbacks is used during restart to avoid # running start/stop callbacks in addition to restart ones - if not block_callbacks and ret and ACTION_ON_START in self.callback: - self.call_nowait(ACTION_ON_START) + ret and not block_callbacks and ret and self.call_nowait(ACTION_ON_START) return ret def stop(self, block_callbacks=False): - try: - is_leader = self.is_leader(check_only=True) - except: - is_leader = None - pass - try: - self.query("CHECKPOINT") - except psycopg2.OperationalError: - # likely PostgreSQL is already stopped. - ret = 0 - else: - ret = subprocess.call(self._pg_ctl + ['stop', '-m', 'fast']) + if block_callbacks: + try: + self.query('SET statement_timeout TO 0') + self.query('CHECKPOINT') + except: + logging.exception('Exception diring CHECKPOINT') + + ret = subprocess.call(self._pg_ctl + ['stop', '-m', 'fast']) == 0 # block_callbacks is used during restart to avoid # running start/stop callbacks in addition to restart ones - if not block_callbacks and ret == 0 and ACTION_ON_STOP in self.callback: - self.call_nowait(ACTION_ON_STOP, is_leader=is_leader) - return ret == 0 + ret and not block_callbacks and self.call_nowait(ACTION_ON_STOP) + return ret def reload(self): - ret = subprocess.call(self._pg_ctl + ['reload']) - if ret == 0 and ACTION_ON_RELOAD in self.callback: - self.call_nowait(ACTION_ON_RELOAD) - return ret == 0 + ret = subprocess.call(self._pg_ctl + ['reload']) == 0 + ret and self.call_nowait(ACTION_ON_RELOAD) + return ret def restart(self): - try: - is_leader = self.is_leader(check_only=True) - except: - is_leader = None - pass - ret = self.stop(block_callbacks=True) - if ret == 0: - ret = self.start(block_callbacks=True) - if ret == 0 and ACTION_ON_RESTART in self.callback: - self.call_nowait(ACTION_ON_RESTART, is_leader=is_leader) - return ret == 0 + ret = self.stop(block_callbacks=True) and self.start(block_callbacks=True) + ret and self.call_nowait(ACTION_ON_RESTART) + return ret def server_options(self): options = "--listen_addresses='{}' --port={}".format(self.listen_addresses, self.port) @@ -325,9 +306,9 @@ 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' self.restart() - if ACTION_ON_ROLE_CHANGE in self.callback: - self.call_nowait(ACTION_ON_ROLE_CHANGE) + run_callback and self.call_nowait(ACTION_ON_ROLE_CHANGE) def save_configuration_files(self): """ @@ -347,7 +328,8 @@ recovery_target_timeline = 'latest' def promote(self): self.is_promoted = subprocess.call(self._pg_ctl + ['promote']) == 0 - if self.is_promoted and ACTION_ON_ROLE_CHANGE in self.callback: + if self.is_promoted: + self._role = 'master' self.call_nowait(ACTION_ON_ROLE_CHANGE) return self.is_promoted @@ -418,7 +400,7 @@ recovery_target_timeline = 'latest' """ ret = False if not current_leader: - ret = self.initialize() and self.start() + ret = self.initialize() and self.start(role='master') if ret: self.create_replication_user() self.create_connection_users() diff --git a/tests/test_patroni.py b/tests/test_patroni.py index 317e1e5c..bdba1fa0 100644 --- a/tests/test_patroni.py +++ b/tests/test_patroni.py @@ -61,7 +61,7 @@ def get_cluster_not_initialized_with_leader(): def get_cluster_initialized_with_leader(): - return get_cluster(True, Leader(0, 0, 0, + return get_cluster(True, Leader(0, 0, 0, Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', None, None, 28))) @@ -132,7 +132,7 @@ class TestPatroni(unittest.TestCase): self.p.ha.dcs.client.read = etcd_read self.p.ha.dcs.watch = time_sleep self.assertRaises(SleepException, self.p.run) - self.p.ha.state_handler.is_leader = lambda: False + self.p.ha.state_handler.is_leader = false self.p.api.start = nop self.assertRaises(SleepException, self.p.run) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 56dbc557..110543da 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -17,10 +17,6 @@ def subprocess_call(cmd, shell=False, env=None): return 0 -def false(*args, **kwargs): - return False - - class MockCursor: def __init__(self): @@ -28,7 +24,7 @@ class MockCursor: self.results = [] def execute(self, sql, *params): - if sql.startswith('blabla'): + if sql.startswith('blabla') or sql == 'CHECKPOINT': raise psycopg2.OperationalError() elif sql.startswith('InterfaceError'): raise psycopg2.InterfaceError() @@ -88,12 +84,11 @@ class MockConnect: def psycopg2_connect(*args, **kwargs): - return MockConnect() -def is_running(): - return False +def raise_exception(*args, **kwargs): + raise Exception class TestPostgresql(unittest.TestCase): @@ -143,7 +138,7 @@ class TestPostgresql(unittest.TestCase): def test_start_stop(self): self.assertFalse(self.p.start()) - self.p.is_running = is_running + self.p.is_running = false with open(os.path.join(self.p.data_dir, 'postmaster.pid'), 'w'): pass self.assertTrue(self.p.start()) @@ -159,6 +154,10 @@ class TestPostgresql(unittest.TestCase): self.p.follow_the_leader(self.leader) self.p.follow_the_leader(Leader(-1, None, 28, self.other)) + def test_create_replica(self): + self.p.delete_trigger_file = raise_exception + self.assertEquals(self.p.create_replica({'host': '', 'port': '', 'user': ''}, ''), 1) + def test_create_connection_users(self): cfg = self.p.config cfg['superuser']['username'] = 'test' @@ -202,7 +201,7 @@ class TestPostgresql(unittest.TestCase): def test_is_healthy(self): self.assertTrue(self.p.is_healthy()) - self.p.is_running = is_running + self.p.is_running = false self.assertFalse(self.p.is_healthy()) def test_promote(self): @@ -211,6 +210,12 @@ class TestPostgresql(unittest.TestCase): def test_last_operation(self): self.assertEquals(self.p.last_operation(), '0') + def test_call_nowait(self): + popen = subprocess.Popen + subprocess.Popen = raise_exception + self.assertFalse(self.p.call_nowait('on_start')) + subprocess.Popen = popen + def test_non_existing_callback(self): self.assertFalse(self.p.call_nowait('foobar')) @@ -220,7 +225,9 @@ class TestPostgresql(unittest.TestCase): self.assertTrue(self.p.stop()) def test_move_data_directory(self): - self.p.is_running = is_running + self.p.is_running = false os.rename = nop os.path.isdir = true self.p.move_data_directory() + os.rename = raise_exception + self.p.move_data_directory() From 0b753d25e1825c266fc2f931d3d5822d535926ed Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 17 Sep 2015 13:57:29 +0200 Subject: [PATCH 3/4] Get rid from is_promoted flag. use role == 'master' instead --- patroni/__init__.py | 9 ++----- patroni/ha.py | 10 ++++---- patroni/postgresql.py | 55 +++++++++++++++++++--------------------- tests/test_ha.py | 2 +- tests/test_postgresql.py | 10 +++----- 5 files changed, 37 insertions(+), 49 deletions(-) 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') From 6530e1f7aada3e280cffa04540e3180e186fe048 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 17 Sep 2015 16:11:08 +0200 Subject: [PATCH 4/4] Remove unused parameter in a is_leader method --- patroni/postgresql.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index c77471d6..398a8ab8 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -152,7 +152,7 @@ class Postgresql: return 1 return ret - def is_leader(self, check_only=False): + def is_leader(self): return not self.query('SELECT pg_is_in_recovery()').fetchone()[0] def is_running(self): @@ -176,7 +176,7 @@ class Postgresql: def start(self, block_callbacks=False): if self.is_running(): - self._role = 'master' if self.is_leader(check_only=True) else 'replica' + self._role = 'master' if self.is_leader() else 'replica' self.schedule_load_slots = self.use_slots logger.error('Cannot start PostgreSQL because one is already running.') return False