From ffde403a0a809ebd354495a1fc11c8340bfa84c4 Mon Sep 17 00:00:00 2001 From: Igor Yanchenko <1504692+yanchenko-igor@users.noreply.github.com> Date: Thu, 20 Feb 2020 09:40:44 +0100 Subject: [PATCH 01/24] Config validator implemented (#1314) --- patroni/__init__.py | 8 +- patroni/postgresql/__init__.py | 6 +- patroni/utils.py | 6 + patroni/validator.py | 379 +++++++++++++++++++++++++++++++++ tests/test_validator.py | 221 +++++++++++++++++++ 5 files changed, 615 insertions(+), 5 deletions(-) create mode 100644 patroni/validator.py create mode 100644 tests/test_validator.py diff --git a/patroni/__init__.py b/patroni/__init__.py index 2d0bc01c..ba7a803e 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -169,15 +169,21 @@ class Patroni(object): def patroni_main(): import argparse from patroni.config import Config, ConfigParseError + from patroni.validator import schema parser = argparse.ArgumentParser() parser.add_argument('--version', action='version', version='%(prog)s {0}'.format(__version__)) + parser.add_argument('--validate-config', action='store_true', help='Run config validator and exit') parser.add_argument('configfile', nargs='?', default='', help='Patroni may also read the configuration from the {0} environment variable' .format(Config.PATRONI_CONFIG_VARIABLE)) args = parser.parse_args() try: - conf = Config(args.configfile) + if args.validate_config: + conf = Config(args.configfile, validator=schema) + sys.exit() + else: + conf = Config(args.configfile) except ConfigParseError as e: if e.value: print(e.value) diff --git a/patroni/postgresql/__init__.py b/patroni/postgresql/__init__.py index 7fba00d3..d4b262d6 100644 --- a/patroni/postgresql/__init__.py +++ b/patroni/postgresql/__init__.py @@ -17,7 +17,7 @@ from patroni.postgresql.misc import parse_history, postgres_major_version_to_int from patroni.postgresql.postmaster import PostmasterProcess from patroni.postgresql.slots import SlotsHandler from patroni.exceptions import PostgresConnectionException -from patroni.utils import Retry, RetryFailedError, polling_loop +from patroni.utils import Retry, RetryFailedError, polling_loop, data_directory_is_empty from threading import current_thread, Lock @@ -266,9 +266,7 @@ class Postgresql(object): def data_directory_empty(self): if self.pg_control_exists(): return False - if not os.path.exists(self._data_dir): - return True - return all(os.name != 'nt' and (n.startswith('.') or n == 'lost+found') for n in os.listdir(self._data_dir)) + return data_directory_is_empty(self._data_dir) def replica_method_options(self, method): return deepcopy(self.config.get(method, {})) diff --git a/patroni/utils.py b/patroni/utils.py index dde96fd4..2d35963e 100644 --- a/patroni/utils.py +++ b/patroni/utils.py @@ -442,3 +442,9 @@ def validate_directory(d, msg="{} {}"): raise PatroniException(msg.format(d, "the directory is not writable")) else: raise PatroniException(msg.format(d, "is not a directory")) + + +def data_directory_is_empty(data_dir): + if not os.path.exists(data_dir): + return True + return all(os.name != 'nt' and (n.startswith('.') or n == 'lost+found') for n in os.listdir(data_dir)) diff --git a/patroni/validator.py b/patroni/validator.py new file mode 100644 index 00000000..d254f0b2 --- /dev/null +++ b/patroni/validator.py @@ -0,0 +1,379 @@ +#!/usr/bin/env python3 +import os +import socket +import re +import subprocess + +from patroni.utils import split_host_port, data_directory_is_empty +from patroni.ctl import find_executable +from patroni.dcs import dcs_modules +from patroni.exceptions import ConfigParseError +from six import string_types + + +def data_directory_empty(data_dir): + if os.path.isfile(os.path.join(data_dir, "global", "pg_control")): + return False + return data_directory_is_empty(data_dir) + + +def validate_connect_address(address): + try: + host, _ = split_host_port(address, None) + except (ValueError, TypeError): + raise ConfigParseError("contains a wrong value") + if host in ["127.0.0.1", "0.0.0.0", "*", "::1"]: + raise ConfigParseError('must not contain "127.0.0.1", "0.0.0.0", "*", "::1"') + return True + + +def validate_host_port(host_port, listen=False, multiple_hosts=False): + try: + hosts, port = split_host_port(host_port, None) + except (ValueError, TypeError): + raise ConfigParseError("contains a wrong value") + else: + if multiple_hosts: + hosts = hosts.split(",") + else: + hosts = [hosts] + for host in hosts: + proto = socket.getaddrinfo(host, "", 0, socket.SOCK_STREAM, 0, socket.AI_PASSIVE) + s = socket.socket(proto[0][0], socket.SOCK_STREAM) + try: + if s.connect_ex((host, port)) == 0: + if listen: + raise ConfigParseError("Port {} is already in use.".format(port)) + elif not listen: + raise ConfigParseError("{} is not reachable".format(host_port)) + except socket.gaierror as e: + raise ConfigParseError(e) + finally: + s.close() + return True + + +def comma_separated_host_port(string): + assert all([validate_host_port(s.strip()) for s in string.split(",")]), "didn't pass the validation" + return True + + +def validate_host_port_listen(host_port): + return validate_host_port(host_port, listen=True) + + +def validate_host_port_listen_multiple_hosts(host_port): + return validate_host_port(host_port, listen=True, multiple_hosts=True) + + +def is_ipv4_address(ip): + try: + socket.inet_aton(ip) + except Exception: + raise ConfigParseError("Is not a valid ipv4 address") + return True + + +def is_ipv6_address(ip): + try: + socket.inet_pton(socket.AF_INET6, ip) + except Exception: + raise ConfigParseError("Is not a valid ipv6 address") + return True + + +def get_major_version(bin_dir=None): + if not bin_dir: + binary = 'postgres' + else: + binary = os.path.join(bin_dir, 'postgres') + version = subprocess.check_output([binary, '--version']).decode() + version = re.match(r'^[^\s]+ [^\s]+ (\d+)(\.(\d+))?', version) + return '.'.join([version.group(1), version.group(3)]) if int(version.group(1)) < 10 else version.group(1) + + +def validate_data_dir(data_dir): + if not data_dir: + raise ConfigParseError("is an empty string") + elif os.path.exists(data_dir) and not os.path.isdir(data_dir): + raise ConfigParseError("is not a directory") + elif not data_directory_empty(data_dir): + if not os.path.exists(os.path.join(data_dir, "PG_VERSION")): + raise ConfigParseError("doesn't look like a valid data directory") + else: + with open(os.path.join(data_dir, "PG_VERSION"), "r") as version: + pgversion = version.read().strip() + waldir = ("pg_wal" if float(pgversion) >= 10 else "pg_xlog") + if not os.path.isdir(os.path.join(data_dir, waldir)): + raise ConfigParseError("data dir for the cluster is not empty, but doesn't contain" + " \"{}\" directory".format(waldir)) + bin_dir = schema.data.get("postgresql", {}).get("bin_dir", None) + major_version = get_major_version(bin_dir) + if pgversion != major_version: + raise ConfigParseError("data_dir directory postgresql version ({}) doesn't match" + "with 'postgres --version' output ({})".format(pgversion, major_version)) + return True + + +class Result(object): + def __init__(self, status, error="didn't pass validation", level=0, path="", data=""): + self.status = status + self.path = path + self.data = data + self.level = level + self._error = error + if not self.status: + self.error = error + else: + self.error = None + + def __repr__(self): + return self.path + (" " + str(self.data) + " " + self._error if self.error else "") + + +class Case(object): + def __init__(self, schema): + self._schema = schema + + +class Or(object): + def __init__(self, *args): + self.args = args + + +class Optional(object): + def __init__(self, name): + self.name = name + + +class Directory(object): + def __init__(self, contains=None, contains_executable=None): + self.contains = contains + self.contains_executable = contains_executable + + def validate(self, name): + if not name: + yield Result(False, "is an empty string") + elif not os.path.exists(name): + yield Result(False, "Directory '{}' does not exist.".format(name)) + elif not os.path.isdir(name): + yield Result(False, "'{}' is not a directory.".format(name)) + else: + if self.contains: + for path in self.contains: + if not os.path.exists(os.path.join(name, path)): + yield Result(False, "'{}' does not contain '{}'".format(name, path)) + if self.contains_executable: + for program in self.contains_executable: + if not find_executable(program, name): + yield Result(False, "'{}' does not contain '{}'".format(name, program)) + + +class Schema(object): + def __init__(self, validator): + self.validator = validator + + def __call__(self, data): + for i in self.validate(data): + if not i.status: + print(i) + + def validate(self, data): + self.data = data + if isinstance(self.validator, string_types): + yield Result(isinstance(self.data, string_types), "is not a string", level=1, data=self.data) + elif issubclass(type(self.validator), type): + validator = self.validator + if self.validator == str: + validator = string_types + yield Result(isinstance(self.data, validator), + "is not {}".format(_get_type_name(self.validator)), level=1, data=self.data) + elif callable(self.validator): + if hasattr(self.validator, "expected_type"): + if not isinstance(data, self.validator.expected_type): + yield Result(False, "is not {}" + .format(_get_type_name(self.validator.expected_type)), level=1, data=self.data) + return + try: + self.validator(data) + yield Result(True, data=self.data) + except Exception as e: + yield Result(False, "didn't pass validation: {}".format(e), data=self.data) + elif isinstance(self.validator, dict): + if not len(self.validator): + yield Result(isinstance(self.data, dict), "is not a dictionary", level=1, data=self.data) + elif isinstance(self.validator, list): + if not isinstance(self.data, list): + yield Result(isinstance(self.data, list), "is not a list", level=1, data=self.data) + return + for i in self.iter(): + yield i + + def iter(self): + if isinstance(self.validator, dict): + if not isinstance(self.data, dict): + yield Result(False, "is not a dictionary.", level=1) + else: + for i in self.iter_dict(): + yield i + elif isinstance(self.validator, list): + if len(self.data) == 0: + yield Result(False, "is an empty list", data=self.data) + if len(self.validator) > 0: + for key, value in enumerate(self.data): + for v in Schema(self.validator[0]).validate(value): + yield Result(v.status, v.error, + path=(str(key) + ("." + v.path if v.path else "")), level=v.level, data=value) + elif isinstance(self.validator, Directory): + for v in self.validator.validate(self.data): + yield v + elif isinstance(self.validator, Or): + for i in self.iter_or(): + yield i + + def iter_dict(self): + for key in self.validator.keys(): + for d in self._data_key(key): + if d not in self.data and not isinstance(key, Optional): + yield Result(False, "is not defined.", path=d) + elif d not in self.data and isinstance(key, Optional): + continue + else: + validator = self.validator[key] + if isinstance(key, Or) and isinstance(self.validator[key], Case): + validator = self.validator[key]._schema[d] + for v in Schema(validator).validate(self.data[d]): + yield Result(v.status, v.error, + path=(d + ("." + v.path if v.path else "")), level=v.level, data=v.data) + + def iter_or(self): + results = [] + for a in self.validator.args: + r = [] + for v in Schema(a).validate(self.data): + r.append(v) + if any([x.status for x in r]) and not all([x.status for x in r]): + results += filter(lambda x: not x.status, r) + else: + results += r + if not any([x.status for x in results]): + max_level = 3 + for v in sorted(results, key=lambda x: x.level): + if v.level > max_level: + break + max_level = v.level + yield Result(v.status, v.error, path=v.path, level=v.level, data=v.data) + + def _data_key(self, key): + if isinstance(self.data, dict) and isinstance(key, str): + yield key + elif isinstance(key, Optional): + yield key.name + elif isinstance(key, Or): + if any([i in self.data for i in key.args]): + for i in key.args: + if i in self.data: + yield i + else: + for i in key.args: + yield i + + +def _get_type_name(python_type): + return {str: 'a string', int: 'and integer', float: 'a number', bool: 'a boolean', + list: 'an array', dict: 'a dictionary', string_types: "a string"}.get( + python_type, getattr(python_type, __name__, "unknown type")) + + +def assert_(condition, message="Wrong value"): + assert condition, message + + +userattributes = {"username": "", Optional("password"): ""} +available_dcs = [m.split(".")[-1] for m in dcs_modules()] +comma_separated_host_port.expected_type = string_types +validate_connect_address.expected_type = string_types +validate_host_port_listen.expected_type = string_types +validate_host_port_listen_multiple_hosts.expected_type = string_types +validate_data_dir.expected_type = string_types + +schema = Schema({ + "name": str, + "scope": str, + "restapi": { + "listen": validate_host_port_listen, + "connect_address": validate_connect_address + }, + Optional("bootstrap"): { + "dcs": { + Optional("ttl"): int, + Optional("loop_wait"): int, + Optional("retry_timeout"): int, + Optional("maximum_lag_on_failover"): int + }, + "pg_hba": [str], + "initdb": [Or(str, dict)] + }, + Or(*available_dcs): Case({ + "consul": { + Or("host", "url"): Case({ + "host": validate_host_port, + "url": str}) + }, + "etcd": { + Or("host", "hosts", "srv", "url", "proxy"): Case({ + "host": validate_host_port, + "hosts": Or(comma_separated_host_port, [validate_host_port]), + "srv": str, + "url": str, + "proxy": str}) + }, + "exhibitor": { + "hosts": [str], + "port": lambda i: assert_(int(i) <= 65535), + Optional("pool_interval"): int + }, + "zookeeper": { + "hosts": Or(comma_separated_host_port, [validate_host_port]), + }, + "kubernetes": { + "labels": {}, + Optional("namespace"): str, + Optional("scope_label"): str, + Optional("role_label"): str, + Optional("use_endpoints"): bool, + Optional("pod_ip"): Or(is_ipv4_address, is_ipv6_address), + Optional("ports"): [{"name": str, "port": int}], + }, + }), + "postgresql": { + "listen": validate_host_port_listen_multiple_hosts, + "connect_address": validate_connect_address, + "authentication": { + "replication": userattributes, + "superuser": userattributes, + "rewind": userattributes + }, + "data_dir": validate_data_dir, + Optional("bin_dir"): Directory(contains_executable=["pg_ctl", "initdb", "pg_controldata", "pg_basebackup", + "postgres", "pg_isready"]), + Optional("parameters"): { + Optional("unix_socket_directories"): lambda s: assert_(all([isinstance(s, string_types), len(s)])) + }, + Optional("pg_hba"): [str], + Optional("pg_ident"): [str], + Optional("pg_ctl_timeout"): int, + Optional("use_pg_rewind"): bool + }, + Optional("watchdog"): { + Optional("mode"): lambda m: assert_(m in ["off", "automatic", "required"]), + Optional("device"): str + }, + Optional("tags"): { + Optional("nofailover"): bool, + Optional("clonefrom"): bool, + Optional("noloadbalance"): bool, + Optional("replicatefrom"): str, + Optional("nosync"): bool + } +}) diff --git a/tests/test_validator.py b/tests/test_validator.py new file mode 100644 index 00000000..dbc3a2e4 --- /dev/null +++ b/tests/test_validator.py @@ -0,0 +1,221 @@ +import unittest +import os +import socket +import copy +from mock import Mock, patch, mock_open +from patroni.validator import schema +from six import StringIO + +config = { + "name": "string", + "scope": "string", + "restapi": { + "listen": "127.0.0.2:800", + "connect_address": "127.0.0.2:800" + }, + "bootstrap": { + "dcs": { + "ttl": 1000, + "loop_wait": 1000, + "retry_timeout": 1000, + "maximum_lag_on_failover": 1000 + }, + "pg_hba": ["string"], + "initdb": ["string", {"key":"value"}] + }, + "consul": { + "host": "127.0.0.1:5000" + }, + "etcd": { + "hosts": "127.0.0.1:2379,127.0.0.1:2380" + }, + "exhibitor": { + "hosts": ["string"], + "port": 4000, + "pool_interval": 1000 + }, + "zookeeper": { + "hosts": "127.0.0.1:3379,127.0.0.1:3380" + }, + "kubernetes": { + "namespace": "string", + "labels": {}, + "scope_label": "string", + "role_label": "string", + "use_endpoints": False, + "pod_ip": "127.0.0.1", + "ports": [{"name": "string", "port": 1000}], + }, + "postgresql": { + "listen": "127.0.0.2,::1:543", + "connect_address": "127.0.0.2:543", + "authentication": { + "replication": {"username": "user"}, + "superuser": {"username": "user"}, + "rewind": {"username": "user"}, + }, + "data_dir": "/tmp/data_dir", + "bin_dir": "/tmp/bin_dir", + "parameters": { + "unix_socket_directories": "." + }, + "pg_hba": [u"string"], + "pg_ident": ["string"], + "pg_ctl_timeout": 1000, + "use_pg_rewind": False + }, + "watchdog": { + "mode": "off", + "device": "string" + }, + "tags": { + "nofailover": False, + "clonefrom": False, + "noloadbalance": False, + "nosync": False + } +} + +directories = [] +files = [] + +def isfile_side_effect(arg): + return arg in files + + +def isdir_side_effect(arg): + return arg in directories + + +def exists_side_effect(arg): + return isfile_side_effect(arg) or isdir_side_effect(arg) + + +def connect_side_effect(host_port): + _, port = host_port + if port < 1000: + return 1 + elif port < 10000: + return 0 + else: + raise socket.gaierror() + + +def parse_output(output): + result = [] + for s in output.split("\n"): + x = s.split(" ")[0] + if x and x not in result: + result.append(x) + result.sort() + return result + + +@patch('socket.socket.connect_ex', Mock(side_effect=connect_side_effect)) +@patch('os.path.exists', Mock(side_effect=exists_side_effect)) +@patch('os.path.isdir', Mock(side_effect=isdir_side_effect)) +@patch('os.path.isfile', Mock(side_effect=isfile_side_effect)) +@patch('sys.stderr', new_callable=StringIO) +@patch('sys.stdout', new_callable=StringIO) +class TestValidator(unittest.TestCase): + + def setUp(self): + del files[:] + del directories[:] + + def test_empty_config(self, mock_out, mock_err): + schema({}) + output = mock_out.getvalue() + self.assertEqual(['consul', 'etcd', 'exhibitor', 'kubernetes', 'name', 'postgresql', 'restapi', 'scope', 'zookeeper'], parse_output(output)) + + def test_complete_config(self, mock_out, mock_err): + schema(config) + output = mock_out.getvalue() + self.assertEqual(['postgresql.bin_dir'], parse_output(output)) + + def test_bin_dir_is_file(self, mock_out, mock_err): + files.append(config["postgresql"]["data_dir"]) + files.append(config["postgresql"]["bin_dir"]) + c = copy.deepcopy(config) + c["restapi"]["connect_address"] = False + c["etcd"]["hosts"] = ["127.0.0.1:2379","1244.0.0.1:2379","127.0.0.1:invalidport"] + c["kubernetes"]["pod_ip"] = "127.0.0.1111" + schema(c) + output = mock_out.getvalue() + self.assertEqual(['etcd.hosts.1', 'etcd.hosts.2', 'kubernetes.pod_ip', 'postgresql.bin_dir', + 'postgresql.data_dir', 'restapi.connect_address'] , parse_output(output)) + + def test_bin_dir_is_empty(self, mock_out, mock_err): + directories.append(config["postgresql"]["data_dir"]) + directories.append(config["postgresql"]["bin_dir"]) + files.append(os.path.join(config["postgresql"]["data_dir"], "global", "pg_control")) + c = copy.deepcopy(config) + c["restapi"]["connect_address"] = "127.0.0.1" + c["kubernetes"]["pod_ip"] = "::1" + c["consul"]["host"] = "127.0.0.1:50000" + c["etcd"]["host"] = "127.0.0.1:237" + c["postgresql"]["listen"] = "127.0.0.1:5432" + with patch('patroni.validator.open', mock_open(read_data='9')): + schema(c) + output = mock_out.getvalue() + self.assertEqual(['consul.host', 'etcd.host', 'postgresql.bin_dir', 'postgresql.data_dir', + 'postgresql.listen', 'restapi.connect_address'], parse_output(output)) + + @patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1")) + def test_data_dir_contains_pg_version(self, mock_out, mock_err): + directories.append(config["postgresql"]["data_dir"]) + directories.append(config["postgresql"]["bin_dir"]) + directories.append(os.path.join(config["postgresql"]["data_dir"], "pg_wal")) + files.append(os.path.join(config["postgresql"]["data_dir"], "global", "pg_control")) + files.append(os.path.join(config["postgresql"]["data_dir"], "PG_VERSION")) + files.append(os.path.join(config["postgresql"]["bin_dir"], "pg_ctl")) + files.append(os.path.join(config["postgresql"]["bin_dir"], "initdb")) + files.append(os.path.join(config["postgresql"]["bin_dir"], "pg_controldata")) + files.append(os.path.join(config["postgresql"]["bin_dir"], "pg_basebackup")) + files.append(os.path.join(config["postgresql"]["bin_dir"], "postgres")) + files.append(os.path.join(config["postgresql"]["bin_dir"], "pg_isready")) + with patch('patroni.validator.open', mock_open(read_data='12')): + schema(config) + output = mock_out.getvalue() + self.assertEqual([], parse_output(output)) + + @patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1")) + def test_pg_version_missmatch(self, mock_out, mock_err): + directories.append(config["postgresql"]["data_dir"]) + directories.append(config["postgresql"]["bin_dir"]) + directories.append(os.path.join(config["postgresql"]["data_dir"], "pg_wal")) + files.append(os.path.join(config["postgresql"]["data_dir"], "global", "pg_control")) + files.append(os.path.join(config["postgresql"]["data_dir"], "PG_VERSION")) + c = copy.deepcopy(config) + c["etcd"]["hosts"] = [] + del c["postgresql"]["bin_dir"] + with patch('patroni.validator.open', mock_open(read_data='11')): + schema(c) + output = mock_out.getvalue() + self.assertEqual(['etcd.hosts', 'postgresql.data_dir'], parse_output(output)) + + @patch('subprocess.check_output', Mock(return_value=b"postgres (PostgreSQL) 12.1")) + def test_pg_wal_doesnt_exist(self, mock_out, mock_err): + directories.append(config["postgresql"]["data_dir"]) + directories.append(config["postgresql"]["bin_dir"]) + files.append(os.path.join(config["postgresql"]["data_dir"], "global", "pg_control")) + files.append(os.path.join(config["postgresql"]["data_dir"], "PG_VERSION")) + c = copy.deepcopy(config) + del c["postgresql"]["bin_dir"] + with patch('patroni.validator.open', mock_open(read_data='11')): + schema(c) + output = mock_out.getvalue() + self.assertEqual(['postgresql.data_dir'], parse_output(output)) + + + def test_data_dir_is_empty_string(self, mock_out, mock_err): + directories.append(config["postgresql"]["data_dir"]) + directories.append(config["postgresql"]["bin_dir"]) + c = copy.deepcopy(config) + c["kubernetes"] = False + c["postgresql"]["pg_hba"] = "" + c["postgresql"]["data_dir"] = "" + c["postgresql"]["bin_dir"] = "" + schema(c) + output = mock_out.getvalue() + self.assertEqual(['kubernetes', 'postgresql.bin_dir', 'postgresql.data_dir', 'postgresql.pg_hba'], parse_output(output)) From 80ce61876e90267c849715ededd59733ba8df6db Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 20 Feb 2020 10:07:43 +0100 Subject: [PATCH 02/24] Don't create permanent physical slot with name of the primary (#1392) It is a regular issue that primary is recycling WALs when one of the replicas is down for a long time. So far there were only two solutions for such a problem and both of them are not perfect: 1. Increase `wal_keep_segments`, but it is hard to guess the good value. 2. Use continuous archiving and PITR, but it is not always possible. This PR is introducing the way to solve the problem for static clusters, with a fixed number of nodes and names that never change. You just need to list the names of all nodes in the `slots` so the primary will not remove the slot when the node is down (not registered in DCS). Of course, the primary will not create the permanent slot which is matching its own name. Usage example: let's assume you have a cluster with nodes named *abc1*, *abc2*, and *abc3*. You have to run `patronictl edit-config` and put the following snippet into the configuration: ```yaml slots: abc1: type: physical abc2: type: physical abc3: type: physical ``` If the node *abc2* is the primary, it will always create slots for *abc1* and *abc3* even if they are not running, but will not create slot *abc2*. Other nodes will behave the same. Close #280 --- docs/SETTINGS.rst | 4 ++-- patroni/dcs/__init__.py | 31 +++++++++++++++---------------- tests/test_postgresql.py | 4 ++-- 3 files changed, 19 insertions(+), 20 deletions(-) diff --git a/docs/SETTINGS.rst b/docs/SETTINGS.rst index c099d2e8..c1198a55 100644 --- a/docs/SETTINGS.rst +++ b/docs/SETTINGS.rst @@ -20,7 +20,7 @@ Dynamic configuration is stored in the DCS (Distributed Configuration Store) and - **synchronous\_mode\_strict**: prevents disabling synchronous replication if no synchronous replicas are available, blocking all client writes to the master. See :ref:`replication modes documentation ` for details. - **postgresql**: - **use\_pg\_rewind**: whether or not to use pg_rewind. Defaults to `false`. - - **use\_slots**: whether or not to use replication_slots. Defaults to `true` on PostgreSQL 9.4+. + - **use\_slots**: whether or not to use replication slots. Defaults to `true` on PostgreSQL 9.4+. - **recovery\_conf**: additional configuration settings written to recovery.conf when configuring follower. There is no recovery.conf anymore in PostgreSQL 12, but you may continue using this section, because Patroni handles it transparently. - **parameters**: list of configuration settings for Postgres. - **standby\_cluster**: if this section is defined, we want to bootstrap a standby cluster. @@ -32,7 +32,7 @@ Dynamic configuration is stored in the DCS (Distributed Configuration Store) and - **archive\_cleanup\_command**: cleanup command for standby leader - **recovery\_min\_apply\_delay**: how long to wait before actually apply WAL records on a standby leader - **slots**: define permanent replication slots. These slots will be preserved during switchover/failover. Patroni will try to create slots before opening connections to the cluster. - - **my_slot_name**: the name of replication slot. It is the responsibility of the operator to make sure that there are no clashes in names between replication slots automatically created by Patroni for members and permanent replication slots. + - **my_slot_name**: the name of replication slot. If the permanent slot name matches with the name of the current primary it will not be created. Everything else is the responsibility of the operator to make sure that there are no clashes in names between replication slots automatically created by Patroni for members and permanent replication slots. - **type**: slot type. Could be ``physical`` or ``logical``. If the slot is logical, you have to additionally define ``database`` and ``plugin``. - **database**: the database name where logical slots should be created. - **plugin**: the plugin name for the logical slot. diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index 4768f894..7faa5429 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -449,21 +449,21 @@ class Cluster(namedtuple('Cluster', 'initialize,config,leader,last_leader_operat def is_synchronous_mode(self): return self.check_mode('synchronous_mode') - def get_replication_slots(self, name, role): + def get_replication_slots(self, my_name, role): # if the replicatefrom tag is set on the member - we should not create the replication slot for it on # the current master, because that member would replicate from elsewhere. We still create the slot if # 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 use_slots = self.config and self.config.data.get('postgresql', {}).get('use_slots', True) if role in ('master', 'standby_leader'): - slot_members = [m.name for m in self.members if use_slots and m.name != name and - (m.replicatefrom is None or m.replicatefrom == name or + slot_members = [m.name for m in self.members if use_slots and m.name != my_name and + (m.replicatefrom is None or m.replicatefrom == my_name or not self.has_member(m.replicatefrom))] permanent_slots = (self.config and self.config.permanent_slots or {}).copy() else: # only manage slots for replicas that replicate from this one, except for the leader among them slot_members = [m.name for m in self.members if use_slots and - m.replicatefrom == name and m.name != self.leader.name] + m.replicatefrom == my_name and m.name != self.leader.name] permanent_slots = {} slots = {slot_name_from_member_name(name): {'type': 'physical'} for name in slot_members} @@ -484,22 +484,21 @@ class Cluster(namedtuple('Cluster', 'initialize,config,leader,last_leader_operat logger.error("Slot name may only contain lower case letters, numbers, and the underscore chars") continue - if name in slots: - logger.error("Permanent replication slot {'%s': %s} is conflicting with" + - " physical replication slot for cluster member", name, value) - continue - - value = deepcopy(value) - if not value: - value = {'type': 'physical'} - + value = deepcopy(value) if value else {'type': 'physical'} if isinstance(value, dict): if 'type' not in value: value['type'] = 'logical' if value.get('database') and value.get('plugin') else 'physical' - if value['type'] == 'physical' or value['type'] == 'logical' \ - and value.get('database') and value.get('plugin'): - slots[name] = value + if value['type'] == 'physical': + if name != my_name: # Don't try to create permanent physical replication slot for yourself + slots[name] = value + continue + elif value['type'] == 'logical' and value.get('database') and value.get('plugin'): + if name in slots: + logger.error("Permanent logical replication slot {'%s': %s} is conflicting with" + + " physical replication slot for cluster member", name, value) + else: + slots[name] = value continue logger.error("Bad value for slot '%s' in permanent_slots: %s", name, permanent_slots[name]) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index a63c572e..00adb8f8 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -281,8 +281,8 @@ class TestPostgresql(BaseTestPostgresql): @patch.object(Postgresql, 'is_running', Mock(return_value=True)) def test_sync_replication_slots(self): self.p.start() - config = ClusterConfig(1, {'slots': {'ls': {'database': 'a', 'plugin': 'b'}, - 'A': 0, 'test_3': 0, 'b': {'type': 'logical', 'plugin': '1'}}}, 1) + config = ClusterConfig(1, {'slots': {'test_3': {'database': 'a', 'plugin': 'b'}, + 'A': 0, 'ls': 0, 'b': {'type': 'logical', 'plugin': '1'}}}, 1) cluster = Cluster(True, config, self.leader, 0, [self.me, self.other, self.leadermem], None, None, None) with mock.patch('patroni.postgresql.Postgresql._query', Mock(side_effect=psycopg2.OperationalError)): self.p.slots_handler.sync_replication_slots(cluster) From 7b0e012f6220fd007069f7648d70e2aa35a3aaae Mon Sep 17 00:00:00 2001 From: Julien Riou Date: Thu, 20 Feb 2020 10:13:41 +0100 Subject: [PATCH 03/24] Disable SSL verification for Consul when it is required (#1399) Consul client uses urllib3 with a verify=True by default. When SSL verification is disabled with verify=False, we can see CERTIFICATE_VERIFY_FAILED exceptions. With urllib3 1.19.1-1 on Debian Stretch, the "cert_reqs" argument must be explicitaly set to ssl.CERT_NONE to effectively disable SSL verification. --- patroni/dcs/consul.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index cea7bb88..26f7b656 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -55,6 +55,8 @@ class HTTPClient(object): kwargs['ca_certs'] = ca_cert if verify or ca_cert: kwargs['cert_reqs'] = ssl.CERT_REQUIRED + else: + kwargs['cert_reqs'] = ssl.CERT_NONE self.http = urllib3.PoolManager(num_pools=10, **kwargs) self._ttl = None From e759a3f2ef483d41a7c86710f15291a980180ec4 Mon Sep 17 00:00:00 2001 From: damien clochard Date: Thu, 20 Feb 2020 10:14:36 +0100 Subject: [PATCH 04/24] [doc] add PATRONICTL_CONFIG_FILE env var (#1397) --- docs/ENVIRONMENT.rst | 1 + 1 file changed, 1 insertion(+) diff --git a/docs/ENVIRONMENT.rst b/docs/ENVIRONMENT.rst index 81203957..f95d125b 100644 --- a/docs/ENVIRONMENT.rst +++ b/docs/ENVIRONMENT.rst @@ -130,6 +130,7 @@ REST API CTL --- +- **PATRONICTL\_CONFIG\_FILE**: location of the configuration file. - **PATRONI\_CTL\_INSECURE**: Allow connections to REST API without verifying SSL certs. - **PATRONI\_CTL\_CACERT**: Specifies the file with the CA_BUNDLE file or directory with certificates of trusted CAs to use while verifying REST API SSL certs. If not provided patronictl will use the value provided for REST API "cafile" parameter. - **PATRONI\_CTL\_CERTFILE**: Specifies the file with the client certificate in the PEM format. If not provided patronictl will use the value provided for REST API "certfile" parameter. From 0fa70e8d881ab6cfe94797df0b6a1bcfc81ef739 Mon Sep 17 00:00:00 2001 From: Steven De Coeyer Date: Thu, 20 Feb 2020 10:15:58 +0100 Subject: [PATCH 05/24] Updates README (#1394) We need to ensure to enable etcd v2, cfr. https://github.com/zalando/patroni/issues/1270 and https://github.com/zalando/patroni/issues/1163. --- docs/README.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/README.rst b/docs/README.rst index 6a86d8f2..2b57eb1e 100644 --- a/docs/README.rst +++ b/docs/README.rst @@ -107,7 +107,7 @@ obtain those files from the git repository and replace `./patroni.py` below with To get started, do the following from different terminals: :: - > etcd --data-dir=data/etcd + > etcd --data-dir=data/etcd --enable-v2=true > ./patroni.py postgres0.yml > ./patroni.py postgres1.yml From bcd75bbeeb0ff7f8d98415c0cfac2b5245e3ce2d Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 27 Feb 2020 12:22:11 +0100 Subject: [PATCH 06/24] Avoid opening replication connection on every cycle of HA loop (#1422) Bug was introduces in the https://github.com/zalando/patroni/pull/1332 --- patroni/ha.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/patroni/ha.py b/patroni/ha.py index fdc2ce07..d38be332 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -204,7 +204,7 @@ class Ha(object): if self.state_handler.role == 'standby_leader': timeline = pg_control_timeline or self.state_handler.pg_control_timeline() else: - timeline = self.state_handler.replica_cached_timeline(timeline) + timeline = self.state_handler.replica_cached_timeline(self._leader_timeline) if timeline: data['timeline'] = timeline except Exception: From 4a29caa9d38de77aa28c5843d9c2614353fa1f75 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 27 Feb 2020 12:22:44 +0100 Subject: [PATCH 07/24] On role change callback didn't fire on failed primary (#1420) Bug was introduced in https://github.com/zalando/patroni/pull/703 Close https://github.com/zalando/patroni/issues/1418 --- patroni/ha.py | 10 ++++++---- patroni/postgresql/__init__.py | 6 +++++- tests/test_ha.py | 1 + tests/test_postgresql.py | 3 +++ 4 files changed, 15 insertions(+), 5 deletions(-) diff --git a/patroni/ha.py b/patroni/ha.py index d38be332..02025fec 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -1143,7 +1143,8 @@ class Ha(object): if not self.state_handler.is_running(): self.watchdog.disable() if self.has_lock(): - self.state_handler.set_role('demoted') + if self.state_handler.role in ('master', 'standby_leader'): + self.state_handler.set_role('demoted') self._delete_leader() return 'removed leader key after trying and failing to start postgres' return 'failed to start postgres' @@ -1172,10 +1173,11 @@ class Ha(object): return ret or 'running post_bootstrap' self.state_handler.bootstrapping = False - self.dcs.set_config_value(json.dumps(self.patroni.config.dynamic_configuration, separators=(',', ':'))) if not self.watchdog.activate(): logger.error('Cancelling bootstrap because watchdog activation failed') self.cancel_initialization() + self.dcs.initialize(create_new=(self.cluster.initialize is None), sysid=self.state_handler.sysid) + self.dcs.set_config_value(json.dumps(self.patroni.config.dynamic_configuration, separators=(',', ':'))) self.state_handler.slots_handler.sync_replication_slots(self.cluster) self.dcs.take_leader() self.set_is_leader(True) @@ -1290,8 +1292,8 @@ class Ha(object): data_sysid = self.state_handler.sysid if not self.sysid_valid(data_sysid): # data directory is not empty, but no valid sysid, cluster must be broken, suggest reinit - return ("data dir for the cluster is not empty, but system ID is invalid; consider doing" - "reinitialize") + return ("data dir for the cluster is not empty, " + "but system ID is invalid; consider doing reinitialize") if self.sysid_valid(self.cluster.initialize): if self.cluster.initialize != data_sysid: diff --git a/patroni/postgresql/__init__.py b/patroni/postgresql/__init__.py index d4b262d6..d31b77d7 100644 --- a/patroni/postgresql/__init__.py +++ b/patroni/postgresql/__init__.py @@ -413,7 +413,11 @@ class Postgresql(object): self.set_state('starting') self._pending_restart = False - configuration = self.config.effective_configuration + try: + configuration = self.config.effective_configuration + except Exception: + return None + self.config.check_directories() self.config.write_postgresql_conf(configuration) self.config.resolve_connection_addresses() diff --git a/tests/test_ha.py b/tests/test_ha.py index 2ade4526..b054cfc2 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -621,6 +621,7 @@ class TestHa(PostgresInit): def test_post_recover(self): self.p.is_running = false self.ha.has_lock = true + self.p.set_role('master') self.assertEqual(self.ha.post_recover(), 'removed leader key after trying and failing to start postgres') self.ha.has_lock = false self.assertEqual(self.ha.post_recover(), 'failed to start postgres') diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 00adb8f8..8a81ddc5 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -134,6 +134,9 @@ class TestPostgresql(BaseTestPostgresql): self.p.cancellable.cancel() self.assertFalse(self.p.start()) + with patch('patroni.postgresql.config.ConfigHandler.effective_configuration', + PropertyMock(side_effect=Exception)): + self.assertIsNone(self.p.start()) @patch.object(Postgresql, 'pg_isready') @patch('patroni.postgresql.polling_loop', Mock(return_value=range(1))) From 613634c26bc7d82e0e71acb889bcd0dc9529ea07 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 27 Feb 2020 12:24:17 +0100 Subject: [PATCH 08/24] Reset rewind state if postgres started after successful pg_rewind (#1408) Close https://github.com/zalando/patroni/issues/1406 --- patroni/ha.py | 2 ++ tests/test_ha.py | 4 ++++ 2 files changed, 6 insertions(+) diff --git a/patroni/ha.py b/patroni/ha.py index 02025fec..ceecb388 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -1149,6 +1149,8 @@ class Ha(object): return 'removed leader key after trying and failing to start postgres' return 'failed to start postgres' self._crash_recovery_executed = False + if self._rewind.executed and not self._rewind.failed: + self._rewind.reset_state() return None def cancel_initialization(self): diff --git a/tests/test_ha.py b/tests/test_ha.py index b054cfc2..1840c901 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -618,6 +618,8 @@ class TestHa(PostgresInit): member = Member(0, 'test', 1, {'api_url': 'http://localhost:8011/patroni'}) self.ha.fetch_node_status(member) + @patch.object(Rewind, 'pg_rewind', true) + @patch.object(Rewind, 'check_leader_is_not_in_recovery', true) def test_post_recover(self): self.p.is_running = false self.ha.has_lock = true @@ -625,6 +627,8 @@ class TestHa(PostgresInit): self.assertEqual(self.ha.post_recover(), 'removed leader key after trying and failing to start postgres') self.ha.has_lock = false self.assertEqual(self.ha.post_recover(), 'failed to start postgres') + leader = Leader(0, 0, Member(0, 'l', 2, {"version": "1.6", "conn_url": "postgres://a", "role": "master"})) + self.ha._rewind.execute(leader) self.p.is_running = true self.assertIsNone(self.ha.post_recover()) From ab38ab2e97e0b7d758398fc45b58e413435b569b Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 10 Mar 2020 12:07:26 +0100 Subject: [PATCH 09/24] Apply 1 second backoff if LIST failed (#1424) It is mostly necessary to avoid flooding logs, but also help to prevent starvation of the main thread. --- patroni/dcs/kubernetes.py | 6 +++++- tests/test_kubernetes.py | 5 +++++ 2 files changed, 10 insertions(+), 1 deletion(-) diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index d84f02c4..6c388102 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -106,7 +106,11 @@ class ObjectCache(Thread): self.start() def _list(self): - return self._func(_request_timeout=(self._retry.deadline, Timeout.DEFAULT_TIMEOUT)) + try: + return self._func(_request_timeout=(self._retry.deadline, Timeout.DEFAULT_TIMEOUT)) + except Exception: + time.sleep(1) + raise def _watch(self, resource_version): return self._func(_request_timeout=(self._retry.deadline, Timeout.DEFAULT_TIMEOUT), diff --git a/tests/test_kubernetes.py b/tests/test_kubernetes.py index 03b7df40..c853fc4a 100644 --- a/tests/test_kubernetes.py +++ b/tests/test_kubernetes.py @@ -174,3 +174,8 @@ class TestCacheBuilder(unittest.TestCase): @patch('patroni.dcs.kubernetes.ObjectCache._build_cache', Mock(side_effect=Exception)) def test_run(self): self.assertRaises(SleepException, self.k._pods.run) + + @patch('time.sleep', Mock()) + def test__list(self): + self.k._pods._func = Mock(side_effect=Exception) + self.assertRaises(Exception, self.k._pods._list) From b0208744862ec24d7f82fb0d78409bba40231562 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 10 Mar 2020 12:07:40 +0100 Subject: [PATCH 10/24] Small improvement in tests (#1423) which actually revealed a small issue in the validator --- patroni/dcs/consul.py | 5 +---- patroni/validator.py | 8 ++++---- tests/test_validator.py | 25 +++++++++++++++---------- 3 files changed, 20 insertions(+), 18 deletions(-) diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index 26f7b656..efa81ea1 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -53,10 +53,7 @@ class HTTPClient(object): kwargs['cert_file'] = cert if ca_cert: kwargs['ca_certs'] = ca_cert - if verify or ca_cert: - kwargs['cert_reqs'] = ssl.CERT_REQUIRED - else: - kwargs['cert_reqs'] = ssl.CERT_NONE + kwargs['cert_reqs'] = ssl.CERT_REQUIRED if verify or ca_cert else ssl.CERT_NONE self.http = urllib3.PoolManager(num_pools=10, **kwargs) self._ttl = None diff --git a/patroni/validator.py b/patroni/validator.py index d254f0b2..44ec1f2a 100644 --- a/patroni/validator.py +++ b/patroni/validator.py @@ -19,11 +19,11 @@ def data_directory_empty(data_dir): def validate_connect_address(address): try: - host, _ = split_host_port(address, None) - except (ValueError, TypeError): + host, _ = split_host_port(address, 1) + except (AttributeError, TypeError, ValueError): raise ConfigParseError("contains a wrong value") - if host in ["127.0.0.1", "0.0.0.0", "*", "::1"]: - raise ConfigParseError('must not contain "127.0.0.1", "0.0.0.0", "*", "::1"') + if host in ["127.0.0.1", "0.0.0.0", "*", "::1", "localhost"]: + raise ConfigParseError('must not contain "127.0.0.1", "0.0.0.0", "*", "::1", "localhost"') return True diff --git a/tests/test_validator.py b/tests/test_validator.py index dbc3a2e4..b562464e 100644 --- a/tests/test_validator.py +++ b/tests/test_validator.py @@ -1,11 +1,14 @@ -import unittest +import copy import os import socket -import copy +import unittest + from mock import Mock, patch, mock_open +from patroni.dcs import dcs_modules from patroni.validator import schema from six import StringIO +available_dcs = [m.split(".")[-1] for m in dcs_modules()] config = { "name": "string", "scope": "string", @@ -21,7 +24,7 @@ config = { "maximum_lag_on_failover": 1000 }, "pg_hba": ["string"], - "initdb": ["string", {"key":"value"}] + "initdb": ["string", {"key": "value"}] }, "consul": { "host": "127.0.0.1:5000" @@ -79,6 +82,7 @@ config = { directories = [] files = [] + def isfile_side_effect(arg): return arg in files @@ -126,7 +130,8 @@ class TestValidator(unittest.TestCase): def test_empty_config(self, mock_out, mock_err): schema({}) output = mock_out.getvalue() - self.assertEqual(['consul', 'etcd', 'exhibitor', 'kubernetes', 'name', 'postgresql', 'restapi', 'scope', 'zookeeper'], parse_output(output)) + expected = list(sorted(['name', 'postgresql', 'restapi', 'scope'] + available_dcs)) + self.assertEqual(expected, parse_output(output)) def test_complete_config(self, mock_out, mock_err): schema(config) @@ -137,20 +142,20 @@ class TestValidator(unittest.TestCase): files.append(config["postgresql"]["data_dir"]) files.append(config["postgresql"]["bin_dir"]) c = copy.deepcopy(config) - c["restapi"]["connect_address"] = False - c["etcd"]["hosts"] = ["127.0.0.1:2379","1244.0.0.1:2379","127.0.0.1:invalidport"] + c["restapi"]["connect_address"] = 'False:blabla' + c["etcd"]["hosts"] = ["127.0.0.1:2379", "1244.0.0.1:2379", "127.0.0.1:invalidport"] c["kubernetes"]["pod_ip"] = "127.0.0.1111" schema(c) output = mock_out.getvalue() self.assertEqual(['etcd.hosts.1', 'etcd.hosts.2', 'kubernetes.pod_ip', 'postgresql.bin_dir', - 'postgresql.data_dir', 'restapi.connect_address'] , parse_output(output)) + 'postgresql.data_dir', 'restapi.connect_address'], parse_output(output)) def test_bin_dir_is_empty(self, mock_out, mock_err): directories.append(config["postgresql"]["data_dir"]) directories.append(config["postgresql"]["bin_dir"]) files.append(os.path.join(config["postgresql"]["data_dir"], "global", "pg_control")) c = copy.deepcopy(config) - c["restapi"]["connect_address"] = "127.0.0.1" + c["restapi"]["connect_address"] = "127.0.0.1:8008" c["kubernetes"]["pod_ip"] = "::1" c["consul"]["host"] = "127.0.0.1:50000" c["etcd"]["host"] = "127.0.0.1:237" @@ -207,7 +212,6 @@ class TestValidator(unittest.TestCase): output = mock_out.getvalue() self.assertEqual(['postgresql.data_dir'], parse_output(output)) - def test_data_dir_is_empty_string(self, mock_out, mock_err): directories.append(config["postgresql"]["data_dir"]) directories.append(config["postgresql"]["bin_dir"]) @@ -218,4 +222,5 @@ class TestValidator(unittest.TestCase): c["postgresql"]["bin_dir"] = "" schema(c) output = mock_out.getvalue() - self.assertEqual(['kubernetes', 'postgresql.bin_dir', 'postgresql.data_dir', 'postgresql.pg_hba'], parse_output(output)) + self.assertEqual(['kubernetes', 'postgresql.bin_dir', + 'postgresql.data_dir', 'postgresql.pg_hba'], parse_output(output)) From 795efc4548e22a5f44603c996a8c373a4d778fd2 Mon Sep 17 00:00:00 2001 From: Michail Nikolaev Date: Tue, 10 Mar 2020 14:08:01 +0300 Subject: [PATCH 11/24] Note about possible data loss while canceling postgres backends. (#1414) Note about possible data loss while canceling postgres backends. Related to zalando#1412 --- docs/replication_modes.rst | 2 ++ 1 file changed, 2 insertions(+) diff --git a/docs/replication_modes.rst b/docs/replication_modes.rst index b6dfe1d0..9446c5a2 100644 --- a/docs/replication_modes.rst +++ b/docs/replication_modes.rst @@ -57,6 +57,8 @@ You can ensure that a standby never becomes the synchronous standby by setting ` Synchronous mode can be switched on and off via Patroni REST interface. See :ref:`dynamic configuration ` for instructions. +Note: Because of the way synchronous replication is implemented in PostgreSQL it is still possible to lose transactions even when using ``synchronous_mode_strict``. If the PostgreSQL backend is cancelled while waiting to acknowledge replication (as a result of packet cancellation due to client timeout or backend failure) transaction changes become visible for other backends. Such changes are not yet replicated and may be lost in case of standby promotion. + Synchronous mode implementation ------------------------------- From d74a4b23a63f9aa63608fc34ed2c95d61c450eb9 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Tue, 10 Mar 2020 12:08:29 +0100 Subject: [PATCH 12/24] Scrub KUBERNETES_ environment from the postmaster (#1407) The KUBERNETES_ environment variables are not required for PostgreSQL, yet having them exposed to the postmaster will also expose them to backends and to regular database users (using pl/perl for example). --- patroni/__init__.py | 1 + patroni/postgresql/postmaster.py | 5 +++-- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/patroni/__init__.py b/patroni/__init__.py index ba7a803e..8a8899b1 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -9,6 +9,7 @@ from patroni.version import __version__ logger = logging.getLogger(__name__) PATRONI_ENV_PREFIX = 'PATRONI_' +KUBERNETES_ENV_PREFIX = 'KUBERNETES_' class Patroni(object): diff --git a/patroni/postgresql/postmaster.py b/patroni/postgresql/postmaster.py index 6ee51d26..b95ead33 100644 --- a/patroni/postgresql/postmaster.py +++ b/patroni/postgresql/postmaster.py @@ -7,7 +7,7 @@ import signal import subprocess import sys -from patroni import PATRONI_ENV_PREFIX +from patroni import PATRONI_ENV_PREFIX, KUBERNETES_ENV_PREFIX # avoid spawning the resource tracker process if sys.version_info >= (3, 8): # pragma: no cover @@ -176,7 +176,8 @@ class PostmasterProcess(psutil.Process): # In order to make everything portable we can't use fork&exec approach here, so we will call # ourselves and pass list of arguments which must be used to start postgres. # On Windows, in order to run a side-by-side assembly the specified env must include a valid SYSTEMROOT. - env = {p: os.environ[p] for p in os.environ if not p.startswith(PATRONI_ENV_PREFIX)} + env = {p: os.environ[p] for p in os.environ if not p.startswith( + PATRONI_ENV_PREFIX) and not p.startswith(KUBERNETES_ENV_PREFIX)} try: proc = PostmasterProcess._from_pidfile(data_dir) if proc and not proc._is_postmaster_process(): From d82301688fa5ccafd2122330e7c979d68c0ca40d Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Fri, 13 Mar 2020 16:36:12 +0100 Subject: [PATCH 13/24] Fix pyinstaller compatibility (#1441) Close https://github.com/zalando/patroni/issues/1440 --- patroni/__init__.py | 4 ++++ patroni/dcs/__init__.py | 7 +++++-- 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/patroni/__init__.py b/patroni/__init__.py index 8a8899b1..c079d61b 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -169,9 +169,13 @@ class Patroni(object): def patroni_main(): import argparse + + from multiprocessing import freeze_support from patroni.config import Config, ConfigParseError from patroni.validator import schema + freeze_support() + parser = argparse.ArgumentParser() parser.add_argument('--version', action='version', version='%(prog)s {0}'.format(__version__)) parser.add_argument('--validate-config', action='store_true', help='Run config validator and exit') diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index 7faa5429..cee96805 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -66,8 +66,11 @@ def dcs_modules(): module_prefix = __package__ + '.' if getattr(sys, 'frozen', False): - importer = pkgutil.get_importer(dcs_dirname) - return [module for module in list(importer.toc) if module.startswith(module_prefix) and module.count('.') == 2] + toc = set() + for importer in pkgutil.iter_importers(dcs_dirname): + if hasattr(importer, 'toc'): + toc |= importer.toc + return [module for module in toc if module.startswith(module_prefix) and module.count('.') == 2] else: return [module_prefix + name for _, name, is_pkg in pkgutil.iter_modules([dcs_dirname]) if not is_pkg] From d2080a3116a4941f5138b510bc47aa1f2446f790 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Fri, 13 Mar 2020 16:36:39 +0100 Subject: [PATCH 14/24] Retry if the retry-after http header is set (#1431) If the K8s API is overwhelmed with requests it might ask to retry. --- patroni/dcs/kubernetes.py | 9 ++++++++- patroni/utils.py | 2 +- 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index 6c388102..9244ca04 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -31,6 +31,13 @@ class KubernetesRetriableException(k8s_client.rest.ApiException): self.body = orig.body self.headers = orig.headers + @property + def sleeptime(self): + try: + return int(self.headers['retry-after']) + except Exception: + return None + class CoreV1ApiProxy(object): @@ -68,7 +75,7 @@ class CoreV1ApiProxy(object): try: return getattr(self._api, func)(*args, **kwargs) except k8s_client.rest.ApiException as e: - if e.status in (502, 503, 504): # XXX + if e.status in (502, 503, 504) or e.headers and 'retry-after' in e.headers: # XXX raise KubernetesRetriableException(e) raise return wrapper diff --git a/patroni/utils.py b/patroni/utils.py index 2d35963e..6ac3a5a9 100644 --- a/patroni/utils.py +++ b/patroni/utils.py @@ -334,7 +334,7 @@ class Retry(object): logger.warning('Retry got exception: %s', e) raise RetryFailedError("Too many retry attempts") self._attempts += 1 - sleeptime = self.sleeptime + sleeptime = hasattr(e, 'sleeptime') and e.sleeptime or self.sleeptime if self._cur_stoptime is not None and time.time() + sleeptime >= self._cur_stoptime: logger.warning('Retry got exception: %s', e) From 810c179592cedc5466d714fe0d54a221e69a762e Mon Sep 17 00:00:00 2001 From: Danyal Prout <672580+danyalprout@users.noreply.github.com> Date: Wed, 18 Mar 2020 06:04:09 -0400 Subject: [PATCH 15/24] Kazoo 2.7.0 Compatibility Close https://github.com/zalando/patroni/issues/1448 --- patroni/dcs/zookeeper.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/patroni/dcs/zookeeper.py b/patroni/dcs/zookeeper.py index 33433d7a..8e9f5637 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -76,7 +76,7 @@ class ZooKeeper(AbstractDCS): self._client.start() - def _kazoo_connect(self, host, port): + def _kazoo_connect(self, *args): """Kazoo is using Ping's to determine health of connection to zookeeper. If there is no response on Ping after Ping interval (1/2 from read_timeout) it will consider current connection dead and try to connect to another node. Without this "magic" it was taking @@ -88,7 +88,7 @@ class ZooKeeper(AbstractDCS): than loop_wait, because we can spend up to 2 seconds when calling `touch_member()` and `write_leader_optime()` methods, which also may hang...""" - ret = self._orig_kazoo_connect(host, port) + ret = self._orig_kazoo_connect(*args) return max(self.loop_wait - 2, 2)*1000, ret[1] def session_listener(self, state): From 0e4d7f01f29c10831b3612db2fd3a40e22bc20db Mon Sep 17 00:00:00 2001 From: Casey Allen Shobe Date: Wed, 1 Apr 2020 07:50:50 -0600 Subject: [PATCH 16/24] Correct documentation for consul.host (#1438) Close #1434 --- docs/ENVIRONMENT.rst | 4 ++-- docs/SETTINGS.rst | 24 ++++++++++++------------ 2 files changed, 14 insertions(+), 14 deletions(-) diff --git a/docs/ENVIRONMENT.rst b/docs/ENVIRONMENT.rst index f95d125b..987b8b34 100644 --- a/docs/ENVIRONMENT.rst +++ b/docs/ENVIRONMENT.rst @@ -35,8 +35,8 @@ Example: defining ``PATRONI_admin_PASSWORD=strongpasswd`` and ``PATRONI_admin_OP Consul ------ -- **PATRONI\_CONSUL\_HOST**: the host:port for the Consul endpoint. -- **PATRONI\_CONSUL\_URL**: url for the Consul, in format: http(s)://host:port +- **PATRONI\_CONSUL\_HOST**: the host:port for the Consul local agent. +- **PATRONI\_CONSUL\_URL**: url for the Consul local agent, in format: http(s)://host:port - **PATRONI\_CONSUL\_PORT**: (optional) Consul port - **PATRONI\_CONSUL\_SCHEME**: (optional) **http** or **https**, defaults to **http** - **PATRONI\_CONSUL\_TOKEN**: (optional) ACL token diff --git a/docs/SETTINGS.rst b/docs/SETTINGS.rst index c1198a55..1dca57b7 100644 --- a/docs/SETTINGS.rst +++ b/docs/SETTINGS.rst @@ -87,20 +87,20 @@ Consul ------ Most of the parameters are optional, but you have to specify one of the **host** or **url** -- **host**: the host:port for the Consul endpoint, in format: http(s)://host:port -- **url**: url for the Consul endpoint -- **port**: (optional) Consul port -- **scheme**: (optional) **http** or **https**, defaults to **http** -- **token**: (optional) ACL token -- **verify**: (optional) whether to verify the SSL certificate for HTTPS requests +- **host**: the host:port for the Consul local agent. +- **url**: url for the Consul local agent, in format: http(s)://host:port. +- **port**: (optional) Consul port. +- **scheme**: (optional) **http** or **https**, defaults to **http**. +- **token**: (optional) ACL token. +- **verify**: (optional) whether to verify the SSL certificate for HTTPS requests. - **cacert**: (optional) The ca certificate. If present it will enable validation. -- **cert**: (optional) file with the client certificate +- **cert**: (optional) file with the client certificate. - **key**: (optional) file with the client key. Can be empty if the key is part of **cert**. - **dc**: (optional) Datacenter to communicate with. By default the datacenter of the host is used. - **consistency**: (optional) Select consul consistency mode. Possible values are ``default``, ``consistent``, or ``stale`` (more details in `consul API reference `__) - **checks**: (optional) list of Consul health checks used for the session. By default an empty list is used. -- **register\_service**: (optional) whether or not to register a service with the name defined by the scope parameter and the tag master, replica or standby-leader depending on the node's role. Defaults to **false** -- **service\_check\_interval**: (optional) how often to perform health check against registered url +- **register\_service**: (optional) whether or not to register a service with the name defined by the scope parameter and the tag master, replica or standby-leader depending on the node's role. Defaults to **false**. +- **service\_check\_interval**: (optional) how often to perform health check against registered url. Etcd ---- @@ -109,8 +109,8 @@ Most of the parameters are optional, but you have to specify one of the **host** - **host**: the host:port for the etcd endpoint. - **hosts**: list of etcd endpoint in format host1:port1,host2:port2,etc... Could be a comma separated string or an actual yaml list. - **use\_proxies**: If this parameter is set to true, Patroni will consider **hosts** as a list of proxies and will not perform a topology discovery of etcd cluster. -- **url**: url for the etcd -- **proxy**: proxy url for the etcd. If you are connecting to the etcd using proxy, use this parameter instead of **url** +- **url**: url for the etcd. +- **proxy**: proxy url for the etcd. If you are connecting to the etcd using proxy, use this parameter instead of **url**. - **srv**: Domain to search the SRV record(s) for cluster autodiscovery. - **protocol**: (optional) http or https, if not specified http is used. If the **url** or **proxy** is specified - will take protocol from them. - **username**: (optional) username for etcd authentication. @@ -126,7 +126,7 @@ ZooKeeper Exhibitor --------- - **hosts**: initial list of Exhibitor (ZooKeeper) nodes in format: 'host1,host2,etc...'. This list updates automatically whenever the Exhibitor (ZooKeeper) cluster topology changes. -- **poll\_interval**: how often the list of ZooKeeper and Exhibitor nodes should be updated from Exhibitor +- **poll\_interval**: how often the list of ZooKeeper and Exhibitor nodes should be updated from Exhibitor. - **port**: Exhibitor port. .. _kubernetes_settings: From 80354f648429e1146ff3f9206555ad3c826000dd Mon Sep 17 00:00:00 2001 From: 0m1xa <43731080+0m1xa@users.noreply.github.com> Date: Wed, 1 Apr 2020 16:51:19 +0300 Subject: [PATCH 17/24] Update Dockerfile (#1461) On postgres:12 find command without this pattern deletes file i18n_ctype which is needed by localdef. And localedef exit with code !=0 --- Dockerfile | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Dockerfile b/Dockerfile index c37ca001..a5e23e33 100644 --- a/Dockerfile +++ b/Dockerfile @@ -30,7 +30,7 @@ RUN set -ex \ \ # Cleanup all locales but en_US.UTF-8 && find /usr/share/i18n/charmaps/ -type f ! -name UTF-8.gz -delete \ - && find /usr/share/i18n/locales/ -type f ! -name en_US ! -name en_GB ! -name i18n ! -name iso14651_t1 ! -name iso14651_t1_common ! -name 'translit_*' -delete \ + && find /usr/share/i18n/locales/ -type f ! -name en_US ! -name en_GB ! -name i18n* ! -name iso14651_t1 ! -name iso14651_t1_common ! -name 'translit_*' -delete \ && echo 'en_US.UTF-8 UTF-8' > /usr/share/i18n/SUPPORTED \ \ # Make sure we have a en_US.UTF-8 locale available From 92d74af06e0b3576729d966b3905481ba6dded6b Mon Sep 17 00:00:00 2001 From: Kaarel Moppel Date: Wed, 1 Apr 2020 16:51:50 +0300 Subject: [PATCH 18/24] Better patronictl help text for the "flush" subcommand (#1466) Right now it's not really clear from --help what it does, flushing in software context usually means persisting...so actually it's the opposite, so make the intention more explicit. --- patroni/ctl.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/patroni/ctl.py b/patroni/ctl.py index b45e38fc..217cb1f7 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -868,7 +868,7 @@ def scaffold(obj, cluster_name, sysid): click.echo("Cluster {0} has been created successfully".format(cluster_name)) -@ctl.command('flush', help='Flush scheduled events') +@ctl.command('flush', help='Discard scheduled events (restarts only currently)') @click.argument('cluster_name') @click.argument('member_names', nargs=-1) @click.argument('target', type=click.Choice(['restart'])) From d58006319b1e72e3049a3a2d68d18e1e66c360d9 Mon Sep 17 00:00:00 2001 From: Kaarel Moppel Date: Wed, 1 Apr 2020 16:52:43 +0300 Subject: [PATCH 19/24] Patronictl - fail if a config file is specified explicitly but not found (#1467) $ python3 patronictl.py -c postgresql0.yml list Error: Provided config file postgresql0.yml not existing or no read rights. Check the -c/--config-file parameter --- patroni/ctl.py | 6 +++++- tests/test_ctl.py | 7 ++++++- 2 files changed, 11 insertions(+), 2 deletions(-) diff --git a/patroni/ctl.py b/patroni/ctl.py index 217cb1f7..08c1969d 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -69,7 +69,11 @@ def load_config(path, dcs): from patroni.config import Config if not (os.path.exists(path) and os.access(path, os.R_OK)): - logging.debug('Ignoring configuration file "%s". It does not exists or is not readable.', path) + if path != CONFIG_FILE_PATH: # bail if non-default config location specified but file not found / readable + raise PatroniCtlException('Provided config file {0} not existing or no read rights.' + ' Check the -c/--config-file parameter'.format(path)) + else: + logging.debug('Ignoring configuration file "%s". It does not exists or is not readable.', path) else: logging.debug('Loading configuration from file %s', path) config = Config(path, validator=None).copy() diff --git a/tests/test_ctl.py b/tests/test_ctl.py index 239058a4..d7d0c689 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -24,7 +24,6 @@ CONFIG_FILE_PATH = './test-ctl.yaml' def test_rw_config(): runner = CliRunner() with runner.isolated_filesystem(): - load_config(CONFIG_FILE_PATH + '/dummy', None) store_config({'etcd': {'host': 'localhost:2379'}}, CONFIG_FILE_PATH + '/dummy') load_config(CONFIG_FILE_PATH + '/dummy', '0.0.0.0') os.remove(CONFIG_FILE_PATH + '/dummy') @@ -43,6 +42,12 @@ class TestCtl(unittest.TestCase): self.runner = CliRunner() self.e = get_dcs({'etcd': {'ttl': 30, 'host': 'ok:2379', 'retry_timeout': 10}}, 'foo') + def test_abort_on_missing_or_unaccessible_config(self): + runner = CliRunner() + with runner.isolated_filesystem(): + with self.assertRaises(PatroniCtlException): + load_config('./non-existing-config-file', None) + @patch('psycopg2.connect', psycopg2_connect) def test_get_cursor(self): self.assertIsNone(get_cursor(get_cluster_initialized_without_leader(), {}, role='master')) From e58680f833334a2ff556df16e618291c8f714024 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 9 Apr 2020 10:33:38 +0200 Subject: [PATCH 20/24] Convert recovery_min_apply_delay to ms before doing comparison (#1484) Fixes https://github.com/zalando/patroni/issues/1483 --- patroni/postgresql/config.py | 1 + 1 file changed, 1 insertion(+) diff --git a/patroni/postgresql/config.py b/patroni/postgresql/config.py index c6ed5ce3..e80505b1 100644 --- a/patroni/postgresql/config.py +++ b/patroni/postgresql/config.py @@ -614,6 +614,7 @@ class ConfigHandler(object): values[match.group(1)] = [value, True] self._recovery_conf_mtime = recovery_conf_mtime values.setdefault('recovery_min_apply_delay', ['0', True]) + values['recovery_min_apply_delay'][0] = parse_int(values['recovery_min_apply_delay'][0], 'ms') values.update({param: ['', True] for param in self._recovery_parameters_to_compare if param not in values}) return values, True From 369a93ce2a4ded82d0b5f8f36e8f826c5f33ae2a Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 9 Apr 2020 10:34:12 +0200 Subject: [PATCH 21/24] Handle cases when conn_url is not defined (#1482) On K8s when one of the Patroni pods in starting there is valid annotation yet, which could cause failure in patronictl. In addition to that handle cases if port isn't specified in the standby_cluster configuration. Close https://github.com/zalando/patroni/issues/1100 Close https://github.com/zalando/patroni/issues/1463 --- patroni/ctl.py | 9 +++++---- patroni/dcs/__init__.py | 13 ++++++++----- patroni/postgresql/__init__.py | 6 +++--- patroni/postgresql/config.py | 2 +- patroni/utils.py | 7 +++++-- tests/test_ctl.py | 7 +++---- 6 files changed, 25 insertions(+), 19 deletions(-) diff --git a/patroni/ctl.py b/patroni/ctl.py index 08c1969d..a8ca193c 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -735,20 +735,21 @@ def output_members(cluster, name, extended=False, fmt='pretty'): columns.append(c) # Show Host as 'host:port' if somebody is running on non-standard port or two nodes are running on the same host - append_port = any(m['port'] != 5432 for m in cluster['members']) or\ - len(set(m['host'] for m in cluster['members'])) < len(cluster['members']) + members = [m for m in cluster['members'] if 'host' in m] + append_port = any('port' in m and m['port'] != 5432 for m in members) or\ + len(set(m['host'] for m in cluster['members'])) < len(members) for m in cluster['members']: logging.debug(m) lag = m.get('lag', '') - m.update(cluster=name, member=m['name'], tl=m.get('timeline', ''), + m.update(cluster=name, member=m['name'], host=m.get('host'), tl=m.get('timeline', ''), role='' if m['role'] == 'replica' else m['role'].replace('_', ' ').title(), lag_in_mb=round(lag/1024/1024) if isinstance(lag, six.integer_types) else lag, pending_restart='*' if m.get('pending_restart') else '', tags=json.dumps(m['tags']) if m.get('tags') else '') - if append_port: + if append_port and m['host'] and m.get('port'): m['host'] = ':'.join([m['host'], str(m['port'])]) if 'scheduled_restart' in m: diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index cee96805..946a65e5 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -140,10 +140,10 @@ class Member(namedtuple('Member', 'index,name,session,data')): @property def conn_url(self): conn_url = self.data.get('conn_url') - conn_kwargs = self.data.get('conn_kwargs') if conn_url: return conn_url + conn_kwargs = self.data.get('conn_kwargs') if conn_kwargs: conn_url = uri('postgresql', (conn_kwargs.get('host'), conn_kwargs.get('port', 5432))) self.data['conn_url'] = conn_url @@ -151,16 +151,19 @@ class Member(namedtuple('Member', 'index,name,session,data')): def conn_kwargs(self, auth=None): defaults = { - "host": "", - "port": "", - "database": "" + "host": None, + "port": None, + "database": None } ret = self.data.get('conn_kwargs') if ret: defaults.update(ret) ret = defaults else: - r = urlparse(self.conn_url) + conn_url = self.conn_url + if not conn_url: + return {} # due to the invalid conn_url we don't care about authentication parameters + r = urlparse(conn_url) ret = { 'host': r.hostname, 'port': r.port or 5432, diff --git a/patroni/postgresql/__init__.py b/patroni/postgresql/__init__.py index d31b77d7..250845de 100644 --- a/patroni/postgresql/__init__.py +++ b/patroni/postgresql/__init__.py @@ -651,10 +651,10 @@ class Postgresql(object): return result @contextmanager - def get_replication_connection_cursor(self, host='localhost', port=5432, database=None, **kwargs): + def get_replication_connection_cursor(self, host='localhost', port=5432, **kwargs): conn_kwargs = self.config.replication.copy() - conn_kwargs.update(host=host, port=int(port), database=database or self._database, connect_timeout=3, - user=conn_kwargs.pop('username'), replication=1, options='-c statement_timeout=2000') + conn_kwargs.update(host=host, port=int(port) if port else None, user=conn_kwargs.pop('username'), + connect_timeout=3, replication=1, options='-c statement_timeout=2000') with get_connection_cursor(**conn_kwargs) as cur: yield cur diff --git a/patroni/postgresql/config.py b/patroni/postgresql/config.py index e80505b1..a358c373 100644 --- a/patroni/postgresql/config.py +++ b/patroni/postgresql/config.py @@ -652,7 +652,7 @@ class ConfigHandler(object): else: return False - return all(primary_conninfo.get(p) == str(v) for p, v in wanted_primary_conninfo.items()) + return all(primary_conninfo.get(p) == str(v) for p, v in wanted_primary_conninfo.items() if v is not None) def check_recovery_conf(self, member): """Returns a tuple. The first boolean element indicates that recovery params don't match diff --git a/patroni/utils.py b/patroni/utils.py index 6ac3a5a9..7886bd6c 100644 --- a/patroni/utils.py +++ b/patroni/utils.py @@ -390,9 +390,12 @@ def cluster_as_json(cluster): else: role = 'replica' + member = {'name': m.name, 'role': role, 'state': m.data.get('state', ''), 'api_url': m.api_url} conn_kwargs = m.conn_kwargs() - member = {'name': m.name, 'host': conn_kwargs['host'], 'port': int(conn_kwargs['port']), - 'role': role, 'state': m.data.get('state', ''), 'api_url': m.api_url} + if conn_kwargs.get('host'): + member['host'] = conn_kwargs['host'] + if conn_kwargs.get('port'): + member['port'] = int(conn_kwargs['port']) optional_attributes = ('timeline', 'pending_restart', 'scheduled_restart', 'tags') member.update({n: m.data[n] for n in optional_attributes if n in m.data}) diff --git a/tests/test_ctl.py b/tests/test_ctl.py index d7d0c689..ca0a4f04 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -72,10 +72,9 @@ class TestCtl(unittest.TestCase): def test_output_members(self): scheduled_at = datetime.now(tzutc) + timedelta(seconds=600) cluster = get_cluster_initialized_with_leader(Failover(1, 'foo', 'bar', scheduled_at)) - self.assertIsNone(output_members(cluster, name='abc', fmt='pretty')) - self.assertIsNone(output_members(cluster, name='abc', fmt='json')) - self.assertIsNone(output_members(cluster, name='abc', fmt='yaml')) - self.assertIsNone(output_members(cluster, name='abc', fmt='tsv')) + del cluster.members[1].data['conn_url'] + for fmt in ('pretty', 'json', 'yaml', 'tsv'): + self.assertIsNone(output_members(cluster, name='abc', fmt=fmt)) @patch('patroni.ctl.get_dcs') @patch.object(PoolManager, 'request', Mock(return_value=MockResponse())) From 27cda08ecefda4dccd545e37574232dd38fb78ba Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 9 Apr 2020 10:34:35 +0200 Subject: [PATCH 22/24] Improve unit-tests (#1479) * tests were failing on windows and macos * improve coverage --- patroni/validator.py | 4 ++-- tests/test_ctl.py | 20 +++++++++++--------- tests/test_kubernetes.py | 10 +++++++++- tests/test_patroni.py | 6 +++++- tests/test_validator.py | 7 +++++-- 5 files changed, 32 insertions(+), 15 deletions(-) diff --git a/patroni/validator.py b/patroni/validator.py index 44ec1f2a..7e777da0 100644 --- a/patroni/validator.py +++ b/patroni/validator.py @@ -110,8 +110,8 @@ def validate_data_dir(data_dir): bin_dir = schema.data.get("postgresql", {}).get("bin_dir", None) major_version = get_major_version(bin_dir) if pgversion != major_version: - raise ConfigParseError("data_dir directory postgresql version ({}) doesn't match" - "with 'postgres --version' output ({})".format(pgversion, major_version)) + raise ConfigParseError("data_dir directory postgresql version ({}) doesn't match with " + "'postgres --version' output ({})".format(pgversion, major_version)) return True diff --git a/tests/test_ctl.py b/tests/test_ctl.py index ca0a4f04..eb2d9793 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -7,7 +7,7 @@ from datetime import datetime, timedelta from mock import patch, Mock from patroni.ctl import ctl, store_config, load_config, output_members, get_dcs, parse_dcs, \ get_all_members, get_any_member, get_cursor, query_member, configure, PatroniCtlException, apply_config_changes, \ - format_config_for_editing, show_diff, invoke_editor, format_pg_version, find_executable + format_config_for_editing, show_diff, invoke_editor, format_pg_version, find_executable, CONFIG_FILE_PATH from patroni.dcs.etcd import Client, Failover from patroni.utils import tzutc from psycopg2 import OperationalError @@ -18,16 +18,18 @@ from .test_etcd import etcd_read, socket_getaddrinfo from .test_ha import get_cluster_initialized_without_leader, get_cluster_initialized_with_leader, \ get_cluster_initialized_with_only_leader, get_cluster_not_initialized_without_leader, get_cluster, Member -CONFIG_FILE_PATH = './test-ctl.yaml' def test_rw_config(): + global CONFIG_FILE_PATH runner = CliRunner() with runner.isolated_filesystem(): - store_config({'etcd': {'host': 'localhost:2379'}}, CONFIG_FILE_PATH + '/dummy') - load_config(CONFIG_FILE_PATH + '/dummy', '0.0.0.0') - os.remove(CONFIG_FILE_PATH + '/dummy') - os.rmdir(CONFIG_FILE_PATH) + load_config(CONFIG_FILE_PATH, None) + CONFIG_PATH = './test-ctl.yaml' + store_config({'etcd': {'host': 'localhost:2379'}}, CONFIG_PATH + '/dummy') + load_config(CONFIG_PATH + '/dummy', '0.0.0.0') + os.remove(CONFIG_PATH + '/dummy') + os.rmdir(CONFIG_PATH) @patch('patroni.ctl.load_config', @@ -42,11 +44,11 @@ class TestCtl(unittest.TestCase): self.runner = CliRunner() self.e = get_dcs({'etcd': {'ttl': 30, 'host': 'ok:2379', 'retry_timeout': 10}}, 'foo') - def test_abort_on_missing_or_unaccessible_config(self): + def test_load_config(self): runner = CliRunner() with runner.isolated_filesystem(): - with self.assertRaises(PatroniCtlException): - load_config('./non-existing-config-file', None) + self.assertRaises(PatroniCtlException, load_config, './non-existing-config-file', None) + self.assertRaises(PatroniCtlException, load_config, './non-existing-config-file', None) @patch('psycopg2.connect', psycopg2_connect) def test_get_cursor(self): diff --git a/tests/test_kubernetes.py b/tests/test_kubernetes.py index c853fc4a..ed9b2141 100644 --- a/tests/test_kubernetes.py +++ b/tests/test_kubernetes.py @@ -33,13 +33,18 @@ def mock_config_map(*args, **kwargs): mock.metadata.resource_version = '2' return mock - +@patch('socket.TCP_KEEPIDLE', 4, create=True) +@patch('socket.TCP_KEEPINTVL', 5, create=True) +@patch('socket.TCP_KEEPCNT', 6, create=True) @patch.object(k8s_client.CoreV1Api, 'patch_namespaced_config_map', mock_config_map) @patch.object(k8s_client.CoreV1Api, 'create_namespaced_config_map', mock_config_map) @patch('kubernetes.client.api_client.ThreadPool', Mock(), create=True) @patch.object(Thread, 'start', Mock()) class TestKubernetes(unittest.TestCase): + @patch('socket.TCP_KEEPIDLE', 4, create=True) + @patch('socket.TCP_KEEPINTVL', 5, create=True) + @patch('socket.TCP_KEEPCNT', 6, create=True) @patch('kubernetes.config.load_kube_config', Mock()) @patch.object(k8s_client.CoreV1Api, 'list_namespaced_config_map', mock_list_namespaced_config_map) @patch.object(k8s_client.CoreV1Api, 'list_namespaced_pod', mock_list_namespaced_pod) @@ -149,6 +154,9 @@ class TestKubernetes(unittest.TestCase): class TestCacheBuilder(unittest.TestCase): + @patch('socket.TCP_KEEPIDLE', 4, create=True) + @patch('socket.TCP_KEEPINTVL', 5, create=True) + @patch('socket.TCP_KEEPCNT', 6, create=True) @patch('kubernetes.config.load_kube_config', Mock()) @patch('kubernetes.client.api_client.ThreadPool', Mock(), create=True) @patch.object(Thread, 'start', Mock()) diff --git a/tests/test_patroni.py b/tests/test_patroni.py index 0b395d59..42370253 100644 --- a/tests/test_patroni.py +++ b/tests/test_patroni.py @@ -44,7 +44,11 @@ class TestPatroni(unittest.TestCase): def test_no_config(self): self.assertRaises(SystemExit, patroni_main) - @patch('pkgutil.get_importer', Mock(return_value=MockFrozenImporter())) + @patch('sys.argv', ['patroni.py', '--validate-config', 'postgres0.yml']) + def test_validate_config(self): + self.assertRaises(SystemExit, patroni_main) + + @patch('pkgutil.iter_importers', Mock(return_value=[MockFrozenImporter()])) @patch('sys.frozen', Mock(return_value=True), create=True) @patch.object(BaseHTTPServer.HTTPServer, '__init__', Mock()) @patch.object(etcd.Client, 'read', etcd_read) diff --git a/tests/test_validator.py b/tests/test_validator.py index b562464e..df225857 100644 --- a/tests/test_validator.py +++ b/tests/test_validator.py @@ -1,6 +1,7 @@ import copy import os import socket +import tempfile import unittest from mock import Mock, patch, mock_open @@ -57,8 +58,8 @@ config = { "superuser": {"username": "user"}, "rewind": {"username": "user"}, }, - "data_dir": "/tmp/data_dir", - "bin_dir": "/tmp/bin_dir", + "data_dir": os.path.join(tempfile.gettempdir(), "data_dir"), + "bin_dir": os.path.join(tempfile.gettempdir(), "bin_dir"), "parameters": { "unix_socket_directories": "." }, @@ -84,6 +85,8 @@ files = [] def isfile_side_effect(arg): + if arg.endswith('.exe'): + arg = arg[:-4] return arg in files From e3335bea1afe765599938d569f7965dc42d4540b Mon Sep 17 00:00:00 2001 From: ksarabu1 <62157128+ksarabu1@users.noreply.github.com> Date: Wed, 15 Apr 2020 06:18:49 -0400 Subject: [PATCH 23/24] Master stop timeout (#1445) ## Feature: Postgres stop timeout Switchover/Failover operation hangs on signal_stop (or checkpoint) call when postmaster doesn't respond or hangs for some reason(Issue described in [1371](https://github.com/zalando/patroni/issues/1371)). This is leading to service loss for an extended period of time until the hung postmaster starts responding or it is killed by some other actor. ### master_stop_timeout The number of seconds Patroni is allowed to wait when stopping Postgres and effective only when synchronous_mode is enabled. When set to > 0 and the synchronous_mode is enabled, Patroni sends SIGKILL to the postmaster if the stop operation is running for more than the value set by master_stop_timeout. Set the value according to your durability/availability tradeoff. If the parameter is not set or set <= 0, master_stop_timeout does not apply. --- docs/SETTINGS.rst | 1 + patroni/config.py | 1 + patroni/ha.py | 19 +++++++++++------ patroni/postgresql/__init__.py | 33 +++++++++++++++++++++++------ patroni/postgresql/postmaster.py | 36 ++++++++++++++++++++++++++++++++ tests/__init__.py | 1 + tests/test_ha.py | 11 ++++++++++ tests/test_postgresql.py | 19 ++++++++++++++++- tests/test_postmaster.py | 31 +++++++++++++++++++++++++++ 9 files changed, 139 insertions(+), 13 deletions(-) diff --git a/docs/SETTINGS.rst b/docs/SETTINGS.rst index 1dca57b7..7a2ed244 100644 --- a/docs/SETTINGS.rst +++ b/docs/SETTINGS.rst @@ -16,6 +16,7 @@ Dynamic configuration is stored in the DCS (Distributed Configuration Store) and - **retry\_timeout**: timeout for DCS and PostgreSQL operation retries (in seconds). DCS or network issues shorter than this will not cause Patroni to demote the leader. Default value: 10 - **maximum\_lag\_on\_failover**: the maximum bytes a follower may lag to be able to participate in leader election. - **master\_start\_timeout**: the amount of time a master is allowed to recover from failures before failover is triggered (in seconds). Default is 300 seconds. When set to 0 failover is done immediately after a crash is detected if possible. When using asynchronous replication a failover can cause lost transactions. Worst case failover time for master failure is: loop\_wait + master\_start\_timeout + loop\_wait, unless master\_start\_timeout is zero, in which case it's just loop\_wait. Set the value according to your durability/availability tradeoff. +- **master\_stop\_timeout**: The number of seconds Patroni is allowed to wait when stopping Postgres and effective only when synchronous_mode is enabled. When set to > 0 and the synchronous_mode is enabled, Patroni sends SIGKILL to the postmaster if the stop operation is running for more than the value set by master_stop_timeout. Set the value according to your durability/availability tradeoff. If the parameter is not set or set <= 0, master_stop_timeout does not apply. - **synchronous\_mode**: turns on synchronous replication mode. In this mode a replica will be chosen as synchronous and only the latest leader and synchronous replica are able to participate in leader election. Synchronous mode makes sure that successfully committed transactions will not be lost at failover, at the cost of losing availability for writes when Patroni cannot ensure transaction durability. See :ref:`replication modes documentation ` for details. - **synchronous\_mode\_strict**: prevents disabling synchronous replication if no synchronous replicas are available, blocking all client writes to the master. See :ref:`replication modes documentation ` for details. - **postgresql**: diff --git a/patroni/config.py b/patroni/config.py index feb785a0..1afce2f6 100644 --- a/patroni/config.py +++ b/patroni/config.py @@ -59,6 +59,7 @@ class Config(object): 'maximum_lag_on_failover': 1048576, 'check_timeline': False, 'master_start_timeout': 300, + 'master_stop_timeout': 0, 'synchronous_mode': False, 'synchronous_mode_strict': False, 'standby_cluster': { diff --git a/patroni/ha.py b/patroni/ha.py index ceecb388..bdc943e4 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -14,7 +14,7 @@ from patroni.exceptions import DCSError, PostgresConnectionException, PatroniExc from patroni.postgresql import ACTION_ON_START, ACTION_ON_ROLE_CHANGE from patroni.postgresql.misc import postgres_version_to_int from patroni.postgresql.rewind import Rewind -from patroni.utils import polling_loop, tzutc, is_standby_cluster as _is_standby_cluster +from patroni.utils import polling_loop, tzutc, is_standby_cluster as _is_standby_cluster, parse_int from patroni.dcs import RemoteMember from threading import RLock @@ -94,6 +94,11 @@ class Ha(object): else: return self.patroni.config.check_mode(mode) + def master_stop_timeout(self): + """ Master stop timeout """ + ret = parse_int(self.patroni.config['master_stop_timeout']) + return ret if ret and ret > 0 and self.is_synchronous_mode() else None + def is_paused(self): return self.check_mode('pause') @@ -772,7 +777,8 @@ class Ha(object): self._rewind.trigger_check_diverged_lsn() self.state_handler.stop(mode_control['stop'], checkpoint=mode_control['checkpoint'], - on_safepoint=self.watchdog.disable if self.watchdog.is_running else None) + on_safepoint=self.watchdog.disable if self.watchdog.is_running else None, + stop_timeout=self.master_stop_timeout()) self.state_handler.set_role('demoted') self.set_is_leader(False) @@ -1083,7 +1089,7 @@ class Ha(object): return (False, 'restart failed') def _do_reinitialize(self, cluster): - self.state_handler.stop('immediate') + self.state_handler.stop('immediate', stop_timeout=self.patroni.config['retry_timeout']) # Commented redundant data directory cleanup here # self.state_handler.remove_data_directory() @@ -1156,7 +1162,7 @@ class Ha(object): def cancel_initialization(self): logger.info('removing initialize key after failed attempt to bootstrap the cluster') self.dcs.cancel_initialization() - self.state_handler.stop('immediate') + self.state_handler.stop('immediate', stop_timeout=self.patroni.config['retry_timeout']) self.state_handler.move_data_directory() raise PatroniException('Failed to bootstrap cluster') @@ -1279,7 +1285,7 @@ class Ha(object): # is data directory empty? if self.state_handler.data_directory_empty(): self.state_handler.set_role('uninitialized') - self.state_handler.stop('immediate') + self.state_handler.stop('immediate', stop_timeout=self.patroni.config['retry_timeout']) # In case datadir went away while we were master. self.watchdog.disable() @@ -1373,7 +1379,8 @@ class Ha(object): # This might not be the desired behavior of users, as a graceful shutdown of the host can mean lost data. # We probably need to something smarter here. disable_wd = self.watchdog.disable if self.watchdog.is_running else None - self.while_not_sync_standby(lambda: self.state_handler.stop(checkpoint=False, on_safepoint=disable_wd)) + self.while_not_sync_standby(lambda: self.state_handler.stop(checkpoint=False, on_safepoint=disable_wd, + stop_timeout=self.master_stop_timeout())) if not self.state_handler.is_running(): if self.has_lock(): self.dcs.delete_leader() diff --git a/patroni/postgresql/__init__.py b/patroni/postgresql/__init__.py index 250845de..8e2d017f 100644 --- a/patroni/postgresql/__init__.py +++ b/patroni/postgresql/__init__.py @@ -19,6 +19,7 @@ from patroni.postgresql.slots import SlotsHandler from patroni.exceptions import PostgresConnectionException from patroni.utils import Retry, RetryFailedError, polling_loop, data_directory_is_empty from threading import current_thread, Lock +from psutil import TimeoutExpired logger = logging.getLogger(__name__) @@ -462,11 +463,13 @@ class Postgresql(object): else: return None - def checkpoint(self, connect_kwargs=None): + def checkpoint(self, connect_kwargs=None, timeout=None): check_not_is_in_recovery = connect_kwargs is not None connect_kwargs = connect_kwargs or self.config.local_connect_kwargs for p in ['connect_timeout', 'options']: connect_kwargs.pop(p, None) + if timeout: + connect_kwargs['connect_timeout'] = timeout try: with get_connection_cursor(**connect_kwargs) as cur: cur.execute("SET statement_timeout = 0") @@ -479,7 +482,7 @@ class Postgresql(object): logger.exception('Exception during CHECKPOINT') return 'not accessible or not healty' - def stop(self, mode='fast', block_callbacks=False, checkpoint=None, on_safepoint=None): + def stop(self, mode='fast', block_callbacks=False, checkpoint=None, on_safepoint=None, stop_timeout=None): """Stop PostgreSQL Supports a callback when a safepoint is reached. A safepoint is when no user backend can return a successful @@ -491,7 +494,7 @@ class Postgresql(object): if checkpoint is None: checkpoint = False if mode == 'immediate' else True - success, pg_signaled = self._do_stop(mode, block_callbacks, checkpoint, on_safepoint) + success, pg_signaled = self._do_stop(mode, block_callbacks, checkpoint, on_safepoint, stop_timeout) if success: # block_callbacks is used during restart to avoid # running start/stop callbacks in addition to restart ones @@ -504,7 +507,7 @@ class Postgresql(object): self.set_state('stop failed') return success - def _do_stop(self, mode, block_callbacks, checkpoint, on_safepoint): + def _do_stop(self, mode, block_callbacks, checkpoint, on_safepoint, stop_timeout): postmaster = self.is_running() if not postmaster: if on_safepoint: @@ -512,7 +515,7 @@ class Postgresql(object): return True, False if checkpoint and not self.is_starting(): - self.checkpoint() + self.checkpoint(timeout=stop_timeout) if not block_callbacks: self.set_state('stopping') @@ -531,10 +534,28 @@ class Postgresql(object): postmaster.wait_for_user_backends_to_close() on_safepoint() - postmaster.wait() + try: + postmaster.wait(timeout=stop_timeout) + except TimeoutExpired: + logger.warning("Timeout during postmaster stop, aborting Postgres.") + if not self.terminate_postmaster(postmaster, mode, stop_timeout): + postmaster.wait() return True, True + def terminate_postmaster(self, postmaster, mode, stop_timeout): + if mode in ['fast', 'smart']: + try: + success = postmaster.signal_stop('immediate', self.pgcommand('pg_ctl')) + if success: + return True + postmaster.wait(timeout=stop_timeout) + return True + except TimeoutExpired: + pass + logger.warning("Sending SIGKILL to Postmaster and its children") + return postmaster.signal_kill() + def terminate_starting_postmaster(self, postmaster): """Terminates a postmaster that has not yet opened ports or possibly even written a pid file. Blocks until the process goes away.""" diff --git a/patroni/postgresql/postmaster.py b/patroni/postgresql/postmaster.py index b95ead33..cac48102 100644 --- a/patroni/postgresql/postmaster.py +++ b/patroni/postgresql/postmaster.py @@ -105,6 +105,42 @@ class PostmasterProcess(psutil.Process): except psutil.NoSuchProcess: return None + def signal_kill(self): + """to suspend and kill postmaster and all children + + :returns True if postmaster and children are killed, False if error + """ + try: + self.suspend() + except psutil.NoSuchProcess: + return True + except psutil.Error as e: + logger.warning('Failed to suspend postmaster: %s', e) + + try: + children = self.children(recursive=True) + except psutil.NoSuchProcess: + return True + except psutil.Error as e: + logger.warning('Failed to get a list of postmaster children: %s', e) + children = [] + + try: + self.kill() + except psutil.NoSuchProcess: + return True + except psutil.Error as e: + logger.warning('Could not kill postmaster: %s', e) + return False + + for child in children: + try: + child.kill() + except psutil.Error: + pass + psutil.wait_procs(children + [self]) + return True + def signal_stop(self, mode, pg_ctl='pg_ctl'): """Signal postmaster process to stop diff --git a/tests/__init__.py b/tests/__init__.py index a646c137..7adea2a0 100644 --- a/tests/__init__.py +++ b/tests/__init__.py @@ -66,6 +66,7 @@ class MockPostmaster(object): self.wait_for_user_backends_to_close = Mock() self.signal_stop = Mock(return_value=None) self.wait = Mock() + self.signal_kill = Mock(return_value=False) class MockCursor(object): diff --git a/tests/test_ha.py b/tests/test_ha.py index 1840c901..d7d16abe 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -817,6 +817,17 @@ class TestHa(PostgresInit): self.assertEqual(self.ha.run_cycle(), 'stopped PostgreSQL to fail over after a crash') demote.assert_called_once() + def test_master_stop_timeout(self): + self.assertEqual(self.ha.master_stop_timeout(), None) + self.ha.patroni.config.set_dynamic_configuration({'master_stop_timeout': 30}) + with patch.object(Ha, 'is_synchronous_mode', Mock(return_value=True)): + self.assertEqual(self.ha.master_stop_timeout(), 30) + self.ha.patroni.config.set_dynamic_configuration({'master_stop_timeout': 30}) + with patch.object(Ha, 'is_synchronous_mode', Mock(return_value=False)): + self.assertEqual(self.ha.master_stop_timeout(), None) + self.ha.patroni.config.set_dynamic_configuration({'master_stop_timeout': None}) + self.assertEqual(self.ha.master_stop_timeout(), None) + @patch('patroni.postgresql.Postgresql.follow') def test_demote_immediate(self, follow): self.ha.has_lock = true diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 8a81ddc5..721c6169 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -1,5 +1,6 @@ import mock # for the mock.call method, importing it without a namespace breaks python3 import os +import psutil import psycopg2 import re import subprocess @@ -177,6 +178,17 @@ class TestPostgresql(BaseTestPostgresql): mock_callback.assert_called() mock_postmaster.signal_stop.assert_called() + # Timed out waiting for fast shutdown triggers immediate shutdown + mock_postmaster.wait.side_effect = [psutil.TimeoutExpired(30), psutil.TimeoutExpired(30), Mock()] + mock_callback.reset_mock() + self.assertTrue(self.p.stop(on_safepoint=mock_callback, stop_timeout=30)) + mock_callback.assert_called() + mock_postmaster.signal_stop.assert_called() + + # Immediate shutdown succeeded + mock_postmaster.wait.side_effect = [psutil.TimeoutExpired(30), Mock()] + self.assertTrue(self.p.stop(on_safepoint=mock_callback, stop_timeout=30)) + # Stop signal failed mock_postmaster.signal_stop.return_value = False self.assertFalse(self.p.stop()) @@ -187,6 +199,11 @@ class TestPostgresql(BaseTestPostgresql): self.assertTrue(self.p.stop(on_safepoint=mock_callback)) mock_callback.assert_called() + # Fast shutdown is timed out but when immediate postmaster is already gone + mock_postmaster.wait.side_effect = [psutil.TimeoutExpired(30), Mock()] + mock_postmaster.signal_stop.side_effect = [None, True] + self.assertTrue(self.p.stop(on_safepoint=mock_callback, stop_timeout=30)) + def test_restart(self): self.p.start = Mock(return_value=False) self.assertFalse(self.p.restart()) @@ -203,7 +220,7 @@ class TestPostgresql(BaseTestPostgresql): self.assertEqual(self.p.checkpoint({'user': 'postgres'}), 'is_in_recovery=true') with patch.object(MockCursor, 'execute', Mock(return_value=None)): self.assertIsNone(self.p.checkpoint()) - self.assertEqual(self.p.checkpoint(), 'not accessible or not healty') + self.assertEqual(self.p.checkpoint(timeout=10), 'not accessible or not healty') @patch('patroni.postgresql.config.mtime', mock_mtime) @patch('patroni.postgresql.config.ConfigHandler._get_pg_settings') diff --git a/tests/test_postmaster.py b/tests/test_postmaster.py index 7189a626..4e4d1397 100644 --- a/tests/test_postmaster.py +++ b/tests/test_postmaster.py @@ -63,6 +63,37 @@ class TestPostmasterProcess(unittest.TestCase): mock_init.side_effect = None self.assertNotEqual(PostmasterProcess.from_pid(123), None) + @patch('psutil.Process.__init__', Mock()) + @patch('psutil.wait_procs', Mock()) + @patch('psutil.Process.suspend') + @patch('psutil.Process.children') + @patch('psutil.Process.kill') + def test_signal_kill(self, mock_kill, mock_children, mock_suspend): + proc = PostmasterProcess(123) + + # all processes successfully stopped + mock_children.return_value = [Mock()] + mock_children.return_value[0].kill.side_effect = psutil.Error + self.assertTrue(proc.signal_kill()) + + # postmaster has gone before suspend + mock_suspend.side_effect = psutil.NoSuchProcess(123) + self.assertTrue(proc.signal_kill()) + + # postmaster has gone before we got a list of children + mock_suspend.side_effect = psutil.Error() + mock_children.side_effect = psutil.NoSuchProcess(123) + self.assertTrue(proc.signal_kill()) + + # postmaster has gone after we got a list of children + mock_children.side_effect = psutil.Error() + mock_kill.side_effect = psutil.NoSuchProcess(123) + self.assertTrue(proc.signal_kill()) + + # failed to kill postmaster + mock_kill.side_effect = psutil.AccessDenied(123) + self.assertFalse(proc.signal_kill()) + @patch('psutil.Process.__init__', Mock()) @patch('psutil.Process.send_signal') @patch('psutil.Process.pid', Mock(return_value=123)) From 337f9efc9ee572fa52af91b46717f1e9a77f1c42 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 15 Apr 2020 12:19:18 +0200 Subject: [PATCH 24/24] Improve patronictl list output (#1486) The redundant column `Column` will be presented in the table header. Depending on output format `Tags` are serialized differently: * For *pretty* format YAML is used, every element on the new line * For *tsv* format for YAML is also used, but all elements and on the same line (similar to JSON) * For *json* and *yaml* formats `Tags` are serialized into an appropriate format.
Examples of output in pretty formats: ```bash $ patronictl list + Cluster: batman (6813309862653668387) +---------+----+-----------+---------------------+ | Member | Host | Role | State | TL | Lag in MB | Tags | +-------------+----------------+--------+---------+----+-----------+---------------------+ | postgresql0 | 127.0.0.1:5432 | Leader | running | 3 | | clonefrom: true | | | | | | | | noloadbalance: true | | | | | | | | nosync: true | +-------------+----------------+--------+---------+----+-----------+---------------------+ | postgresql1 | 127.0.0.1:5433 | | running | 3 | 0.0 | | +-------------+----------------+--------+---------+----+-----------+---------------------+ $ patronictl list badclustername + Cluster: badclustername (uninitialized) ------+ | Member | Host | Role | State | TL | Lag in MB | +--------+------+------+-------+----+-----------+ +--------+------+------+-------+----+-----------+ ```
Example in tsv format: ```bash Cluster Member Host Role State TL Lag in MB Pending restart Tags batman postgresql0 127.0.0.1:5432 Leader running 2 batman postgresql1 127.0.0.1:5433 running 2 0 {clonefrom: true, nofailover: true, noloadbalance: true, replicatefrom: postgresql0} batman postgresql2 127.0.0.1:5434 running 2 0 * {replicatefrom: postgres1} ```
In addition to that, `patronictl list` command will stop showing keys with empty values in `json` and `yaml` formats.
Examples: ```yaml $ patronictl list -f yaml - Cluster: batman Host: 127.0.0.1:5432 Member: postgresql0 Role: Leader State: running TL: 2 - Cluster: batman Host: 127.0.0.1:5433 Lag in MB: 0 Member: postgresql1 State: running TL: 2 Tags: clonefrom: true nofailover: true noloadbalance: true replicatefrom: postgresql0 - Cluster: batman Host: 127.0.0.1:5434 Lag in MB: 0 Member: postgresql2 Pending restart: '*' State: running TL: 2 Tags: replicatefrom: postgres1 ``` ```json $ patronictl list -f json | jq . [ { "Cluster": "batman", "Member": "postgresql0", "Host": "127.0.0.1:5432", "Role": "Leader", "State": "running", "TL": 2 }, { "Cluster": "batman", "Member": "postgresql1", "Host": "127.0.0.1:5433", "State": "running", "TL": 2, "Lag in MB": 0, "Tags": { "nofailover": true, "noloadbalance": true, "replicatefrom": "postgresql0", "clonefrom": true } }, { "Cluster": "batman", "Member": "postgresql2", "Host": "127.0.0.1:5434", "State": "running", "TL": 2, "Lag in MB": 0, "Pending restart": "*", "Tags": { "replicatefrom": "postgres1" } } ] ```
--- patroni/ctl.py | 109 ++++++++++++++++++++++++++++------------------ tests/__init__.py | 1 + tests/test_ctl.py | 3 -- 3 files changed, 67 insertions(+), 46 deletions(-) diff --git a/patroni/ctl.py b/patroni/ctl.py index a8ca193c..4f2fcd0b 100644 --- a/patroni/ctl.py +++ b/patroni/ctl.py @@ -31,7 +31,7 @@ from patroni.postgresql.misc import postgres_version_to_int from patroni.utils import cluster_as_json, patch_config, polling_loop from patroni.request import PatroniRequest from patroni.version import __version__ -from prettytable import PrettyTable +from prettytable import ALL, FRAME, PrettyTable from six.moves.urllib_parse import urlparse CONFIG_DIR_PATH = click.get_app_dir('patroni') @@ -46,6 +46,34 @@ class PatroniCtlException(ClickException): pass +class PatronictlPrettyTable(PrettyTable): + + def __init__(self, header, *args, **kwargs): + PrettyTable.__init__(self, *args, **kwargs) + self.__table_header = header + self.__hline_num = 0 + self.__hline = None + + def _is_first_hline(self): + return self.__hline_num == 0 + + def _set_hline(self, value): + self.__hline = value + + def _get_hline(self): + ret = self.__hline + + # Inject nice table header + if self._is_first_hline() and self.__table_header: + header = self.__table_header[:len(ret) - 2] + ret = "".join([ret[0], header, ret[1 + len(header):]]) + + self.__hline_num += 1 + return ret + + _hrule = property(_get_hline, _set_hline) + + def parse_dcs(dcs): if dcs is None: return None @@ -94,7 +122,7 @@ def store_config(config, path): yaml.dump(config, fd) -option_format = click.option('--format', '-f', 'fmt', help='Output format (pretty, json, yaml)', default='pretty') +option_format = click.option('--format', '-f', 'fmt', help='Output format (pretty, tsv, json, yaml)', default='pretty') option_watchrefresh = click.option('-w', '--watch', type=float, help='Auto update the screen every X seconds') option_watch = click.option('-W', is_flag=True, help='Auto update the screen every 2 seconds') option_force = click.option('--force', is_flag=True, help='Do not ask for confirmation at any point') @@ -137,31 +165,33 @@ def request_patroni(member, method='GET', endpoint=None, data=None): return request_executor(member, method, endpoint, data) -def print_output(columns, rows=None, alignment=None, fmt='pretty', header=True, delimiter='\t'): - rows = rows or [] - if fmt == 'pretty': - t = PrettyTable(columns) - for k, v in (alignment or {}).items(): - t.align[k] = v - for r in rows: - t.add_row(r) - click.echo(t) - return +def print_output(columns, rows, alignment=None, fmt='pretty', header=None, delimiter='\t'): + if fmt in {'json', 'yaml', 'yml'}: + elements = [{k: v for k, v in zip(columns, r) if not header or str(v)} for r in rows] + func = json.dumps if fmt == 'json' else format_config_for_editing + click.echo(func(elements)) + elif fmt in {'pretty', 'tsv'}: + list_cluster = bool(header and columns and columns[0] == 'Cluster') + if list_cluster and 'Tags' in columns: # we want to format member tags as YAML + i = columns.index('Tags') + for row in rows: + if row[i]: + row[i] = format_config_for_editing(row[i], fmt == 'tsv').strip() + if list_cluster and fmt == 'pretty': # skip cluster name if pretty-printing + columns = columns[1:] if columns else [] + rows = [row[1:] for row in rows] - if fmt in ['json', 'yaml', 'yml']: - elements = [dict(zip(columns, r)) for r in rows] - if fmt == 'json': - click.echo(json.dumps(elements)) - elif fmt in ('yaml', 'yml'): - click.echo(yaml.safe_dump(elements, encoding=None, default_flow_style=False, allow_unicode=True, width=200)) - - if fmt == 'tsv': - if columns is not None and header: - click.echo(delimiter.join(columns)) - - for r in rows: - c = [str(c) for c in r] - click.echo(delimiter.join(c)) + if fmt == 'tsv': + for r in ([columns] if columns else []) + rows: + click.echo(delimiter.join(map(str, r))) + else: + hrules = ALL if any(any(isinstance(c, six.string_types) and '\n' in c for c in r) for r in rows) else FRAME + table = PatronictlPrettyTable(header, columns, hrules=hrules) + for k, v in (alignment or {}).items(): + table.align[k] = v + for r in rows: + table.add_row(r) + click.echo(table) def watching(w, watch, max_count=None, clear=True): @@ -310,8 +340,7 @@ def dsn(obj, cluster_name, role, member): @ctl.command('query', help='Query a Patroni PostgreSQL member') @arg_cluster_name -@option_format -@click.option('--format', 'fmt', help='Output format (pretty, json)', default='tsv') +@click.option('--format', 'fmt', help='Output format (pretty, tsv, json, yaml)', default='tsv') @click.option('--file', '-f', 'p_file', help='Execute the SQL commands from this file', type=click.File('rb')) @click.option('--password', help='force password prompt', is_flag=True) @click.option('-U', '--username', help='database user name', type=str) @@ -368,8 +397,8 @@ def query( if cursor is None: cluster = dcs.get_cluster() - output, cursor = query_member(cluster, cursor, member, role, command, connect_parameters) - print_output(None, output, fmt=fmt, delimiter=delimiter) + output, header = query_member(cluster, cursor, member, role, command, connect_parameters) + print_output(header, output, fmt=fmt, delimiter=delimiter) def query_member(cluster, cursor, member, role, command, connect_parameters): @@ -386,15 +415,8 @@ def query_member(cluster, cursor, member, role, command, connect_parameters): logging.debug(message) return [[timestamp(0), message]], None - cursor.execute('SELECT pg_catalog.pg_is_in_recovery()') - in_recovery = cursor.fetchone()[0] - - if in_recovery and role == 'master' or not in_recovery and role == 'replica': - cursor.connection.close() - return None, None - cursor.execute(command) - return cursor.fetchall(), cursor + return cursor.fetchall(), [d.name for d in cursor.description] except (psycopg2.OperationalError, psycopg2.DatabaseError) as oe: logging.debug(oe) if cursor is not None and not cursor.connection.closed: @@ -727,6 +749,7 @@ def switchover(obj, cluster_name, master, candidate, force, scheduled): def output_members(cluster, name, extended=False, fmt='pretty'): rows = [] logging.debug(cluster) + initialize = {None: 'uninitialized', '': 'initializing'}.get(cluster.initialize, cluster.initialize) cluster = cluster_as_json(cluster) columns = ['Cluster', 'Member', 'Host', 'Role', 'State', 'TL', 'Lag in MB'] @@ -746,8 +769,7 @@ def output_members(cluster, name, extended=False, fmt='pretty'): m.update(cluster=name, member=m['name'], host=m.get('host'), tl=m.get('timeline', ''), role='' if m['role'] == 'replica' else m['role'].replace('_', ' ').title(), lag_in_mb=round(lag/1024/1024) if isinstance(lag, six.integer_types) else lag, - pending_restart='*' if m.get('pending_restart') else '', - tags=json.dumps(m['tags']) if m.get('tags') else '') + pending_restart='*' if m.get('pending_restart') else '') if append_port and m['host'] and m.get('port'): m['host'] = ':'.join([m['host'], str(m['port'])]) @@ -760,7 +782,8 @@ def output_members(cluster, name, extended=False, fmt='pretty'): rows.append([m.get(n.lower().replace(' ', '_'), '') for n in columns]) - print_output(columns, rows, {'Lag in MB': 'r', 'TL': 'r', 'Tags': 'l'}, fmt) + print_output(columns, rows, {'Lag in MB': 'r', 'TL': 'r', 'Tags': 'l'}, + fmt, ' Cluster: {0} ({1}) '.format(name, initialize)) if fmt != 'pretty': # Omit service info when using machine-readable formats return @@ -1006,12 +1029,12 @@ def show_diff(before_editing, after_editing): click.echo(line.rstrip('\n')) -def format_config_for_editing(data): +def format_config_for_editing(data, default_flow_style=False): """Formats configuration as YAML for human consumption. :param data: configuration as nested dictionaries :returns unicode YAML of the configuration""" - return yaml.safe_dump(data, default_flow_style=False, encoding=None, allow_unicode=True) + return yaml.safe_dump(data, default_flow_style=default_flow_style, encoding=None, allow_unicode=True, width=200) def apply_config_changes(before_editing, data, kvpairs): diff --git a/tests/__init__.py b/tests/__init__.py index 7adea2a0..8a075d61 100644 --- a/tests/__init__.py +++ b/tests/__init__.py @@ -76,6 +76,7 @@ class MockCursor(object): self.closed = False self.rowcount = 0 self.results = [] + self.description = [Mock()] def execute(self, sql, *params): if sql.startswith('blabla'): diff --git a/tests/test_ctl.py b/tests/test_ctl.py index eb2d9793..b72785b6 100644 --- a/tests/test_ctl.py +++ b/tests/test_ctl.py @@ -206,9 +206,6 @@ class TestCtl(unittest.TestCase): rows = query_member(None, None, None, 'master', 'SELECT pg_catalog.pg_is_in_recovery()', {}) self.assertTrue('False' in str(rows)) - rows = query_member(None, None, None, 'replica', 'SELECT pg_catalog.pg_is_in_recovery()', {}) - self.assertEqual(rows, (None, None)) - with patch.object(MockCursor, 'execute', Mock(side_effect=OperationalError('bla'))): rows = query_member(None, None, None, 'replica', 'SELECT pg_catalog.pg_is_in_recovery()', {})