From 1bcc2b5fa6875157da1385458818d45348fc3957 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 1 Jun 2015 16:05:35 +0200 Subject: [PATCH 1/5] Extend postgres?.yml with pg_hba section to give possibility to customize pg_hba.conf --- helpers/postgresql.py | 4 ++++ postgres0.yml | 11 +++++++++++ postgres1.yml | 11 +++++++++++ tests/test_postgresql.py | 3 +++ 4 files changed, 29 insertions(+) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 7677c6c4..58038788 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -199,6 +199,10 @@ class Postgresql: def write_pg_hba(self): with open(os.path.join(self.data_dir, 'pg_hba.conf'), 'a') as f: f.write('\nhost replication {username} {network} md5\n'.format(**self.replication)) + for line in self.config.get('pg_hba', []): + if line['type'] == 'hostssl' and self.config['parameters'].get('ssl', 'off').lower() != 'on': + continue + f.write('{type} {database} {user} {address} {method}\n'. format(**line)) @staticmethod def primary_conninfo(leader_url): diff --git a/postgres0.yml b/postgres0.yml index 8757538f..7c1c8dbf 100644 --- a/postgres0.yml +++ b/postgres0.yml @@ -12,6 +12,17 @@ postgresql: connect_address: 127.0.0.1:5432 data_dir: data/postgresql0 maximum_lag_on_failover: 1048576 # 1 megabyte in bytes + pg_hba: + - type: host + database: all + user: all + address: 0.0.0.0/0 + method: md5 + - type: hostssl + database: all + user: all + address: 0.0.0.0/0 + method: md5 replication: username: replicator password: rep-pass diff --git a/postgres1.yml b/postgres1.yml index a0afe20f..2e6edf33 100644 --- a/postgres1.yml +++ b/postgres1.yml @@ -12,6 +12,17 @@ postgresql: connect_address: 127.0.0.1:5433 data_dir: data/postgresql1 maximum_lag_on_failover: 1048576 # 1 megabyte in bytes + pg_hba: + - type: host + database: all + user: all + address: 0.0.0.0/0 + method: md5 + - type: hostssl + database: all + user: all + address: 0.0.0.0/0 + method: md5 replication: username: replicator password: rep-pass diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index f0b04904..4fcec552 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -93,6 +93,9 @@ class TestPostgresql(unittest.TestCase): 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', + 'pg_hba': [{'type': 'hostssl', 'database': 'all', 'user': 'all', 'address': '0.0.0.0/0', + 'method': 'md5'}, {'type': 'host', 'database': 'all', 'user': 'all', + 'address': '0.0.0.0/0', 'method': 'md5'}], 'replication': {'username': 'replicator', 'password': 'rep-pass', 'network': '127.0.0.1/32'}, From efccd777b2860b147d708999cc23261db1b28c72 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 1 Jun 2015 16:47:56 +0200 Subject: [PATCH 2/5] rename pid_path to postmaster_pid --- helpers/postgresql.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 58038788..06e8b1a1 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -40,7 +40,7 @@ class Postgresql: self.data_dir = config['data_dir'] 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.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.is_promoted = False @@ -139,9 +139,9 @@ class Postgresql: logger.error('Cannot start PostgreSQL because one is already running.') return False - if os.path.exists(self.pid_path): - os.remove(self.pid_path) - logger.info('Removed %s', self.pid_path) + 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() From bacd05d99e371fabed8ce201b1d3d3d5f9347f28 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 1 Jun 2015 17:00:15 +0200 Subject: [PATCH 3/5] Determine preferable local address to connect through. If listen contains '*' or 0.0.0.0 - connect via localhost In all other cases pick the first one. --- helpers/postgresql.py | 10 ++++++++-- tests/test_postgresql.py | 2 +- 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 06e8b1a1..7f027b2b 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -57,8 +57,14 @@ class Postgresql: self.members = [] # list of already existing replication slots def get_local_address(self): - # TODO: try to get unix_socket_directory from postmaster.pid - return self.listen_addresses.split(',')[0].strip() + ':' + self.port + listen_addresses = self.listen_addresses.split(',') + local_address = listen_addresses[0].strip() # take first address from listen_addresses + + for la in listen_addresses: + if la.strip() in ['*', '0.0.0.0']: # we are listening on * + local_address = 'localhost' # connection via localhost is preferred + break + return local_address + ':' + self.port def connection(self): if not self._connection or self._connection.closed != 0: diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 4fcec552..4a4b7416 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -91,7 +91,7 @@ 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', + self.p = Postgresql({'name': 'test0', 'data_dir': 'data/test0', 'listen': '127.0.0.1, *:5432', 'connect_address': '127.0.0.2:5432', 'pg_hba': [{'type': 'hostssl', 'database': 'all', 'user': 'all', 'address': '0.0.0.0/0', 'method': 'md5'}, {'type': 'host', 'database': 'all', 'user': 'all', From f9ce29d49f944d0592f1669103995637964a47cf Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 2 Jun 2015 08:52:34 +0200 Subject: [PATCH 4/5] rename address to conn_url in a Member obj and introduce new field: api_url --- helpers/etcd.py | 6 +++--- helpers/postgresql.py | 10 +++++----- tests/test_postgresql.py | 14 +++++++------- 3 files changed, 15 insertions(+), 15 deletions(-) diff --git a/helpers/etcd.py b/helpers/etcd.py index 48c02d85..8aeb8ea7 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -8,7 +8,7 @@ from helpers.utils import sleep logger = logging.getLogger(__name__) -Member = namedtuple('Member', 'hostname,address,ttl') +Member = namedtuple('Member', 'hostname,conn_url,api_url,ttl') class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members')): @@ -91,7 +91,7 @@ class Etcd: initialize = True if node else False # get list of members node = self.find_node(response['node'], '/members') or {'nodes': []} - members = [Member(n['key'].split('/')[-1], n['value'], n.get('ttl', None)) for n in node['nodes']] + members = [Member(n['key'].split('/')[-1], n['value'], None, n.get('ttl', None)) for n in node['nodes']] # get last leader operation last_leader_operation = 0 @@ -110,7 +110,7 @@ class Etcd: leader = m break if not leader: - leader = Member(node['value'], None, None) + leader = Member(node['value'], None, None, None) return Cluster(initialize, leader, last_leader_operation, members) elif status_code == 404: diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 7f027b2b..ad987359 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -115,7 +115,7 @@ class Postgresql: os.path.exists(self.trigger_file) and os.unlink(self.trigger_file) def sync_from_leader(self, leader): - r = parseurl(leader.address) + r = parseurl(leader.conn_url) pgpass = 'pgpass' with open(pgpass, 'w') as f: @@ -185,7 +185,7 @@ class Postgresql: if member.hostname == self.name: continue try: - r = parseurl(member.address) + r = parseurl(member.conn_url) member_conn = psycopg2.connect(**r) member_conn.autocommit = True member_cursor = member_conn.cursor() @@ -219,7 +219,7 @@ class Postgresql: if not os.path.isfile(self.recovery_conf): return False - pattern = leader and leader.address and self.primary_conninfo(leader.address) + pattern = leader and leader.conn_url and self.primary_conninfo(leader.conn_url) with open(self.recovery_conf, 'r') as f: for line in f: @@ -235,11 +235,11 @@ class Postgresql: f.write("""standby_mode = 'on' recovery_target_timeline = 'latest' """) - if leader and leader.address: + if leader and leader.conn_url: f.write(""" primary_slot_name = '{}' primary_conninfo = '{}' -""".format(self.name, self.primary_conninfo(leader.address))) +""".format(self.name, self.primary_conninfo(leader.conn_url))) for name, value in self.config.get('recovery_conf', {}).items(): f.write("{} = '{}'\n".format(name, value)) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 4a4b7416..b7da49aa 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -103,7 +103,7 @@ class TestPostgresql(unittest.TestCase): psycopg2.connect = psycopg2_connect if not os.path.exists(self.p.data_dir): os.makedirs(self.p.data_dir) - self.leader = Member('leader', 'postgres://replicator:rep-pass@127.0.0.1:5434/postgres', 28) + self.leader = Member('leader', 'postgres://replicator:rep-pass@127.0.0.1:5434/postgres', None, 28) def tear_down(self): shutil.rmtree('data') @@ -130,12 +130,12 @@ class TestPostgresql(unittest.TestCase): self.p.follow_the_leader(None) self.p.demote(self.leader) self.p.follow_the_leader(self.leader) - self.p.follow_the_leader(Member('leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', 28)) + self.p.follow_the_leader(Member('leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', None, 28)) def test_create_replication_slots(self): self.p.start() - 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) + me = Member('test0', 'postgres://replicator:rep-pass@127.0.0.1:5434/postgres', None, 28) + other = Member('test1', 'postgres://replicator:rep-pass@127.0.0.1:5433/postgres', None, 28) cluster = Cluster(True, self.leader, 0, [me, other, self.leader]) self.p.create_replication_slots(cluster) @@ -150,9 +150,9 @@ class TestPostgresql(unittest.TestCase): self.assertRaises(psycopg2.OperationalError, self.p.query, 'blabla') def test_is_healthiest_node(self): - leader = Member('leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', 28) - 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) + leader = Member('leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', None, 28) + me = Member('test0', 'postgres://replicator:rep-pass@127.0.0.1:5434/postgres', None, 28) + other = Member('test1', 'postgres://replicator:rep-pass@127.0.0.1:5433/postgres', None, 28) cluster = Cluster(True, leader, 0, [me, other, leader]) self.assertTrue(self.p.is_healthiest_node(cluster)) self.p.is_leader = false From fca987deda961643495fecc5f0f67af6687feac6 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 2 Jun 2015 09:45:37 +0200 Subject: [PATCH 5/5] set value of api_url from application_name parameter of conn_url. Reset parameters of conn_url --- helpers/etcd.py | 22 ++++++++++++++++++++-- tests/test_etcd.py | 4 ++-- 2 files changed, 22 insertions(+), 4 deletions(-) diff --git a/helpers/etcd.py b/helpers/etcd.py index 8aeb8ea7..32345106 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -1,14 +1,32 @@ import logging import requests +import sys from requests.exceptions import RequestException from collections import namedtuple from helpers.errors import CurrentLeaderError, EtcdError from helpers.utils import sleep +if sys.hexversion >= 0x03000000: + from urllib.parse import urlparse, urlunparse, parse_qsl +else: + from urlparse import urlparse, urlunparse, parse_qsl + logger = logging.getLogger(__name__) -Member = namedtuple('Member', 'hostname,conn_url,api_url,ttl') + +class Member(namedtuple('Member', 'hostname,conn_url,api_url,ttl')): + + @staticmethod + def fromNode(node): + scheme, netloc, path, params, query, fragment = urlparse(node['value']) + conn_url = urlunparse((scheme, netloc, path, params, '', fragment)) + api_url = None + for name, value in parse_qsl(query): + if name == 'application_name' and value: + api_url = value + break + return Member(node['key'].split('/')[-1], conn_url, api_url, node.get('ttl', None)) class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members')): @@ -91,7 +109,7 @@ class Etcd: initialize = True if node else False # get list of members node = self.find_node(response['node'], '/members') or {'nodes': []} - members = [Member(n['key'].split('/')[-1], n['value'], None, n.get('ttl', None)) for n in node['nodes']] + members = [Member.fromNode(n) for n in node['nodes']] # get last leader operation last_leader_operation = 0 diff --git a/tests/test_etcd.py b/tests/test_etcd.py index 8f2280e6..a7ab93ee 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -29,11 +29,11 @@ def requests_get(url, **kwargs): raise requests.exceptions.RequestException() response = MockResponse() if url.startswith('http://remote') or url.startswith('http://127.0.0.1'): - response.content = '{"action":"get","node":{"key":"/service/batman5","dir":true,"nodes":[{"key":"/service/batman5/initialize","value":"postgresql0","modifiedIndex":1582,"createdIndex":1582},{"key":"/service/batman5/leader","value":"postgresql1","expiration":"2015-05-15T09:11:00.037397538Z","ttl":21,"modifiedIndex":20728,"createdIndex":20434},{"key":"/service/batman5/optime","dir":true,"nodes":[{"key":"/service/batman5/optime/leader","value":"2164261704","modifiedIndex":20729,"createdIndex":20729}],"modifiedIndex":20437,"createdIndex":20437},{"key":"/service/batman5/members","dir":true,"nodes":[{"key":"/service/batman5/members/postgresql1","value":"postgres://replicator:rep-pass@127.0.0.1:5434/postgres","expiration":"2015-05-15T09:10:59.949384522Z","ttl":21,"modifiedIndex":20727,"createdIndex":20727},{"key":"/service/batman5/members/postgresql0","value":"postgres://replicator:rep-pass@127.0.0.1:5433/postgres","expiration":"2015-05-15T09:11:09.611860899Z","ttl":30,"modifiedIndex":20730,"createdIndex":20730}],"modifiedIndex":1581,"createdIndex":1581}],"modifiedIndex":1581,"createdIndex":1581}}' + response.content = '{"action":"get","node":{"key":"/service/batman5","dir":true,"nodes":[{"key":"/service/batman5/initialize","value":"postgresql0","modifiedIndex":1582,"createdIndex":1582},{"key":"/service/batman5/leader","value":"postgresql1","expiration":"2015-05-15T09:11:00.037397538Z","ttl":21,"modifiedIndex":20728,"createdIndex":20434},{"key":"/service/batman5/optime","dir":true,"nodes":[{"key":"/service/batman5/optime/leader","value":"2164261704","modifiedIndex":20729,"createdIndex":20729}],"modifiedIndex":20437,"createdIndex":20437},{"key":"/service/batman5/members","dir":true,"nodes":[{"key":"/service/batman5/members/postgresql1","value":"postgres://replicator:rep-pass@127.0.0.1:5434/postgres?application_name=http://127.0.0.1:8009/governor","expiration":"2015-05-15T09:10:59.949384522Z","ttl":21,"modifiedIndex":20727,"createdIndex":20727},{"key":"/service/batman5/members/postgresql0","value":"postgres://replicator:rep-pass@127.0.0.1:5433/postgres?application_name=http://127.0.0.1:8008/governor","expiration":"2015-05-15T09:11:09.611860899Z","ttl":30,"modifiedIndex":20730,"createdIndex":20730}],"modifiedIndex":1581,"createdIndex":1581}],"modifiedIndex":1581,"createdIndex":1581}}' elif url.startswith('http://other'): response.status_code = 404 elif url.startswith('http://noleader'): - response.content = '{"action":"get","node":{"key":"/service/batman5","dir":true,"nodes":[{"key":"/service/batman5/initialize","value":"postgresql0","modifiedIndex":1582,"createdIndex":1582},{"key":"/service/batman5/leader","value":"postgresql1","expiration":"2015-05-15T09:11:00.037397538Z","ttl":21,"modifiedIndex":20728,"createdIndex":20434},{"key":"/service/batman5/optime","dir":true,"nodes":[{"key":"/service/batman5/optime/leader","value":"2164261704","modifiedIndex":20729,"createdIndex":20729}],"modifiedIndex":20437,"createdIndex":20437},{"key":"/service/batman5/members","dir":true,"nodes":[{"key":"/service/batman5/members/postgresql0","value":"postgres://replicator:rep-pass@127.0.0.1:5433/postgres","expiration":"2015-05-15T09:11:09.611860899Z","ttl":30,"modifiedIndex":20730,"createdIndex":20730}],"modifiedIndex":1581,"createdIndex":1581}],"modifiedIndex":1581,"createdIndex":1581}}' + response.content = '{"action":"get","node":{"key":"/service/batman5","dir":true,"nodes":[{"key":"/service/batman5/initialize","value":"postgresql0","modifiedIndex":1582,"createdIndex":1582},{"key":"/service/batman5/leader","value":"postgresql1","expiration":"2015-05-15T09:11:00.037397538Z","ttl":21,"modifiedIndex":20728,"createdIndex":20434},{"key":"/service/batman5/optime","dir":true,"nodes":[{"key":"/service/batman5/optime/leader","value":"2164261704","modifiedIndex":20729,"createdIndex":20729}],"modifiedIndex":20437,"createdIndex":20437},{"key":"/service/batman5/members","dir":true,"nodes":[{"key":"/service/batman5/members/postgresql0","value":"postgres://replicator:rep-pass@127.0.0.1:5433/postgres?application_name=http://127.0.0.1:8008/governor","expiration":"2015-05-15T09:11:09.611860899Z","ttl":30,"modifiedIndex":20730,"createdIndex":20730}],"modifiedIndex":1581,"createdIndex":1581}],"modifiedIndex":1581,"createdIndex":1581}}' return response