From b85262637df86221752dfd6c2d3ec978d1125844 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 1 Jun 2015 09:20:24 +0200 Subject: [PATCH 1/5] Bugfix: do not create replication slot for master --- helpers/api.py | 2 +- helpers/ha.py | 2 +- helpers/postgresql.py | 3 ++- tests/test_postgresql.py | 13 ++++++++++--- 4 files changed, 14 insertions(+), 6 deletions(-) diff --git a/helpers/api.py b/helpers/api.py index c723e47b..06ff07da 100644 --- a/helpers/api.py +++ b/helpers/api.py @@ -20,7 +20,7 @@ class RestApiHandler(BaseHTTPRequestHandler): try: response = self.get_postgresql_status() except (psycopg2.OperationalError, psycopg2.InterfaceError): - logging.exception('get_postgresql_status') + logger.exception('get_postgresql_status') response = {'running': False} path = '/master' if self.path == '/' else self.path diff --git a/helpers/ha.py b/helpers/ha.py index 29979363..f8cf5576 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -78,7 +78,7 @@ class Ha: return 'no action. i am the leader with the lock' finally: # create replication slots - self.state_handler.create_replication_slots([m.hostname for m in self.cluster.members]) + 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 48df8345..bdbfabe0 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -247,7 +247,8 @@ 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, members): + def create_replication_slots(self, cluster): + members = [m.hostname for m in cluster.members if m.hostname != self.name] # drop unused slots for slot in set(self.members) - set(members): self.query("""SELECT pg_drop_replication_slot(%s) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index dc6fb512..c766046d 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -91,8 +91,12 @@ class TestPostgresql(unittest.TestCase): def set_up(self): subprocess.call = subprocess_call - self.p = Postgresql({'name': 'test0', 'data_dir': 'data/test0', 'listen': '127.0.0.1, 127.0.0.2:5432', 'connect_address': '127.0.0.2:5432', 'replication': { - 'username': 'replicator', 'password': 'rep-pass', 'network': '127.0.0.1/32'}, 'parameters': {'foo': 'bar'}, 'recovery_conf': {'foo': 'bar'}}) + self.p = Postgresql({'name': 'test0', 'data_dir': 'data/test0', 'listen': '127.0.0.1, 127.0.0.2:5432', + 'connect_address': '127.0.0.2:5432', + 'replication': {'username': 'replicator', + 'password': 'rep-pass', + 'network': '127.0.0.1/32'}, + 'parameters': {'foo': 'bar'}, 'recovery_conf': {'foo': 'bar'}}) psycopg2.connect = psycopg2_connect if not os.path.exists(self.p.data_dir): os.makedirs(self.p.data_dir) @@ -127,7 +131,10 @@ class TestPostgresql(unittest.TestCase): def test_create_replication_slots(self): self.p.start() - self.p.create_replication_slots('qaz') + me = Member('test0', 'postgres://replicator:rep-pass@127.0.0.1:5434/postgres', 28) + other = Member('test1', 'postgres://replicator:rep-pass@127.0.0.1:5433/postgres', 28) + cluster = Cluster(True, self.leader, 0, [me, other, self.leader]) + self.p.create_replication_slots(cluster) def test_query(self): self.p.query('select 1') From f8c3582715699b81c2bcf763bb63f57beef5aa81 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 1 Jun 2015 10:48:16 +0200 Subject: [PATCH 2/5] Refactor api: Call is_running only when it's not possible to execute query against governed postgres. --- helpers/api.py | 57 ++++++++++++++++++++++------------------------- tests/test_api.py | 9 +------- 2 files changed, 28 insertions(+), 38 deletions(-) diff --git a/helpers/api.py b/helpers/api.py index 06ff07da..53e3f2fb 100644 --- a/helpers/api.py +++ b/helpers/api.py @@ -10,21 +10,16 @@ if sys.hexversion >= 0x03000000: else: from BaseHTTPServer import BaseHTTPRequestHandler, HTTPServer - logger = logging.getLogger(__name__) class RestApiHandler(BaseHTTPRequestHandler): def do_GET(self): - try: - response = self.get_postgresql_status() - except (psycopg2.OperationalError, psycopg2.InterfaceError): - logger.exception('get_postgresql_status') - response = {'running': False} + response = self.get_postgresql_status() path = '/master' if self.path == '/' else self.path - status_code = 200 if response['running'] and response['role'] in path else 503 + status_code = 200 if response['running'] and 'role' in response and response['role'] in path else 503 self.send_response(status_code) self.send_header('Content-Type', 'application/json') @@ -32,29 +27,31 @@ class RestApiHandler(BaseHTTPRequestHandler): self.wfile.write(json.dumps(response).encode('utf-8')) def get_postgresql_status(self): - if not self.server.governor.postgresql.is_running(): - return {'running': False} - cursor = self.server._cursor() - cursor.execute("""SELECT to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'), - pg_is_in_recovery(), - CASE WHEN pg_is_in_recovery() - THEN null - ELSE pg_current_xlog_location() END, - pg_last_xlog_receive_location(), - pg_last_xlog_replay_location(), - pg_is_in_recovery() AND pg_is_xlog_replay_paused()""") - row = cursor.fetchone() - return { - 'running': True, - 'postmaster_start_time': row[0], - 'role': 'slave' if row[1] else 'master', - 'xlog': ({ - 'received_location': row[3], - 'replayed_location': row[4], - 'paused': row[5]} if row[1] else { - 'location': row[2] - }) - } + try: + cursor = self.server._cursor() + cursor.execute("""SELECT to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'), + pg_is_in_recovery(), + CASE WHEN pg_is_in_recovery() + THEN null + ELSE pg_current_xlog_location() END, + pg_last_xlog_receive_location(), + pg_last_xlog_replay_location(), + pg_is_in_recovery() AND pg_is_xlog_replay_paused()""") + row = cursor.fetchone() + return { + 'running': True, + 'postmaster_start_time': row[0], + 'role': 'slave' if row[1] else 'master', + 'xlog': ({ + 'received_location': row[3], + 'replayed_location': row[4], + 'paused': row[5]} if row[1] else { + 'location': row[2] + }) + } + except (psycopg2.OperationalError, psycopg2.InterfaceError): + logger.exception('get_postgresql_status') + return {'running': self.server.governor.postgresql.is_running()} class RestApiServer(HTTPServer, Thread): diff --git a/tests/test_api.py b/tests/test_api.py index 6b117b48..eff848a3 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -11,10 +11,6 @@ else: from StringIO import StringIO as IO -def false(*args, **kwargs): - return False - - def throws(*args, **kwargs): raise psycopg2.OperationalError() @@ -48,7 +44,7 @@ class MockRestApiServer(RestApiServer): def __init__(self, Handler, path, *args): self.governor = MockGovernor() if len(args) > 0: - self.governor.postgresql.is_running = args[0] + self._cursor = args[0] self._cursor_holder = None Handler(MockRequest(path), ('0.0.0.0', 8080), self) @@ -61,6 +57,3 @@ class TestRestApiHandler(unittest.TestCase): def test_do_GET(self): MockRestApiServer(RestApiHandler, b'GET /') MockRestApiServer(RestApiHandler, b'GET /', throws) - - def test_get_postgresql_status(self): - MockRestApiServer(RestApiHandler, b'GET /', false) From 939254021e750feee3e35025359b3feeb24c2cf8 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 1 Jun 2015 11:03:06 +0200 Subject: [PATCH 3/5] Code cleanup: do not modify os.environ but pass copy of it to subprocess.call --- helpers/postgresql.py | 17 +++++++---------- tests/test_postgresql.py | 2 +- 2 files changed, 8 insertions(+), 11 deletions(-) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index bdbfabe0..e8a15e4f 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -109,12 +109,10 @@ class Postgresql: os.fchmod(f.fileno(), 0o600) f.write('{host}:{port}:*:{user}:{password}\n'.format(**r)) - try: - os.environ['PGPASSFILE'] = pgpass - return subprocess.call(['pg_basebackup', '-R', '-D', self.data_dir, - '--host=' + r['host'], '--port=' + str(r['port']), '-U', r['user']]) == 0 - finally: - os.environ.pop('PGPASSFILE') + 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 def is_leader(self): return not self.query('SELECT pg_is_in_recovery()').fetchone()[0] @@ -223,10 +221,9 @@ primary_conninfo = '{}' f.write("{} = '{}'\n".format(name, value)) def follow_the_leader(self, leader): - if self.check_recovery_conf(leader): - return - self.write_recovery_conf(leader) - self.restart() + if not self.check_recovery_conf(leader): + self.write_recovery_conf(leader) + self.restart() def promote(self): return subprocess.call(self._pg_ctl + ['promote']) == 0 diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index c766046d..6ded6a23 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -8,7 +8,7 @@ from helpers.etcd import Cluster, Member from helpers.postgresql import Postgresql -def subprocess_call(cmd, shell=False): +def subprocess_call(cmd, shell=False, env=None): return 0 From 8c4e547d7ad4cff6f382a8131f384559e5454d67 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 1 Jun 2015 12:02:46 +0200 Subject: [PATCH 4/5] Handle KeyboardInterrupt in the main loop --- governor.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/governor.py b/governor.py index 55cb4931..5267d07c 100755 --- a/governor.py +++ b/governor.py @@ -87,6 +87,8 @@ def main(): try: governor.initialize() governor.run() + except KeyboardInterrupt: + pass finally: governor.touch_member(300) # schedule member removal governor.postgresql.stop() From adcc7ac2560578684807c0dee721ce278d1f8016 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 1 Jun 2015 13:48:44 +0200 Subject: [PATCH 5/5] 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. --- helpers/ha.py | 15 +++++++-------- helpers/postgresql.py | 22 ++++++++++++++++++---- tests/test_ha.py | 1 + tests/test_postgresql.py | 2 ++ 4 files changed, 28 insertions(+), 12 deletions(-) diff --git a/helpers/ha.py b/helpers/ha.py index f8cf5576..771e0bfd 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -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) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index e8a15e4f..7677c6c4 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -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) diff --git a/tests/test_ha.py b/tests/test_ha.py index 65d7bd6a..59495b95 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -19,6 +19,7 @@ class MockPostgresql: def __init__(self): self.name = 'postgresql0' + self.is_promoted = False def is_healthy(self): return True diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 6ded6a23..f0b04904 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -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())