From 494565e6bb184323afaafbfdd05cf2da99c63983 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Fri, 18 Dec 2015 17:15:40 +0100 Subject: [PATCH] Patroni changes to support BDR. Currently, BDR and physical replication at the same time is not supported. BDR requires additional postgresql configuration and a patched version of PostgreSQL 9.4 + BDR plugin. Included is the Docker image to try it locally. --- Dockerfile | 13 +++- README.rst | 28 ++++++++ docker/entrypoint.sh | 14 ++-- patroni/__init__.py | 3 +- patroni/api.py | 2 +- patroni/ha.py | 130 +++++++++++++++++++++++++++++------ patroni/postgresql.py | 144 ++++++++++++++++++++++++++++++++------- postgres0.yml | 14 ++-- postgres1.yml | 14 ++-- tests/test_ha.py | 1 + tests/test_postgresql.py | 4 +- 11 files changed, 305 insertions(+), 62 deletions(-) diff --git a/Dockerfile b/Dockerfile index 84b9cdf6..36a8d82e 100644 --- a/Dockerfile +++ b/Dockerfile @@ -6,15 +6,22 @@ MAINTAINER Feike Steenbergen # We need curl RUN apt-get update -y && apt-get install curl -y -# Add PGDG repositories +# Add PGDG and BDR repositories RUN echo "deb http://apt.postgresql.org/pub/repos/apt/ $(lsb_release -cs)-pgdg main" > /etc/apt/sources.list.d/pgdg.list +RUN echo "deb http://packages.2ndquadrant.com/bdr/apt/ $(lsb_release -cs)-2ndquadrant main" >> /etc/apt/sources.list.d/pgdg.list RUN curl https://www.postgresql.org/media/keys/ACCC4CF8.asc | apt-key add - +# import the BDR key +RUN curl -s -o - http://packages.2ndquadrant.com/bdr/apt/AA7A6805.asc | sudo apt-key add - + RUN apt-get update -y RUN apt-get upgrade -y ENV PGVERSION 9.4 -RUN apt-get install python python-yaml python-requests python-boto postgresql-${PGVERSION} python-dnspython python-kazoo python-pip -y -RUN apt-get install python-dev postgresql-server-dev-${PGVERSION} -y +ENV PGTYPE postgresql-bdr +ENV PGCOMPATIBLETYPE postgresql + +RUN apt-get install python python-yaml python-requests python-boto ${PGTYPE}-${PGVERSION} ${PGTYPE}-${PGVERSION}-bdr-plugin python-dnspython python-kazoo python-pip -y +RUN apt-get install python-dev ${PGTYPE}-server-dev-${PGVERSION} -y RUN pip install python-etcd psycopg2 ENV PATH /usr/lib/postgresql/${PGVERSION}/bin:$PATH diff --git a/README.rst b/README.rst index 422e0849..a95af16c 100644 --- a/README.rst +++ b/README.rst @@ -75,6 +75,10 @@ For an example file, see ``postgres0.yml``. Regarding settings: - *port*: Exhibitor port. - *hosts*: initial list of Exhibitor (ZooKeeper) nodes in format: ['host1', 'host2', 'etc...' ]. This list updates automatically whenever the Exhibitor (ZooKeeper) cluster topology changes. +- *bdr*: + - *enable*: on if you want to enable BDR + - *database*: database name to support BDR (only a single database is supported) + - *postgresql*: - *name*: the name of the Postgres host. Must be unique for the cluster. - *listen*: IP address + port that Postgres listens to; must be accessible from other nodes in the cluster, if you're using streaming replication. @@ -163,6 +167,30 @@ Choosing your replication schema is dependent on your business considerations. Investigate both async and sync replication, as well as other HA solutions, to determine which solution is best for you. +You can also use BDR (bi-directional replication) if you have a compatible +PostgreSQL version with the BDR plugin installed (see http://bdr-project.org/docs/next/installation.html). +It will require adding a BDR session to your configuration, as well as +setting the following options for postgresql (see http://bdr-project.org/docs/next/settings-prerequisite.html): + +.. code:: YAML + max_worker_processes: 10 + max_replication_slots: 10 + max_wal_senders: 10 + shared_preload_libraries: 'bdr' + track_commit_timestamp: 'on' + wal_level: 'logical' + +At the moment Patroni BDR is not compatible with a streaming replication, +if BDR is enabled normal replica node won't be able to join. This is not +a principal limitation of BDR, and we might resolve this in the future +(although 'promotion' will only work between nodes running a physical +replication, i.e. a replica won't be able to attach to a different +multimaster node). + +Another limitation is that only one database is supported at the moment. +BDR requires the replication user to be also a superuser, so you might +want to excersie extra caution when choosing the password for this user. + Applications Should Not Use Superusers -------------------------------------- diff --git a/docker/entrypoint.sh b/docker/entrypoint.sh index 9edb2120..4ec35268 100755 --- a/docker/entrypoint.sh +++ b/docker/entrypoint.sh @@ -92,6 +92,9 @@ etcd: scope: *scope ttl: *ttl host: ${ETCD_CLUSTER} +bdr: + enable: 'on' + database: 'bdrtest' postgresql: name: ${HOSTNAME} scope: *scope @@ -124,16 +127,19 @@ postgresql: archive_timeout: 1800s max_replication_slots: 20 hot_standby: "on" + max_worker_processes: 10 + max_replication_slots: 10 + max_wal_senders: 10 + shared_preload_libraries: 'bdr' + track_commit_timestamp: true + wal_level: 'logical' __EOF__ cat /patroni/postgres.yml if [ ! -z $CHEAT ] then - while : - do - sleep 60 - done + exec bash else exec python /patroni.py /patroni/postgres.yml fi diff --git a/patroni/__init__.py b/patroni/__init__.py index 045d0d59..2ed7fd93 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -19,7 +19,8 @@ class Patroni: def __init__(self, config): self.nap_time = config['loop_wait'] self.tags = config.get('tags', dict()) - self.postgresql = Postgresql(config['postgresql']) + self.bdr = config.get('bdr', dict()) + self.postgresql = Postgresql(config['postgresql'], self.bdr) self.dcs = self.get_dcs(self.postgresql.name, config) self.api = RestApiServer(self, config['restapi']) self.ha = Ha(self) diff --git a/patroni/api.py b/patroni/api.py index 47416b10..314e199a 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -281,7 +281,7 @@ class RestApiServer(ThreadingMixIn, HTTPServer, Thread): def query(self, sql, *params): cursor = None try: - with self.patroni.postgresql.connection().cursor() as cursor: + with self.patroni.postgresql.connection('postgres').cursor() as cursor: cursor.execute(sql, params) return [r for r in cursor] except psycopg2.Error as e: diff --git a/patroni/ha.py b/patroni/ha.py index 3a8d46d7..1f0ee643 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -20,6 +20,11 @@ class Ha: self.cluster = None self.old_cluster = None self._async_executor = AsyncExecutor() + self.bdr = patroni.bdr + self.use_bdr = self.bdr and\ + self.bdr.get('enable') and\ + (self.bdr.get('database') is not None) + self.bdr_may_need_reinitialize = False def load_cluster_from_dcs(self): cluster = self.dcs.get_cluster() @@ -46,7 +51,7 @@ class Ha: logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name) return lock_owner == self.state_handler.name - def touch_member(self): + def touch_member(self, ttl=None): data = { 'conn_url': self.state_handler.connection_string, 'api_url': self.patroni.api.connection_string, @@ -59,7 +64,7 @@ class Ha: data['xlog_location'] = self.state_handler.xlog_position() except: pass - self.dcs.touch_member(json.dumps(data, separators=(',', ':'))) + self.dcs.touch_member(json.dumps(data, separators=(',', ':')), ttl=ttl) def copy_backup_from_leader(self, leader): if self.state_handler.bootstrap(leader): @@ -69,15 +74,29 @@ class Ha: self.state_handler.remove_data_directory() logger.error('failed to bootstrap from leader') - def bootstrap(self): + def initialize_bdr(self, leader): + if self.state_handler.initialize_bdr_database(leader): + logger.info("bootstrapped from leader") + else: + self.state_handler.stop('immediate') + self.state_handler.remove_data_directory() + logger.error('failed to bootstrap from leader') + + def bootstrap(self, bdr=False): if not self.cluster.is_unlocked(): # cluster already has leader self._async_executor.schedule('bootstrap from leader') - self._async_executor.run_async(self.copy_backup_from_leader, args=(self.cluster.leader, )) + if not bdr: + self._async_executor.run_async(self.copy_backup_from_leader, args=(self.cluster.leader, )) + else: + self._async_executor.run_async(self.initialize_bdr, args=(self.cluster.leader, )) return 'trying to bootstrap from leader' elif not self.cluster.initialize and not self.patroni.nofailover: # no initialize key if self.dcs.initialize(create_new=True): # race for initialization try: - self.state_handler.bootstrap() + if not bdr: + self.state_handler.bootstrap() + else: + self.state_handler.initialize_bdr_database() self.dcs.initialize(create_new=False, sysid=self.state_handler.sysid) except: # initdb or start failed # remove initialization key and give a chance to other members @@ -93,6 +112,9 @@ class Ha: else: return 'waiting for leader to bootstrap' + def bootstrap_bdr(self): + return self.bootstrap(bdr=True) + def recover(self): has_lock = self.has_lock() @@ -118,6 +140,22 @@ class Ha: logger.info('started as readonly because i had the session lock') self.load_cluster_from_dcs() + def recover_bdr(self): + if not self.bdr_may_need_reinitialize: + self.state_handler.start_bdr() + return "started as a member of BDR group" + has_lock = self.has_lock() + if self.cluster.is_unlocked(): # no leader yet, and we need to reinitialize from the leader + return "waiting for another member of the BDR group" + if has_lock: + self.dcs.delete_leader() + return "released the lock to another member of the BDR group" + if self.state_handler.initialize_bdr_database(self.cluster.leader, empty=False): + # we reinitialized this node, make sure the cluster is aware of it. + self.touch_member() + self.bdr_may_need_reinitialize = False + return "reinitialized and started as a member of the BDR group" + def follow_the_leader(self, demote_reason, follow_reason, refresh=True): refresh and self.load_cluster_from_dcs() ret = demote_reason if self.state_handler.is_leader() else follow_reason @@ -268,14 +306,18 @@ class Ha: self.dcs.reset_cluster() self.state_handler.follow_the_leader(None) - def process_manual_failover_from_leader(self): + def process_manual_failover_from_leader(self, bdr=False): failover = self.cluster.failover if not failover.leader or failover.leader == self.state_handler.name: if not failover.member or failover.member != self.state_handler.name: members = [m for m in self.cluster.members if not failover.member or m.name == failover.member] if self.is_failover_possible(members): # check that there are healthy members - self._async_executor.schedule('manual failover: demote') - self._async_executor.run_async(self.demote) + if not bdr: + self._async_executor.schedule('manual failover: demote') + self._async_executor.run_async(self.demote) + else: + self.dcs.delete_leader() + self.dcs.reset_cluster() return 'manual failover: demoting myself' else: logger.warning('manual failover: no healthy members found, failover is not possible') @@ -288,6 +330,9 @@ class Ha: logger.info('Trying to clean up failover key') self.dcs.manual_failover('', '', self.cluster.failover.index) + def process_manual_failover_from_leader_bdr(self): + self.process_manual_failover_from_leader(bdr=True) + def process_unhealthy_cluster(self): if self.is_healthiest_node(): if self.acquire_lock(): @@ -326,6 +371,31 @@ class Ha: return self.follow_the_leader('demoting self because i do not have the lock and i was a leader', 'no action. i am a secondary and i am following a leader', False) + def process_unhealthy_cluster_bdr(self): + if self.acquire_lock(): + if self.cluster.failover: + logger.info('Cleaning up failover key after acquiring leader lock...') + self.dcs.manual_failover('', '') + self.dcs.get_cluster() + return "acquired the session lock as a leader" + else: + return "lost the race to acquire the session lock" + + def process_healthy_cluster_bdr(self): + if self.has_lock(): + if self.cluster.failover: + msg = self.process_manual_failover_from_leader_bdr() + if msg is not None: + return msg + if self.update_lock(): + return "no action. I am the leader with the lock" + else: + logger.error("failed to update leader lock") + self.load_cluster_from_dcs() + else: + logger.info("does not have lock") + return "no action. I do not have the lock" + def schedule(self, action): with self._async_executor: return self._async_executor.schedule(action) @@ -386,7 +456,14 @@ class Ha: try: self.load_cluster_from_dcs() - self.touch_member() + # if member key is missing - we may need to destory + # and re-create the BDR database unless we are the first node in the cluster + if self.use_bdr and len([m for m in self.cluster.members if m.name == self.state_handler.name]) == 0: + self.bdr_may_need_reinitialize = True + + # Create a member key with a very short expiration time if we may need to reinitialize this member + # in order to avoid lossing the "may need to reinitialize status" in case Patroni is restarted + self.touch_member(ttl=None if not self.bdr_may_need_reinitialize else self.bdr.get('min_ttl', 5)) # cluster has leader key but not initialize key if not self.cluster.is_unlocked() and not self.sysid_valid(self.cluster.initialize) and self.has_lock(): @@ -402,33 +479,46 @@ class Ha: # is data directory empty? if self.state_handler.data_directory_empty(): - return self.bootstrap() # new node + if self.use_bdr: + # member key was missing, but it's a new database - no need to reinitialize + self.bdr_may_need_reinitialize = False + self.touch_member() + return self.bootstrap(self.use_bdr) # new node # "bootstrap", but data directory is not empty - elif not self.sysid_valid(self.cluster.initialize) and self.cluster.is_unlocked(): + elif not self.sysid_valid(self.cluster.initialize) and self.cluster.is_unlocked() and \ + not self.bdr_may_need_reinitialize: self.dcs.initialize(create_new=(self.cluster.initialize is None), sysid=self.state_handler.sysid) else: # check if we are allowed to join - if self.sysid_valid(self.cluster.initialize) and self.cluster.initialize != self.state_handler.sysid: + if not self.use_bdr and \ + self.sysid_valid(self.cluster.initialize) and self.cluster.initialize != self.state_handler.sysid: logger.fatal("system ID mismatch, node {0} belongs to a different cluster". format(self.state_handler.name)) sys.exit(1) # try to start dead postgres - if not self.state_handler.is_healthy(): - msg = self.recover() - if msg is not None: - return msg + if self.use_bdr: + if self.bdr_may_need_reinitialize or not self.state_handler.is_healthy(): + msg = self.recover_bdr() + if msg is not None: + return msg + else: + if not self.state_handler.is_healthy(): + msg = self.recover() + if msg is not None: + return msg try: if self.cluster.is_unlocked(): - return self.process_unhealthy_cluster() + return self.process_unhealthy_cluster() if not self.use_bdr else self.process_unhealthy_cluster_bdr() else: - return self.process_healthy_cluster() + return self.process_healthy_cluster() if not self.use_bdr else self.process_healthy_cluster_bdr() finally: - self.state_handler.sync_replication_slots(self.cluster) + if not self.use_bdr: + self.state_handler.sync_replication_slots(self.cluster) except DCSError: logger.error('Error communicating with DCS') - if self.state_handler.is_running() and self.state_handler.is_leader(): + if self.state_handler.is_running() and self.state_handler.is_leader() and not self.use_bdr: self.demote(delete_leader=False) return 'demoted self because DCS is not accessible and i was a leader' except (psycopg2.Error, PostgresConnectionException): diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 45b871c2..817b30a4 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -1,6 +1,7 @@ import logging import os import psycopg2 +import re import shlex import shutil import subprocess @@ -41,8 +42,9 @@ def parseurl(url): class Postgresql: - def __init__(self, config): + def __init__(self, config, config_bdr={}): self.config = config + self.bdr = config_bdr self.name = config['name'] self.server_parameters = config.get('parameters', {}) self.scope = config['scope'] @@ -72,6 +74,7 @@ class Postgresql: connect_address=connect_address, **self.replication) self._connection = None + self._connection_db = None self._cursor_holder = None self._need_rewind = False self._sysid = None @@ -128,28 +131,32 @@ class Postgresql: break return local_address + ':' + self.port - def connection(self): + def connection(self, dbname): if not self._connection or self._connection.closed != 0: - r = parseurl('postgres://{}/postgres'.format(self.local_address)) + r = parseurl('postgres://{0}/{1}'.format(self.local_address, dbname)) self._connection = psycopg2.connect(**r) + self._connection_db = dbname self._connection.autocommit = True return self._connection - def _cursor(self): + def _cursor(self, dbname='postgres'): + if self._connection_db != dbname: + self.close_connection() if not self._cursor_holder or self._cursor_holder.closed or self._cursor_holder.connection.closed != 0: - logger.info("established a new patroni connection to the postgres cluster") - self._cursor_holder = self.connection().cursor() + logger.info("established a new patroni connection to the {0} cluster".format(dbname)) + self._cursor_holder = self.connection(dbname).cursor() return self._cursor_holder def close_connection(self): if self._cursor_holder and self._cursor_holder.connection and self._cursor_holder.connection.closed == 0: self._cursor_holder.connection.close() - logger.info("closed patroni connection to the postgresql cluster") + logger.info("closed patroni connection to the {0} cluster".format(self._connection_db)) + self._connection_db = None - def _query(self, sql, *params): + def _query(self, sql, dbname, *params): cursor = None try: - cursor = self._cursor() + cursor = self._cursor(dbname) cursor.execute(sql, params) return cursor except psycopg2.Error as e: @@ -159,9 +166,10 @@ class Postgresql: raise RetryFailedError('cluster is being restarted') raise PostgresConnectionException('connection problems') - def query(self, sql, *params): + def query(self, sql, *params, **kwargs): + dbname = kwargs.get('dbname', 'postgres') try: - return self.retry(self._query, sql, *params) + return self.retry(self._query, sql, dbname, *params) except RetryFailedError as e: raise PostgresConnectionException(str(e)) @@ -207,6 +215,81 @@ class Postgresql: self.set_state('initdb failed') return ret + def initialize_bdr_database(self, leader=None, empty=True): + """ initialize BDR database. If empty = False, reinitialize an existing one. + If leader is None, create a new group, otherwise, join an exising one. + non-empty database with an empty leader is not supported + """ + if not empty and not leader: + ret = False + else: + try: + ret = (self.initialize() and self.start_bdr()) if empty else\ + (self.drop_bdr_database() and self.start_bdr()) + except: + logger.exception("exception") + ret = False + if ret: + # convert all connection URLs to name=value strings + leader_dsn = self.primary_conninfo(leader.conn_url, include_db=self.bdr['database']) if leader else None + self_dsn = self.primary_conninfo(self.connection_string, include_db=self.bdr['database']) + try: + self.create_replication_user(superuser=True) + ret = self.create_bdr_database() and (self.join_bdr_group(self_dsn, leader_dsn) if leader + else self.create_bdr_group(self_dsn)) + except: + logger.exception("exception") + ret = False + if not ret and not leader: + raise Exception("Could not bootstrap initial PostgreSQL BDR node") + return ret + + def start_bdr(self): + logger.info("starting new BDR instance {0}".format(self.name)) + return self.start(bdr=True) + + def drop_bdr_database(self): + logger.info("Dropping BDR database for instance {0}".format(self.name)) + # stop the database if it's running in order to turn off bdr + if self.is_running(): + self.stop() + ret = self.start(bdr=False) + if ret: + self.query("DROP DATABASE IF EXISTS {0}".format(self.bdr['database'])) + return ret and self.restart(bdr=True) + + def create_bdr_database(self): + logger.info("Creating BDR database and extensions for instance {0}".format(self.name)) + self.query("CREATE DATABASE {0}".format(self.bdr['database'])) + self.query("CREATE EXTENSION IF NOT EXISTS btree_gist", dbname=self.bdr['database']) + self.query("CREATE EXTENSION IF NOT EXISTS bdr", dbname=self.bdr['database']) + return True + + def join_bdr_group(self, self_dsn, join_dsn): + logger.info("Joining BDR group {0}, host {1}, dsn {2}".format(join_dsn, self.name, self_dsn)) + self.query("SELECT bdr.bdr_group_join(local_node_name:=%s, join_using_dsn:=%s, node_external_dsn:=%s)", + self.name, join_dsn, self_dsn, dbname=self.bdr['database']) + return True + + def create_bdr_group(self, self_dsn): + logger.info("Creating new BDR group {0}, host {1}".format(self_dsn, self.name)) + self.query("SELECT bdr.bdr_group_create(local_node_name:=%s, node_external_dsn:=%s)", + self.name, self_dsn, dbname=self.bdr['database']) + return True + + def remove_bdr_nodes(self, nodes): + logger.info("BDR node {0}, removing nodes {1}".format(self.name, nodes)) + self.query("SELECT bdr.part_by_node_names(%s)", nodes, dbname=self.bdr['database']) + + def remove_bdr_disconnected_members(self, cluster): + # get all members + result = self.query("SELECT node_name FROM bdr.bdr_nodes WHERE node_status = 'r'") + if result: + nodes_db = set([node[0] for node in result.fetchall()]) + nodes_etcd = set([member.name for member in self.members]) + nodes_remove = nodes_db.difference(nodes_etcd) + self.remove_bdr_nodes(list(nodes_remove)) + def delete_trigger_file(self): os.path.exists(self.trigger_file) and os.unlink(self.trigger_file) @@ -318,12 +401,15 @@ class Postgresql: with self._state_lock: self._state = value - def start(self, block_callbacks=False): + def start(self, block_callbacks=False, bdr=False): if self.is_running(): logger.error('Cannot start PostgreSQL because one is already running.') return True - self.set_role('replica' if os.path.exists(self.recovery_conf) else 'master') + if not bdr: + self.set_role('replica' if os.path.exists(self.recovery_conf) else 'master') + else: + self.set_role('master') if os.path.exists(self.postmaster_pid): os.remove(self.postmaster_pid) logger.info('Removed %s', self.postmaster_pid) @@ -331,12 +417,13 @@ class Postgresql: if not block_callbacks: self.set_state('starting') - ret = subprocess.call(self._pg_ctl + ['start', '-o', self.server_options()]) == 0 + ret = subprocess.call(self._pg_ctl + ['start', '-o', self.server_options(bdr)]) == 0 self.set_state('running' if ret else 'start failed') - self.schedule_load_slots = ret and self.use_slots - self.save_configuration_files() + if not bdr: + self.schedule_load_slots = ret and self.use_slots + self.save_configuration_files() # block_callbacks is used during restart to avoid # running start/stop callbacks in addition to restart ones ret and not block_callbacks and self.call_nowait(ACTION_ON_START) @@ -385,19 +472,28 @@ class Postgresql: ret and self.call_nowait(ACTION_ON_RELOAD) return ret - def restart(self): + def restart(self, bdr=False): self.set_state('restarting') - ret = self.stop(block_callbacks=True) and self.start(block_callbacks=True) + ret = self.stop(block_callbacks=True) and self.start(block_callbacks=True, bdr=bdr) if ret: self.call_nowait(ACTION_ON_RESTART) else: self.set_state('restart failed ({})'.format(self.state)) return ret - def server_options(self): + def server_options(self, bdr=False): + bdr_set = False options = "--listen_addresses='{}' --port={}".format(self.listen_addresses, self.port) for setting, value in self.server_parameters.items(): + if setting == 'shared_preload_libraries': + if not bdr and 'bdr' in value: + value = ','.join([x for x in value.split(',') if x != 'bdr']) + elif bdr and 'bdr' not in value: + bdr_set = True + value = ','.join([x for x in value.split(',')]+['bdr']) options += " --{}='{}'".format(setting, value) + if bdr and not bdr_set: + options += " --shared_preload_libraries='bdr'" return options def is_healthy(self): @@ -419,9 +515,10 @@ class Postgresql: f.write(line + '\n') @staticmethod - def primary_conninfo(leader_url): + def primary_conninfo(leader_url, include_db=None): r = parseurl(leader_url) - return 'user={user} password={password} host={host} port={port} sslmode=prefer sslcompression=1'.format(**r) + ret = 'user={user} password={password} host={host} port={port} sslmode=prefer sslcompression=1'.format(**r) + return ret if not include_db else ret + ' dbname={0}'.format(include_db) def check_recovery_conf(self, leader): if not os.path.isfile(self.recovery_conf): @@ -612,8 +709,9 @@ BEGIN END; $$""".format(name, options), name, password, password) - def create_replication_user(self): - self.create_or_update_role(self.replication['username'], self.replication['password'], 'REPLICATION') + def create_replication_user(self, superuser=False): + options = 'REPLICATION SUPERUSER' if superuser else 'REPLICATION' + self.create_or_update_role(self.replication['username'], self.replication['password'], '{0}'.format(options)) def create_connection_users(self): if 'username' in self.superuser: diff --git a/postgres0.yml b/postgres0.yml index 36018ad5..4cf43dc6 100644 --- a/postgres0.yml +++ b/postgres0.yml @@ -4,7 +4,7 @@ scope: &scope batman restapi: listen: 127.0.0.1:8008 connect_address: 127.0.0.1:8008 - auth: 'username:password' + #auth: 'username:password' # certfile: /etc/ssl/certs/ssl-cert-snakeoil.pem # keyfile: /etc/ssl/private/ssl-cert-snakeoil.key etcd: @@ -26,6 +26,9 @@ etcd: # - host1 # - host2 # - host3 +bdr: + enable: "on" + database: "bdrtest" postgresql: name: postgresql0 scope: *scope @@ -86,12 +89,15 @@ postgresql: restore_command: cp ../wal_archive/%f %p parameters: archive_mode: "on" - wal_level: hot_standby + wal_level: logical archive_command: mkdir -p ../wal_archive && test ! -f ../wal_archive/%f && cp %p ../wal_archive/%f - max_wal_senders: 5 + max_wal_senders: 10 wal_keep_segments: 8 archive_timeout: 1800s - max_replication_slots: 5 + max_replication_slots: 10 + max_worker_processes: 10 + track_commit_timestamp: 'on' + shared_preload_libraries: 'bdr' hot_standby: "on" wal_log_hints: "on" tags: diff --git a/postgres1.yml b/postgres1.yml index e1b61b3b..2cf62845 100644 --- a/postgres1.yml +++ b/postgres1.yml @@ -4,7 +4,7 @@ scope: &scope batman restapi: listen: 127.0.0.1:8009 connect_address: 127.0.0.1:8009 - auth: 'username:password' + #auth: 'username:password' # certfile: /etc/ssl/certs/ssl-cert-snakeoil.pem # keyfile: /etc/ssl/private/ssl-cert-snakeoil.key etcd: @@ -26,6 +26,9 @@ etcd: # - host1 # - host2 # - host3 +bdr: + enable: "on" + database: "bdrtest" postgresql: name: postgresql1 scope: *scope @@ -86,12 +89,15 @@ postgresql: restore_command: cp ../wal_archive/%f %p parameters: archive_mode: "on" - wal_level: hot_standby + wal_level: logical archive_command: mkdir -p ../wal_archive && test ! -f ../wal_archive/%f && cp %p ../wal_archive/%f - max_wal_senders: 5 + max_wal_senders: 10 + max_worker_processes: 10 + track_commit_timestamp: 'on' + shared_preload_libraries: 'bdr' wal_keep_segments: 8 archive_timeout: 1800s - max_replication_slots: 5 + max_replication_slots: 10 hot_standby: "on" wal_log_hints: "on" tags: diff --git a/tests/test_ha.py b/tests/test_ha.py index d9a408a4..daab936a 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -88,6 +88,7 @@ class MockPatroni: self.api = Mock() self.tags = {} self.nofailover = None + self.bdr = {} self.api.connection_string = 'http://127.0.0.1:8008' diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 2097bcb7..721e1ff4 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -310,9 +310,9 @@ class TestPostgresql(unittest.TestCase): @patch.object(MockConnect, 'closed', 2) def test__query(self): - self.assertRaises(PostgresConnectionException, self.p._query, 'blabla') + self.assertRaises(PostgresConnectionException, self.p._query, 'blabla', 'postgres') self.p._state = 'restarting' - self.assertRaises(RetryFailedError, self.p._query, 'blabla') + self.assertRaises(RetryFailedError, self.p._query, 'blabla', 'postgres') def test_query(self): self.p.query('select 1')