From 1fc8b43b366df1b234ea41029d98a4710a2ff691 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Wed, 24 Aug 2016 09:28:58 +0200 Subject: [PATCH 1/6] Return replication information on the api To enable better monitoring, it is useful to have replication statistics. Addresses issue #261 --- patroni/api.py | 10 ++++++++-- tests/test_postgresql.py | 3 ++- 2 files changed, 10 insertions(+), 3 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index 6f3a6ae9..be57c52b 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -368,7 +368,11 @@ class RestApiHandler(BaseHTTPRequestHandler): def get_postgresql_status(self, retry=False): try: - row = self.query("""SELECT to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'), + row = self.query("""WITH replication_info AS ( + SELECT 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'), pg_is_in_recovery(), CASE WHEN pg_is_in_recovery() THEN 0 @@ -377,12 +381,14 @@ class RestApiHandler(BaseHTTPRequestHandler): pg_xlog_location_diff(pg_last_xlog_receive_location(), '0/0')::bigint, pg_xlog_location_diff(pg_last_xlog_replay_location(), '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_xlog_replay_paused()""", retry=retry)[0] + pg_is_in_recovery() AND pg_is_xlog_replay_paused(), + (SELECT json_agg(row_to_json(ri)) FROM replication_info ri)""", retry=retry)[0] return { 'state': self.server.patroni.postgresql.state, 'postmaster_start_time': row[0], 'role': 'replica' if row[1] else 'master', 'server_version': self.server.patroni.postgresql.server_version, + 'replication': row[7], 'xlog': ({ 'received_location': row[3], 'replayed_location': row[4], diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 2ab29124..f6aaa869 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -33,7 +33,8 @@ class MockCursor(object): elif sql == 'SELECT pg_is_in_recovery()': self.results = [(False, )] elif sql.startswith('SELECT to_char(pg_postmaster_start_time'): - self.results = [('', True, '', '', '', '', False)] + replication_info = '[{"application_name":"walreceiver","client_addr":"1.2.3.4","state":"streaming","sync_state":"async","sync_priority":0}]' + self.results = [('', True, '', '', '', '', False, replication_info)] elif sql.startswith('SELECT name, setting'): self.results = [('wal_segment_size', '2048', '8kB', 'integer', 'internal'), ('search_path', 'public', None, 'string', 'user'), From a573983753e494e6a114808a5561897f1ddf46c1 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Wed, 24 Aug 2016 11:54:40 +0200 Subject: [PATCH 2/6] Include usename in replication information Also only return the key if any replication information is known --- patroni/api.py | 20 ++++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index be57c52b..1d3aa966 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -369,7 +369,7 @@ class RestApiHandler(BaseHTTPRequestHandler): def get_postgresql_status(self, retry=False): try: row = self.query("""WITH replication_info AS ( - SELECT application_name, client_addr, state, sync_state, sync_priority + 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'), @@ -383,20 +383,24 @@ class RestApiHandler(BaseHTTPRequestHandler): to_char(pg_last_xact_replay_timestamp(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'), pg_is_in_recovery() AND pg_is_xlog_replay_paused(), (SELECT json_agg(row_to_json(ri)) FROM replication_info ri)""", retry=retry)[0] - return { + result = { 'state': self.server.patroni.postgresql.state, 'postmaster_start_time': row[0], 'role': 'replica' if row[1] else 'master', 'server_version': self.server.patroni.postgresql.server_version, - 'replication': row[7], - 'xlog': ({ + 'replication': row[7]} + if result['role'] == 'replica': + result['xlog'] = { 'received_location': row[3], 'replayed_location': row[4], 'replayed_timestamp': row[5], - 'paused': row[6]} if row[1] else { - 'location': row[2] - }) - } + 'paused': row[6]} + else: + result['xlog'] = {'location': row[2]} + if not result['replication']: + del result['replication'] + + return result except (psycopg2.Error, RetryFailedError, PostgresConnectionException): state = self.server.patroni.postgresql.state if state == 'running': From a09f905a78f9ba4545a960a01ba3d915627dd192 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Wed, 24 Aug 2016 12:28:31 +0200 Subject: [PATCH 3/6] Only add replication info if it is found --- patroni/api.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index 1d3aa966..8293c051 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -387,8 +387,7 @@ class RestApiHandler(BaseHTTPRequestHandler): 'state': self.server.patroni.postgresql.state, 'postmaster_start_time': row[0], 'role': 'replica' if row[1] else 'master', - 'server_version': self.server.patroni.postgresql.server_version, - 'replication': row[7]} + 'server_version': self.server.patroni.postgresql.server_version} if result['role'] == 'replica': result['xlog'] = { 'received_location': row[3], @@ -397,8 +396,8 @@ class RestApiHandler(BaseHTTPRequestHandler): 'paused': row[6]} else: result['xlog'] = {'location': row[2]} - if not result['replication']: - del result['replication'] + if row[7]: + result['replication'] = row[7] return result except (psycopg2.Error, RetryFailedError, PostgresConnectionException): From 74166e996c7abe1ac5ffc786c3de1928e6fe728b Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 25 Aug 2016 10:09:32 +0200 Subject: [PATCH 4/6] Fix tests and formatting --- patroni/api.py | 14 ++++++++------ tests/test_postgresql.py | 5 +++-- 2 files changed, 11 insertions(+), 8 deletions(-) diff --git a/patroni/api.py b/patroni/api.py index 8293c051..dc737a4f 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -383,19 +383,21 @@ class RestApiHandler(BaseHTTPRequestHandler): to_char(pg_last_xact_replay_timestamp(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'), pg_is_in_recovery() AND pg_is_xlog_replay_paused(), (SELECT json_agg(row_to_json(ri)) FROM replication_info ri)""", retry=retry)[0] + result = { 'state': self.server.patroni.postgresql.state, 'postmaster_start_time': row[0], 'role': 'replica' if row[1] else 'master', - 'server_version': self.server.patroni.postgresql.server_version} - if result['role'] == 'replica': - result['xlog'] = { + 'server_version': self.server.patroni.postgresql.server_version, + 'xlog': ({ 'received_location': row[3], 'replayed_location': row[4], 'replayed_timestamp': row[5], - 'paused': row[6]} - else: - result['xlog'] = {'location': row[2]} + 'paused': row[6]} if row[1] else { + 'location': row[2] + }) + } + if row[7]: result['replication'] = row[7] diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index f6aaa869..0d931d54 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -32,8 +32,9 @@ class MockCursor(object): self.results = [(0,)] elif sql == 'SELECT pg_is_in_recovery()': self.results = [(False, )] - elif sql.startswith('SELECT to_char(pg_postmaster_start_time'): - replication_info = '[{"application_name":"walreceiver","client_addr":"1.2.3.4","state":"streaming","sync_state":"async","sync_priority":0}]' + elif sql.startswith('WITH replication_info AS ('): + replication_info = '[{"application_name":"walreceiver","client_addr":"1.2.3.4",' +\ + '"state":"streaming","sync_state":"async","sync_priority":0}]' self.results = [('', True, '', '', '', '', False, replication_info)] elif sql.startswith('SELECT name, setting'): self.results = [('wal_segment_size', '2048', '8kB', 'integer', 'internal'), From e428c8d0faaad8fcfda55f78820bb1b802881314 Mon Sep 17 00:00:00 2001 From: Ants Aasma Date: Tue, 30 Aug 2016 00:21:30 +0300 Subject: [PATCH 5/6] Replace invalid characters in member names for replication slot names PostgreSQL replication slot names only allow names consisting of [a-z0-9_]. Invalid characters cause replication slot creation and standby startup to fail. This change substitutes the invalid characters with underscores or unicode codepoints. In case multiple member names map to identical replication slots master log will contain a corresponding error message. Motivated by wanting to use hostnames as member names. Hostnames often contain periods and dashes. --- patroni/postgresql.py | 45 +++++++++++++++++++++++++++++++++------- tests/test_postgresql.py | 10 ++++++++- 2 files changed, 46 insertions(+), 9 deletions(-) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index da4b932a..dab5b88f 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -1,6 +1,8 @@ +from collections import defaultdict import logging import os import psycopg2 +import re import shlex import shutil import subprocess @@ -21,6 +23,22 @@ ACTION_ON_RELOAD = "on_reload" ACTION_ON_ROLE_CHANGE = "on_role_change" +def slot_name_from_member_name(member_name): + """Translate member name to valid PostgreSQL slot name. + + PostgreSQL replication slot names must be valid PostgreSQL names. This function maps the wider space of + member names to valid PostgreSQL names. Names are lowercased, dashes and periods common in hostnames + are replaced with underscores, other characters are encoded as their unicode codepoint. Name is truncated + to 64 characters. Multiple different member names may map to a single slot name.""" + + def replace_char(match): + c = match.group(0) + return '_' if c in '-.' else "u%04d" % ord(c) + + slot_name = re.sub('[^a-z0-9_]', replace_char, member_name.lower()) + return slot_name[0:64] + + class Postgresql(object): # List of parameters which must be always passed to postmaster as command line options @@ -651,7 +669,7 @@ class Postgresql(object): if primary_conninfo: f.write("primary_conninfo = '{0}'\n".format(primary_conninfo)) if self.use_slots: - f.write("primary_slot_name = '{0}'\n".format(self.name)) + f.write("primary_slot_name = '{0}'\n".format(slot_name_from_member_name(self.name))) for name, value in self.config.get('recovery_conf', {}).items(): if name not in ('standby_mode', 'recovery_target_timeline', 'primary_conninfo', 'primary_slot_name'): f.write("{0} = '{1}'\n".format(name, value)) @@ -889,21 +907,32 @@ $$""".format(name, ' '.join(options)), name, password, password) # the replicatefrom destination member is currently not a member of the cluster (fallback to the # master), or if replicatefrom destination member happens to be the current master if self.role == 'master': - slots = [m.name for m in cluster.members if m.name != self.name and - (m.replicatefrom is None or m.replicatefrom == self.name or - not cluster.has_member(m.replicatefrom))] + slot_members = [m.name for m in cluster.members if m.name != self.name and + (m.replicatefrom is None or m.replicatefrom == self.name or + not cluster.has_member(m.replicatefrom))] else: # only manage slots for replicas that replicate from this one, except for the leader among them - slots = [m.name for m in cluster.members if m.replicatefrom == self.name and - m.name != cluster.leader.name] + slot_members = [m.name for m in cluster.members if m.replicatefrom == self.name and + m.name != cluster.leader.name] + slots = set(slot_name_from_member_name(name) for name in slot_members) + + if len(slots) < len(slot_members): + # Find which names are conflicting for a nicer error message + slot_conflicts = defaultdict(list) + for name in slot_members: + slot_conflicts[slot_name_from_member_name(name)].append(name) + logger.error("Following cluster members share a replication slot name: %s", + "; ".join("%s map to %s" % (", ".join(v), k) + for k, v in slot_conflicts.items() if len(v) > 1)) + # drop unused slots - for slot in set(self._replication_slots) - set(slots): + for slot in set(self._replication_slots) - slots: self.query("""SELECT pg_drop_replication_slot(%s) WHERE EXISTS(SELECT 1 FROM pg_replication_slots WHERE slot_name = %s AND NOT active)""", slot, slot) # create new slots - for slot in set(slots) - set(self._replication_slots): + for slot in slots - set(self._replication_slots): self.query("""SELECT pg_create_physical_replication_slot(%s) WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots WHERE slot_name = %s)""", slot, slot) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 2ab29124..ab973514 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -182,7 +182,7 @@ class TestPostgresql(unittest.TestCase): 'restore': 'true'}) self.leadermem = Member(0, 'leader', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres'}) self.leader = Leader(-1, 28, self.leadermem) - self.other = Member(0, 'test1', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5433/postgres', + self.other = Member(0, 'test-1', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5433/postgres', 'tags': {'replicatefrom': 'leader'}}) self.me = Member(0, 'test0', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5434/postgres'}) @@ -320,6 +320,14 @@ class TestPostgresql(unittest.TestCase): self.p.schedule_load_slots = False with mock.patch('patroni.postgresql.Postgresql.role', new_callable=PropertyMock(return_value='replica')): self.p.sync_replication_slots(cluster) + with mock.patch('patroni.postgresql.logger.error', new_callable=Mock()) as errorlog_mock: + self.p.query = Mock() + alias1 = Member(0, 'test-3', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5436/postgres'}) + alias2 = Member(0, 'test.3', 28, {'conn_url': 'postgres://replicator:rep-pass@127.0.0.1:5436/postgres'}) + cluster.members.extend([alias1, alias2]) + self.p.sync_replication_slots(cluster) + errorlog_mock.assert_called_once() + assert "test-3" in errorlog_mock.call_args[0][1] and "test.3" in errorlog_mock.call_args[0][1] @patch.object(MockConnect, 'closed', 2) def test__query(self): From fa6bd51ad125d7512e1235f28087e2c0490ff836 Mon Sep 17 00:00:00 2001 From: Ants Aasma Date: Tue, 30 Aug 2016 00:40:19 +0300 Subject: [PATCH 6/6] Appease Quantifiedcode about stylistic issues --- patroni/postgresql.py | 4 ++-- tests/test_postgresql.py | 3 ++- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index dab5b88f..37d09460 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -33,7 +33,7 @@ def slot_name_from_member_name(member_name): def replace_char(match): c = match.group(0) - return '_' if c in '-.' else "u%04d" % ord(c) + return '_' if c in '-.' else "u{:04d}".format(ord(c)) slot_name = re.sub('[^a-z0-9_]', replace_char, member_name.lower()) return slot_name[0:64] @@ -922,7 +922,7 @@ $$""".format(name, ' '.join(options)), name, password, password) for name in slot_members: slot_conflicts[slot_name_from_member_name(name)].append(name) logger.error("Following cluster members share a replication slot name: %s", - "; ".join("%s map to %s" % (", ".join(v), k) + "; ".join("{} map to {}".format(", ".join(v), k) for k, v in slot_conflicts.items() if len(v) > 1)) # drop unused slots diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index ab973514..a9d729f6 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -327,7 +327,8 @@ class TestPostgresql(unittest.TestCase): cluster.members.extend([alias1, alias2]) self.p.sync_replication_slots(cluster) errorlog_mock.assert_called_once() - assert "test-3" in errorlog_mock.call_args[0][1] and "test.3" in errorlog_mock.call_args[0][1] + assert "test-3" in errorlog_mock.call_args[0][1] + assert "test.3" in errorlog_mock.call_args[0][1] @patch.object(MockConnect, 'closed', 2) def test__query(self):