From e428c8d0faaad8fcfda55f78820bb1b802881314 Mon Sep 17 00:00:00 2001 From: Ants Aasma Date: Tue, 30 Aug 2016 00:21:30 +0300 Subject: [PATCH 1/2] 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 2/2] 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):