diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 09571b80..a68b8dfd 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -9,7 +9,6 @@ import time from patroni.exceptions import PostgresConnectionException, PostgresException from patroni.utils import Retry, RetryFailedError from six.moves.urllib_parse import urlparse -from threading import Lock logger = logging.getLogger(__name__) @@ -48,8 +47,6 @@ class Postgresql: self.replication = config['replication'] self.superuser = config['superuser'] self.admin = config['admin'] - self.pgpass = config.get('pgpass', None) or os.path.join(os.path.expanduser('~'), 'pgpass') - self.pg_rewind = config.get('pg_rewind', {}) self.callback = config.get('callbacks', {}) self.use_slots = config.get('use_slots', True) self.schedule_load_slots = self.use_slots @@ -59,6 +56,7 @@ class Postgresql: 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._role = 'replica' self._pg_ctl = ['pg_ctl', '-w', '-D', self.data_dir] @@ -69,50 +67,8 @@ class Postgresql: self._connection = None self._cursor_holder = None - self._need_rewind = False - self._sysid = None - self.replication_slots = [] # list of already existing replication slots - self.retry = Retry(max_tries=-1, deadline=5, max_delay=1, retry_exceptions=PostgresConnectionException) - - self._state = 'stopped' - self._state_lock = Lock() - self._role = 'replica' - self._role_lock = Lock() - - if self.is_running(): - self._state = 'running' - self._role = 'master' if self.is_leader() else 'replica' - - @property - def can_rewind(self): - """ check if pg_rewind executable is there and that pg_controldata indicates - we have either wal_log_hints or checksums turned on - """ - # low-hanging fruit: check if pg_rewind configuration is there - if not self.pg_rewind or\ - not (self.pg_rewind.get('username', '') and self.pg_rewind.get('password', '')): - return False - - cmd = ['pg_rewind', '--help'] - try: - ret = subprocess.call(cmd, stdout=open(os.devnull, 'w'), stderr=subprocess.STDOUT) - if ret != 0: # pg_rewind is not there, close up the shop and go home - return False - except OSError: - return False - # check if the cluster's configuration permits pg_rewind - data = self.controldata() - return data.get('wal_log_hints setting', 'off') == 'on' or data.get('Data page checksum version', '0') != '0' - - @property - def sysid(self): - if not self._sysid: - data = self.controldata() - self._sysid = data.get('Database system identifier', "") - return self._sysid - - def require_rewind(self): - self._need_rewind = True + self.members = [] # list of already existing replication slots + self.retry = Retry(max_tries=-1, deadline=10, max_delay=1, retry_exceptions=PostgresConnectionException) def get_local_address(self): listen_addresses = self.listen_addresses.split(',') @@ -133,15 +89,9 @@ class Postgresql: def _cursor(self): 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() 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") - def _query(self, sql, *params): cursor = None try: @@ -151,8 +101,6 @@ class Postgresql: except psycopg2.Error as e: if cursor and cursor.connection.closed == 0: raise e - if self.state == 'restarting': - raise RetryFailedError('cluster is being restarted') raise PostgresConnectionException('connection problems') def query(self, sql, *params): @@ -165,30 +113,23 @@ class Postgresql: return not os.path.exists(self.data_dir) or os.listdir(self.data_dir) == [] def initialize(self): - self.set_state('initalizing new cluster') ret = subprocess.call(self._pg_ctl + ['initdb', '-o', '--encoding=UTF8']) == 0 - if ret: - self.write_pg_hba() - else: - self.set_state('initdb failed') + ret and self.write_pg_hba() return ret def delete_trigger_file(self): os.path.exists(self.trigger_file) and os.unlink(self.trigger_file) - def write_pgpass(self, record): - with open(self.pgpass, 'w') as f: - os.fchmod(f.fileno(), 0o600) - f.write('{host}:{port}:*:{user}:{password}\n'.format(**record)) - - env = os.environ.copy() - env['PGPASSFILE'] = self.pgpass - return env - def sync_from_leader(self, leader): r = parseurl(leader.conn_url) - env = self.write_pgpass(r) + pgpass = 'pgpass' + with open(pgpass, 'w') as f: + os.fchmod(f.fileno(), 0o600) + f.write('{host}:{port}:*:{user}:{password}\n'.format(**r)) + + env = os.environ.copy() + env['PGPASSFILE'] = pgpass return self.create_replica(r, env) == 0 @staticmethod @@ -272,82 +213,40 @@ class Postgresql: @property def role(self): - with self._role_lock: - return self._role - - def set_role(self, value): - with self._role_lock: - self._role = value - - @property - def state(self): - with self._state_lock: - return self._state - - def set_state(self, value): - with self._state_lock: - self._state = value + return self._role def start(self, block_callbacks=False): if self.is_running(): + self._role = 'master' if self.is_leader() else 'replica' + self.schedule_load_slots = self.use_slots logger.error('Cannot start PostgreSQL because one is already running.') - return True + return False - self.set_role('replica' if os.path.exists(self.recovery_conf) else 'master') + self._role = 'replica' if os.path.exists(self.recovery_conf) else 'master' if os.path.exists(self.postmaster_pid): os.remove(self.postmaster_pid) logger.info('Removed %s', self.postmaster_pid) - if not block_callbacks: - self.set_state('starting') - ret = subprocess.call(self._pg_ctl + ['start', '-o', self.server_options()]) == 0 - - self.set_state('running' if ret else 'start failed') - 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) + ret and not block_callbacks and ret and self.call_nowait(ACTION_ON_START) return ret - def checkpoint(self): - try: - r = parseurl('postgres://{}/postgres'.format(self.local_address)) - r['options'] = '-c statement_timeout=0' - with psycopg2.connect(**r) as conn: - conn.autocommit = True - with conn.cursor() as cur: - cur.execute('CHECKPOINT') - except: - logging.exception('Exception during CHECKPOINT') - def stop(self, mode='fast', block_callbacks=False): - # make sure we close all connections established against - # the former node, otherwise, we might get a stalled one - # after kill -9, which would report incorrect data to - # patroni. - - self.close_connection() - if not self.is_running(): - if not block_callbacks: - self.set_state('stopped') - return True - if block_callbacks: - self.checkpoint() - else: - self.set_state('stopping') + try: + self.query('SET statement_timeout TO 0') + self.query('CHECKPOINT') + except: + logging.exception('Exception diring CHECKPOINT') ret = subprocess.call(self._pg_ctl + ['stop', '-m', mode]) == 0 # block_callbacks is used during restart to avoid # running start/stop callbacks in addition to restart ones - if not ret: - self.set_state('stop failed') - elif not block_callbacks: - self.set_state('stopped') - self.call_nowait(ACTION_ON_STOP) + ret and not block_callbacks and self.call_nowait(ACTION_ON_STOP) return ret def reload(self): @@ -356,12 +255,8 @@ class Postgresql: return ret def restart(self): - self.set_state('restarting') ret = self.stop(block_callbacks=True) and self.start(block_callbacks=True) - if ret: - self.call_nowait(ACTION_ON_RESTART) - else: - self.set_state('restart failed ({})'.format(self.state)) + ret and self.call_nowait(ACTION_ON_RESTART) return ret def server_options(self): @@ -376,9 +271,36 @@ class Postgresql: return False return True - def check_replication_lag(self, last_leader_operation): - return (last_leader_operation if last_leader_operation else 0) - self.xlog_position() <=\ - self.config.get('maximum_lag_on_failover', 0) + def is_healthiest_node(self, cluster): + if self.is_leader(): + return True + + if cluster.last_leader_operation - self.xlog_position() > self.config.get('maximum_lag_on_failover', 0): + return False + + for member in cluster.members: + if member.name == self.name: + continue + try: + r = parseurl(member.conn_url) + member_conn = psycopg2.connect(**r) + member_conn.autocommit = True + member_cursor = member_conn.cursor() + member_cursor.execute( + "SELECT pg_is_in_recovery(), %s - pg_xlog_location_diff(pg_last_xlog_replay_location(), '0/0')", + (self.xlog_position(),)) + row = member_cursor.fetchone() + member_cursor.close() + member_conn.close() + logger.error([self.name, member.name, row]) + if not row[0]: + logger.warning('Master (%s) is still alive', member.name) + return False + if row[1] < 0: + return False + except psycopg2.Error: + continue + return True def write_pg_hba(self): with open(os.path.join(self.data_dir, 'pg_hba.conf'), 'a') as f: @@ -402,7 +324,10 @@ class Postgresql: with open(self.recovery_conf, 'r') as f: for line in f: if line.startswith('primary_conninfo'): - return pattern and (pattern in line) + if not pattern: + return False + return pattern in line + return not pattern def write_recovery_conf(self, leader): @@ -417,154 +342,40 @@ recovery_target_timeline = 'latest' for name, value in self.config.get('recovery_conf', {}).items(): f.write("{} = '{}'\n".format(name, value)) - def rewind(self, leader): - # prepare pg_rewind connection - r = parseurl(leader.conn_url) - r.update(self.pg_rewind) - r['user'] = r['username'] - env = self.write_pgpass(r) - pc = "user={user} host={host} port={port} dbname=postgres sslmode=prefer sslcompression=1".format(**r) - logger.info("running pg_rewind from {}".format(pc)) - pg_rewind = ['pg_rewind', '-D', self.data_dir, '--source-server', pc] - try: - ret = (subprocess.call(pg_rewind, env=env) == 0) - except: - ret = False - if ret: + def follow_the_leader(self, leader): + if not self.check_recovery_conf(leader): self.write_recovery_conf(leader) - return ret - - def controldata(self): - """ return the contents of pg_controldata, or non-True value if pg_controldata call failed """ - result = {} - try: - data = subprocess.check_output(['pg_controldata', self.data_dir]) - if data: - data = data.decode().splitlines() - result = {l.split(':')[0].replace('Current ', '', 1): l.split(':')[1].strip() for l in data if l} - except subprocess.CalledProcessError: - logger.exception("Error when calling pg_controldata") - finally: - return result - - def read_postmaster_opts(self): - """ returns the list of option names/values from postgres.opts, Empty dict if read failed or no file """ - result = {} - try: - with open(os.path.join(self.data_dir, "postmaster.opts")) as f: - data = f.read() - opts = [opt.strip('"\n') for opt in data.split(' "')] - for opt in opts: - if '=' in opt and opt.startswith('--'): - name, val = opt.split('=', 1) - name = name.strip('-') - result[name] = val - except IOError: - logger.exception('Error when reading postmaster.opts') - finally: - return result - - def single_user_mode(self, command=None, options={}): - """ run a given command in a single-user mode. If the command is empty - then just start and stop """ - cmd = ['postgres', '--single', '-D', self.data_dir] - for opt in sorted(options): - cmd.extend(['-c', '{0}={1}'.format(opt, options[opt])]) - # need a database name to connect - cmd.append('postgres') - p = subprocess.Popen(cmd, stdin=subprocess.PIPE, stdout=open(os.devnull, 'w'), stderr=subprocess.STDOUT) - if p: - command and p.communicate('{}\n'.format(command)) - p.stdin.close() - return p.wait() - return 1 - - def cleanup_archive_status(self): - status_dir = os.path.join(self.data_dir, 'pg_xlog', 'archive_status') - if os.path.isdir(status_dir): - for f in os.listdir(status_dir): - path = os.path.join(status_dir, f) - try: - if os.path.islink(path): - os.unlink(path) - elif os.path.isfile(path): - os.remove(path) - except: - logger.exception("Unable to remove {}".format(path)) - - def follow_the_leader(self, leader, recovery=False): - if not self.check_recovery_conf(leader) or recovery: - change_role = (self.role == 'master') - - self._need_rewind = (self._need_rewind or change_role) and self.can_rewind - if self._need_rewind: - logger.info("set the rewind flag after demote") - self.write_recovery_conf(leader) - if not leader or not self._need_rewind: # do not rewind until the leader becomes available - ret = self.restart() - else: # we have a leader and need to rewind - if self.is_running(): - self.stop() - # at present, pg_rewind only runs when the cluster is shut down cleanly - # and not shutdown in recovery. We have to remove the recovery.conf if present - # and start/shutdown in a single user mode to emulate this. - # XXX: if recovery.conf is linked, it will be written anew as a normal file. - if os.path.islink(self.recovery_conf): - os.unlink(self.recovery_conf) - else: - os.remove(self.recovery_conf) - # Archived segments might be useful to pg_rewind, - # clean the flags that tell we should remove them. - self.cleanup_archive_status() - # Start in a single user mode and stop to produce a clean shutdown - opts = self.read_postmaster_opts() - opts['archive_mode'] = 'on' - opts['archive_command'] = 'false' - self.single_user_mode(options=opts) - if self.rewind(leader): - ret = self.start() - else: - logger.error("unable to rewind the former master") - self.remove_data_directory() - ret = True - self._need_rewind = False - change_role and ret and self.call_nowait(ACTION_ON_ROLE_CHANGE) - return ret - else: - return True + run_callback = self.role == 'master' + self.restart() + run_callback and self.call_nowait(ACTION_ON_ROLE_CHANGE) def save_configuration_files(self): """ - copy postgresql.conf to postgresql.conf.backup to be able to retrive configuration files - - originally stored as symlinks, those are normally skipped by pg_basebackup - - in case of WAL-E basebackup (see http://comments.gmane.org/gmane.comp.db.postgresql.wal-e/239) + copy postgresql.conf to postgresql.conf.backup to preserve it in the WAL-e backup. + see http://comments.gmane.org/gmane.comp.db.postgresql.wal-e/239 """ - try: - for f in self.configuration_to_save: - os.path.isfile(f) and shutil.copy(f, f + '.backup') - except: - logger.exception('unable to create backup copies of configuration files') + for f in self.configuration_to_save: + shutil.copy(f, f + '.backup') def restore_configuration_files(self): """ restore a previously saved postgresql.conf """ try: for f in self.configuration_to_save: - not os.path.isfile(f) and os.path.isfile(f+'.backup') and shutil.copy(f + '.backup', f) + shutil.copy(f + '.backup', f) except: - logger.exception('unable to restore configuration files from backup') + logger.exception('unable to restore configuration from WAL-E backup') def promote(self): if self.role == 'master': return True ret = subprocess.call(self._pg_ctl + ['promote']) == 0 if ret: - self.set_role('master') - logger.info("cleared rewind flag after becoming the leader") - self._need_rewind = False + self._role = 'master' self.call_nowait(ACTION_ON_ROLE_CHANGE) return ret - def demote(self): - self.follow_the_leader(None) + def demote(self, leader): + self.follow_the_leader(leader) def create_replication_user(self): self.query('CREATE USER "{}" WITH REPLICATION ENCRYPTED PASSWORD %s'.format( @@ -586,34 +397,31 @@ recovery_target_timeline = 'latest' return self.query("""SELECT pg_xlog_location_diff(CASE WHEN pg_is_in_recovery() THEN pg_last_xlog_replay_location() ELSE pg_current_xlog_location() - END, '0/0')::bigint""").fetchone()[0] + END, '0/0')""").fetchone()[0] def load_replication_slots(self): if self.use_slots and self.schedule_load_slots: cursor = self.query("SELECT slot_name FROM pg_replication_slots WHERE slot_type='physical'") - self.replication_slots = [r[0] for r in cursor] + self.members = [r[0] for r in cursor] self.schedule_load_slots = False def sync_replication_slots(self, cluster): if self.use_slots: - try: - self.load_replication_slots() - slots = [m.name for m in cluster.members if m.name != self.name] if self.role == 'master' else [] - # drop unused slots - for slot in set(self.replication_slots) - set(slots): - self.query("""SELECT pg_drop_replication_slot(%s) - WHERE EXISTS(SELECT 1 FROM pg_replication_slots - WHERE slot_name = %s)""", slot, slot) + self.load_replication_slots() + members = [m.name for m in cluster.members if m.name != self.name] if self.role == 'master' else [] + # drop unused slots + for slot in set(self.members) - set(members): + self.query("""SELECT pg_drop_replication_slot(%s) + WHERE EXISTS(SELECT 1 FROM pg_replication_slots + WHERE slot_name = %s)""", slot, slot) - # create new slots - for slot in set(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) + # create new slots + for slot in set(members) - set(self.members): + self.query("""SELECT pg_create_physical_replication_slot(%s) + WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots + WHERE slot_name = %s)""", slot, slot) - self.replication_slots = slots - except: - logger.exception('Exception when changing replication slots') + self.members = members def last_operation(self): return str(self.xlog_position()) @@ -638,7 +446,6 @@ recovery_target_timeline = 'latest' raise PostgresException("Could not bootstrap master PostgreSQL") else: if self.sync_from_leader(current_leader): - self.restore_configuration_files() self.write_recovery_conf(current_leader) ret = self.start() return ret