From 0f666e69f347d9f75ae66f68d42312322a4e67f8 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 1 Nov 2018 16:17:40 +0100 Subject: [PATCH] Prefix system tables, views and functions with pg_catalog (#845) and implement missing unit tests --- patroni/api.py | 25 ++++++++------- patroni/ctl.py | 4 +-- patroni/postgresql.py | 57 ++++++++++++++++----------------- patroni/scripts/wale_restore.py | 13 ++++---- tests/test_ctl.py | 12 +++---- tests/test_postgresql.py | 14 +++++--- 6 files changed, 65 insertions(+), 60 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index e3791771..5bda448c 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -412,17 +412,20 @@ class RestApiHandler(BaseHTTPRequestHandler): raise RetryFailedError('') stmt = ("WITH replication_info AS (" "SELECT usename, application_name, client_addr, state, sync_state, sync_priority" - " FROM pg_stat_replication) SELECT" - " to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ')," - " CASE WHEN pg_is_in_recovery() THEN 0" - " ELSE ('x' || SUBSTR(pg_{0}file_name(pg_current_{0}_{1}()), 1, 8))::bit(32)::int END," - " CASE WHEN pg_is_in_recovery() THEN 0" - " ELSE pg_{0}_{1}_diff(pg_current_{0}_{1}(), '0/0')::bigint END," - " pg_{0}_{1}_diff(COALESCE(pg_last_{0}_receive_{1}(), pg_last_{0}_replay_{1}()), '0/0')::bigint," - " pg_{0}_{1}_diff(pg_last_{0}_replay_{1}(), '0/0')::bigint," - " to_char(pg_last_xact_replay_timestamp(), 'YYYY-MM-DD HH24:MI:SS.MS TZ')," - " pg_is_in_recovery() AND pg_is_{0}_replay_paused()," - " (SELECT array_to_json(array_agg(row_to_json(ri))) FROM replication_info ri)") + " FROM pg_catalog.pg_stat_replication) SELECT" + " pg_catalog.to_char(pg_catalog.pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ')," + " CASE WHEN pg_catalog.pg_is_in_recovery() THEN 0" + " ELSE ('x' || pg_catalog.substr(pg_catalog.pg_{0}file_name(" + "pg_catalog.pg_current_{0}_{1}()), 1, 8))::bit(32)::int END," + " CASE WHEN pg_catalog.pg_is_in_recovery() THEN 0" + " ELSE pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_current_{0}_{1}(), '0/0')::bigint END," + " pg_catalog.pg_{0}_{1}_diff(COALESCE(pg_catalog.pg_last_{0}_receive_{1}()," + " pg_catalog.pg_last_{0}_replay_{1}()), '0/0')::bigint," + " pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_last_{0}_replay_{1}(), '0/0')::bigint," + " pg_catalog.to_char(pg_catalog.pg_last_xact_replay_timestamp(), 'YYYY-MM-DD HH24:MI:SS.MS TZ')," + " pg_catalog.pg_is_in_recovery() AND pg_catalog.pg_is_{0}_replay_paused(), " + "(SELECT pg_catalog.array_to_json(pg_catalog.array_agg(" + "pg_catalog.row_to_json(ri))) FROM replication_info ri)") row = self.query(stmt.format(self.server.patroni.postgresql.wal_name, self.server.patroni.postgresql.lsn_name), retry=retry)[0] diff --git a/patroni/ctl.py b/patroni/ctl.py index 58d75d1e..32ee99b9 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -252,7 +252,7 @@ def get_cursor(cluster, connect_parameters, role='master', member=None): if role == 'any': return cursor - cursor.execute('SELECT pg_is_in_recovery()') + cursor.execute('SELECT pg_catalog.pg_is_in_recovery()') in_recovery = cursor.fetchone()[0] if in_recovery and role == 'replica' or not in_recovery and role == 'master': @@ -392,7 +392,7 @@ def query_member(cluster, cursor, member, role, command, connect_parameters): logging.debug(message) return [[timestamp(0), message]], None - cursor.execute('SELECT pg_is_in_recovery()') + cursor.execute('SELECT pg_catalog.pg_is_in_recovery()') in_recovery = cursor.fetchone()[0] if in_recovery and role == 'master' or not in_recovery and role == 'replica': diff --git a/patroni/postgresql.py b/patroni/postgresql.py index f3dae055..232a3891 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -39,12 +39,14 @@ STOP_POLLING_INTERVAL = 1 REWIND_STATUS = type('Enum', (), {'INITIAL': 0, 'CHECK': 1, 'NEED': 2, 'NOT_NEED': 3, 'SUCCESS': 4, 'FAILED': 5}) sync_standby_name_re = re.compile(r'^[A-Za-z_][A-Za-z_0-9\$]*$') -cluster_info_query = ("SELECT CASE WHEN pg_is_in_recovery() THEN 0 " - "ELSE ('x' || SUBSTR(pg_{0}file_name(pg_current_{0}_{1}()), 1, 8))::bit(32)::int END, " - "CASE WHEN pg_is_in_recovery() THEN GREATEST(" - " pg_{0}_{1}_diff(COALESCE(pg_last_{0}_receive_{1}(), '0/0'), '0/0')::bigint," - " pg_{0}_{1}_diff(pg_last_{0}_replay_{1}(), '0/0')::bigint)" - "ELSE pg_{0}_{1}_diff(pg_current_{0}_{1}(), '0/0')::bigint END") +cluster_info_query = ("SELECT CASE WHEN pg_catalog.pg_is_in_recovery() THEN 0 " + "ELSE ('x' || pg_catalog.substr(pg_catalog.pg_{0}file_name(" + "pg_catalog.pg_current_{0}_{1}()), 1, 8))::bit(32)::int END, " + "CASE WHEN pg_catalog.pg_is_in_recovery() THEN GREATEST(" + " pg_catalog.pg_{0}_{1}_diff(COALESCE(" + "pg_catalog.pg_last_{0}_receive_{1}(), '0/0'), '0/0')::bigint," + " pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_last_{0}_replay_{1}(), '0/0')::bigint)" + "ELSE pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_current_{0}_{1}(), '0/0')::bigint END") def quote_ident(value): @@ -308,8 +310,8 @@ class Postgresql(object): changes['wal_segment_size'] = '16384kB' # XXX: query can raise an exception for r in self.query("""SELECT name, setting, unit, vartype, context - FROM pg_settings - WHERE LOWER(name) IN (""" + ', '.join(['%s'] * len(changes)) + """) + FROM pg_catalog.pg_settings + WHERE pg_catalog.lower(name) IN (""" + ', '.join(['%s'] * len(changes)) + """) ORDER BY 1 DESC""", *(k.lower() for k in changes.keys())): if r[4] == 'internal': if r[0] == 'wal_segment_size': @@ -712,12 +714,8 @@ class Postgresql(object): "datadir": self._data_dir, "connstring": connstring}) else: - if 'no_params' in method_config: - del method_config['no_params'] - if 'no_master' in method_config: - del method_config['no_master'] - if 'keep_data' in method_config: - del method_config['keep_data'] + for param in ('no_params', 'no_master', 'keep_data'): + method_config.pop(param, None) params = ["--{0}={1}".format(arg, val) for arg, val in method_config.items()] try: # call script with the full set of parameters @@ -945,7 +943,7 @@ class Postgresql(object): with self._get_connection_cursor(**connect_kwargs) as cur: cur.execute("SET statement_timeout = 0") if check_not_is_in_recovery: - cur.execute('SELECT pg_is_in_recovery()') + cur.execute('SELECT pg_catalog.pg_is_in_recovery()') if cur.fetchone()[0]: return 'is_in_recovery=true' return cur.execute('CHECKPOINT') @@ -1253,7 +1251,7 @@ class Postgresql(object): def check_leader_is_not_in_recovery(self, **kwargs): try: with self._get_connection_cursor(connect_timeout=3, options='-c statement_timeout=2000', **kwargs) as cur: - cur.execute('SELECT pg_is_in_recovery()') + cur.execute('SELECT pg_catalog.pg_is_in_recovery()') if not cur.fetchone()[0]: return True logger.info('Leader is still in_recovery and therefore can\'t be used for rewind') @@ -1364,10 +1362,10 @@ class Postgresql(object): history_path = 'pg_{0}/{1:08X}.history'.format(self.wal_name, timeline) try: cursor = self._cursor() - cursor.execute('SELECT isdir, modification FROM pg_stat_file(%s)', (history_path,)) + cursor.execute('SELECT isdir, modification FROM pg_catalog.pg_stat_file(%s)', (history_path,)) isdir, modification = cursor.fetchone() if not isdir: - cursor.execute('SELECT pg_read_file(%s)', (history_path,)) + cursor.execute('SELECT pg_catalog.pg_read_file(%s)', (history_path,)) history = list(self.parse_history(cursor.fetchone()[0])) if history[-1][0] == timeline - 1: history[-1].append(modification.isoformat()) @@ -1539,7 +1537,7 @@ $$""".format(name, ' '.join(options)), name, password, password) def load_replication_slots(self): if self.use_slots and self._schedule_load_slots: replication_slots = {} - cursor = self._query('SELECT slot_name, slot_type, plugin, database FROM pg_replication_slots') + cursor = self._query('SELECT slot_name, slot_type, plugin, database FROM pg_catalog.pg_replication_slots') for r in cursor: value = {'type': r[1]} if r[1] == 'logical': @@ -1550,14 +1548,15 @@ $$""".format(name, ' '.join(options)), name, password, password) def postmaster_start_time(self): try: - cursor = self.query("""SELECT to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ')""") + cursor = self.query("SELECT pg_catalog.to_char(pg_catalog.pg_postmaster_start_time()," + " 'YYYY-MM-DD HH24:MI:SS.MS TZ')") return cursor.fetchone()[0] except psycopg2.Error: return None def drop_replication_slot(self, name): - cursor = self._query(('SELECT pg_drop_replication_slot(%s) WHERE EXISTS (SELECT 1 ' + - 'FROM pg_replication_slots WHERE slot_name = %s AND NOT active)'), name, name) + cursor = self._query(('SELECT pg_catalog.pg_drop_replication_slot(%s) WHERE EXISTS (SELECT 1 ' + + 'FROM pg_catalog.pg_replication_slots WHERE slot_name = %s AND NOT active)'), name, name) # In normal situation rowcount should be 1, otherwise either slot doesn't exists or it is still active return cursor.rowcount == 1 @@ -1594,8 +1593,8 @@ $$""".format(name, ' '.join(options)), name, password, password) if name not in self._replication_slots: if value['type'] == 'physical': try: - self._query(("SELECT pg_create_physical_replication_slot(%s{0})" + - " WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots" + + self._query(("SELECT pg_catalog.pg_create_physical_replication_slot(%s{0})" + + " WHERE NOT EXISTS (SELECT 1 FROM pg_catalog.pg_replication_slots" + " WHERE slot_type = 'physical' AND slot_name = %s)").format( immediately_reserve), name, name) except Exception: @@ -1611,8 +1610,8 @@ $$""".format(name, ' '.join(options)), name, password, password) with self._get_connection_cursor(**conn_kwargs) as cur: for name, value in values.items(): try: - cur.execute("SELECT pg_create_logical_replication_slot(%s, %s)" + - " WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots" + + cur.execute("SELECT pg_catalog.pg_create_logical_replication_slot(%s, %s)" + + " WHERE NOT EXISTS (SELECT 1 FROM pg_catalog.pg_replication_slots" + " WHERE slot_type = 'logical' AND slot_name = %s)", (name, value['plugin'], name)) except Exception: @@ -1775,9 +1774,9 @@ $$""".format(name, ' '.join(options)), name, password, password) # Pick candidates based on who has flushed WAL farthest. # TODO: for synchronous_commit = remote_write we actually want to order on write_location for app_name, state, sync_state in self.query( - """SELECT LOWER(application_name), state, sync_state - FROM pg_stat_replication - ORDER BY flush_{0} DESC""".format(self.lsn_name)): + "SELECT pg_catalog.lower(application_name), state, sync_state" + " FROM pg_catalog.pg_stat_replication" + " ORDER BY flush_{0} DESC".format(self.lsn_name)): member = members.get(app_name) if state != 'streaming' or not member or member.tags.get('nosync', False): continue diff --git a/patroni/scripts/wale_restore.py b/patroni/scripts/wale_restore.py index ad7335fb..b4c70991 100755 --- a/patroni/scripts/wale_restore.py +++ b/patroni/scripts/wale_restore.py @@ -224,13 +224,12 @@ class WALERestore(object): lsn_name = 'location' con.autocommit = True with con.cursor() as cur: - cur.execute("""SELECT CASE WHEN pg_is_in_recovery() - THEN GREATEST( - pg_{0}_{1}_diff(COALESCE( - pg_last_{0}_receive_{1}(), '0/0'), %s)::bigint, - pg_{0}_{1}_diff(pg_last_{0}_replay_{1}(), %s)::bigint) - ELSE pg_{0}_{1}_diff(pg_current_{0}_{1}(), %s)::bigint - END""".format(wal_name, lsn_name), + cur.execute(("SELECT CASE WHEN pg_catalog.pg_is_in_recovery()" + " THEN GREATEST(pg_catalog.pg_{0}_{1}_diff(COALESCE(" + "pg_last_{0}_receive_{1}(), '0/0'), %s)::bigint, " + "pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_last_{0}_replay_{1}(), %s)::bigint)" + " ELSE pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_current_{0}_{1}(), %s)::bigint" + " END").format(wal_name, lsn_name), (backup_start_lsn, backup_start_lsn, backup_start_lsn)) diff_in_bytes = int(cur.fetchone()[0]) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index c0af302b..c24b4ecc 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -192,24 +192,24 @@ class TestCtl(unittest.TestCase): def test_query_member(self): with patch('patroni.ctl.get_cursor', Mock(return_value=MockConnect().cursor())): - rows = query_member(None, None, None, 'master', 'SELECT pg_is_in_recovery()', {}) + rows = query_member(None, None, None, 'master', 'SELECT pg_catalog.pg_is_in_recovery()', {}) self.assertTrue('False' in str(rows)) - rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()', {}) + rows = query_member(None, None, None, 'replica', 'SELECT pg_catalog.pg_is_in_recovery()', {}) self.assertEqual(rows, (None, None)) with patch('test_postgresql.MockCursor.execute', Mock(side_effect=OperationalError('bla'))): - rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()', {}) + rows = query_member(None, None, None, 'replica', 'SELECT pg_catalog.pg_is_in_recovery()', {}) with patch('patroni.ctl.get_cursor', Mock(return_value=None)): - rows = query_member(None, None, None, None, 'SELECT pg_is_in_recovery()', {}) + rows = query_member(None, None, None, None, 'SELECT pg_catalog.pg_is_in_recovery()', {}) self.assertTrue('No connection to' in str(rows)) - rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()', {}) + rows = query_member(None, None, None, 'replica', 'SELECT pg_catalog.pg_is_in_recovery()', {}) self.assertTrue('No connection to' in str(rows)) with patch('patroni.ctl.get_cursor', Mock(side_effect=OperationalError('bla'))): - rows = query_member(None, None, None, 'replica', 'SELECT pg_is_in_recovery()', {}) + rows = query_member(None, None, None, 'replica', 'SELECT pg_catalog.pg_is_in_recovery()', {}) @patch('patroni.ctl.get_dcs') def test_dsn(self, mock_get_dcs): diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index eac9e5b6..873d3519 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -28,15 +28,15 @@ class MockCursor(object): def execute(self, sql, *params): if sql.startswith('blabla'): raise psycopg2.ProgrammingError() - elif sql == 'CHECKPOINT' or sql.startswith('SELECT pg_create_'): + elif sql == 'CHECKPOINT' or sql.startswith('SELECT pg_catalog.pg_create_'): raise psycopg2.OperationalError() elif sql.startswith('RetryFailedError'): raise RetryFailedError('retry') elif sql.startswith('SELECT slot_name'): self.results = [('blabla', 'physical'), ('foobar', 'physical'), ('ls', 'logical', 'a', 'b')] - elif sql.startswith('SELECT CASE WHEN pg_is_in_recovery()'): + elif sql.startswith('SELECT CASE WHEN pg_catalog.pg_is_in_recovery()'): self.results = [(1, 2)] - elif sql.startswith('SELECT pg_is_in_recovery()'): + elif sql.startswith('SELECT pg_catalog.pg_is_in_recovery()'): self.results = [(False, 2)] elif sql.startswith('WITH replication_info AS ('): replication_info = '[{"application_name":"walreceiver","client_addr":"1.2.3.4",' +\ @@ -53,7 +53,7 @@ class MockCursor(object): self.results = [('1', 2, '0/402EEC0', '')] elif sql.startswith('SELECT isdir, modification'): self.results = [(False, datetime.datetime.now())] - elif sql.startswith('SELECT pg_read_file'): + elif sql.startswith('SELECT pg_catalog.pg_read_file'): self.results = [('1\t0/40159C0\tno recovery target specified\n\n' + '2\t1/40159C0\tno recovery target specified\n',)] elif sql.startswith('TIMELINE_HISTORY '): @@ -421,9 +421,13 @@ class TestPostgresql(unittest.TestCase): def test_create_replica(self, mock_cancellable_subprocess_call): self.p.delete_trigger_file = Mock(side_effect=OSError) + self.p.config['create_replica_methods'] = ['pgBackRest'] + self.p.config['pgBackRest'] = {'command': 'pgBackRest', 'keep_data': True, 'no_params': True} + mock_cancellable_subprocess_call.return_value = 0 + self.assertEqual(self.p.create_replica(self.leader), 0) + self.p.config['create_replica_methods'] = ['wale', 'basebackup'] self.p.config['wale'] = {'command': 'foo'} - mock_cancellable_subprocess_call.return_value = 0 self.assertEqual(self.p.create_replica(self.leader), 0) del self.p.config['wale'] self.assertEqual(self.p.create_replica(self.leader), 0)