mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-31 16:49:46 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
710afd5952 | ||
|
|
096ee8f36f | ||
|
|
91e2be092c | ||
|
|
4148e0b5b2 | ||
|
|
5ceba81269 | ||
|
|
dbbe065a27 |
@@ -3,6 +3,28 @@
|
|||||||
Release notes
|
Release notes
|
||||||
=============
|
=============
|
||||||
|
|
||||||
|
Version 3.1.2
|
||||||
|
-------------
|
||||||
|
|
||||||
|
**Bugfixes**
|
||||||
|
|
||||||
|
- Fixed bug with ``wal_keep_size`` checks (Alexander Kukushkin)
|
||||||
|
|
||||||
|
The ``wal_keep_size`` is a GUC that normally has a unit and Patroni was failing to cast its value to ``int``. As a result the value of ``bootstrap.dcs`` was not written to the ``/config`` key afterwards.
|
||||||
|
|
||||||
|
- Detect and resolve inconsistencies between ``/sync`` key and ``synchronous_standby_names`` (Alexander Kukushkin)
|
||||||
|
|
||||||
|
Normally, Patroni updates ``/sync`` and ``synchronous_standby_names`` in a very specific order, but in case of a bug or when someone manually reset ``synchronous_standby_names``, Patroni was getting into an inconsistent state. As a result it was possible that the failover happens to an asynchronous node.
|
||||||
|
|
||||||
|
- Read GUC's values when joining running Postgres (Alexander Kukushkin)
|
||||||
|
|
||||||
|
When restarted in ``pause``, Patroni was discarding the ``synchronous_standby_names`` GUC from the ``postgresql.conf``. To solve it and avoid similar issues, Patroni will read GUC's value if it is joining an already running Postgres.
|
||||||
|
|
||||||
|
- Silenced annoying warnings when checking for node uniqueness (Alexander Kukushkin)
|
||||||
|
|
||||||
|
``WARNING`` messages are produced by ``urllib3`` if Patroni is quickly restarted.
|
||||||
|
|
||||||
|
|
||||||
Version 3.1.1
|
Version 3.1.1
|
||||||
-------------
|
-------------
|
||||||
|
|
||||||
|
|||||||
@@ -68,6 +68,7 @@ Scenario: check API requests for the primary-replica pair in the pause mode
|
|||||||
When I kill postmaster on postgres1
|
When I kill postmaster on postgres1
|
||||||
And I issue a GET request to http://127.0.0.1:8009/replica
|
And I issue a GET request to http://127.0.0.1:8009/replica
|
||||||
Then I receive a response code 503
|
Then I receive a response code 503
|
||||||
|
And "members/postgres1" key in DCS has state=stopped after 10 seconds
|
||||||
When I run patronictl.py restart batman postgres1 --force
|
When I run patronictl.py restart batman postgres1 --force
|
||||||
Then I receive a response returncode 0
|
Then I receive a response returncode 0
|
||||||
Then replication works from postgres0 to postgres1 after 20 seconds
|
Then replication works from postgres0 to postgres1 after 20 seconds
|
||||||
@@ -76,7 +77,7 @@ Scenario: check API requests for the primary-replica pair in the pause mode
|
|||||||
Then I receive a response code 200
|
Then I receive a response code 200
|
||||||
And I receive a response state running
|
And I receive a response state running
|
||||||
And I receive a response role replica
|
And I receive a response role replica
|
||||||
When I run patronictl.py reinit batman postgres1 --force
|
When I run patronictl.py reinit batman postgres1 --force --wait
|
||||||
Then I receive a response returncode 0
|
Then I receive a response returncode 0
|
||||||
And I receive a response output "Success: reinitialize for member postgres1"
|
And I receive a response output "Success: reinitialize for member postgres1"
|
||||||
And postgres1 role is the secondary after 30 seconds
|
And postgres1 role is the secondary after 30 seconds
|
||||||
|
|||||||
+8
-3
@@ -65,6 +65,8 @@ class Patroni(AbstractPatroniDaemon):
|
|||||||
|
|
||||||
def ensure_unique_name(self) -> None:
|
def ensure_unique_name(self) -> None:
|
||||||
"""A helper method to prevent splitbrain from operator naming error."""
|
"""A helper method to prevent splitbrain from operator naming error."""
|
||||||
|
from urllib.parse import urlparse
|
||||||
|
from urllib3.connection import HTTPConnection
|
||||||
from patroni.dcs import Member
|
from patroni.dcs import Member
|
||||||
|
|
||||||
cluster = self.dcs.get_cluster()
|
cluster = self.dcs.get_cluster()
|
||||||
@@ -74,9 +76,12 @@ class Patroni(AbstractPatroniDaemon):
|
|||||||
if not isinstance(member, Member):
|
if not isinstance(member, Member):
|
||||||
return
|
return
|
||||||
try:
|
try:
|
||||||
_ = self.request(member, endpoint="/liveness", timeout=3)
|
parts = urlparse(member.api_url)
|
||||||
logger.fatal("Can't start; there is already a node named '%s' running", self.config['name'])
|
if isinstance(parts.hostname, str):
|
||||||
sys.exit(1)
|
connection = HTTPConnection(parts.hostname, port=parts.port or 80, timeout=3)
|
||||||
|
connection.connect()
|
||||||
|
logger.fatal("Can't start; there is already a node named '%s' running", self.config['name'])
|
||||||
|
sys.exit(1)
|
||||||
except Exception:
|
except Exception:
|
||||||
return
|
return
|
||||||
|
|
||||||
|
|||||||
+2
-1
@@ -342,7 +342,8 @@ class Config(object):
|
|||||||
elif not is_local:
|
elif not is_local:
|
||||||
validator = ConfigHandler.CMDLINE_OPTIONS[name][1]
|
validator = ConfigHandler.CMDLINE_OPTIONS[name][1]
|
||||||
if validator(value):
|
if validator(value):
|
||||||
pg_params[name] = int(value) if isinstance(validator, IntValidator) else value
|
int_val = parse_int(value) if isinstance(validator, IntValidator) else None
|
||||||
|
pg_params[name] = int_val if isinstance(int_val, int) else value
|
||||||
else:
|
else:
|
||||||
logger.warning("postgresql parameter %s=%s failed validation, defaulting to %s",
|
logger.warning("postgresql parameter %s=%s failed validation, defaulting to %s",
|
||||||
name, value, ConfigHandler.CMDLINE_OPTIONS[name][0])
|
name, value, ConfigHandler.CMDLINE_OPTIONS[name][0])
|
||||||
|
|||||||
+22
-15
@@ -393,21 +393,28 @@ class ZooKeeper(AbstractDCS):
|
|||||||
cluster = self.cluster
|
cluster = self.cluster
|
||||||
member = cluster and cluster.get_member(self._name, fallback_to_leader=False)
|
member = cluster and cluster.get_member(self._name, fallback_to_leader=False)
|
||||||
member_data = self.__last_member_data or member and member.data
|
member_data = self.__last_member_data or member and member.data
|
||||||
# We want to notify leader if some important fields in the member key changed by removing ZNode
|
if member and member_data:
|
||||||
if member and (self._client.client_id is not None and member.session != self._client.client_id[0]
|
is_leader = data.get('role') in ('master', 'primary', 'standby_leader')
|
||||||
or not (member_data and deep_compare(member_data.get('tags', {}), data.get('tags', {}))
|
checkpoint_after_promote_changed = member_data.get('checkpoint_after_promote') \
|
||||||
and (member_data.get('state') == data.get('state')
|
!= data.get('checkpoint_after_promote')
|
||||||
or 'running' not in (member_data.get('state'), data.get('state')))
|
state_running_changed = member_data.get('state') != data.get('state') \
|
||||||
and member_data.get('version') == data.get('version')
|
and 'running' in (member_data.get('state'), data.get('state'))
|
||||||
and member_data.get('checkpoint_after_promote')
|
tags_changed = not deep_compare(member_data.get('tags', {}), data.get('tags', {}))
|
||||||
== data.get('checkpoint_after_promote'))):
|
|
||||||
try:
|
# We want delete the member ZNode if:
|
||||||
self._client.delete_async(self.member_path).get(timeout=1)
|
# - our session doesn't match with session id on our member key; or
|
||||||
except NoNodeError:
|
# - we want to notify leader if some important fields in the member key changed; or
|
||||||
pass
|
# - if we are the leader and want to notify replicas about checkpoint_after_promote;
|
||||||
except Exception:
|
if self._client.client_id is not None and member.session != self._client.client_id[0] \
|
||||||
return False
|
or is_leader and checkpoint_after_promote_changed \
|
||||||
member = None
|
or not is_leader and (state_running_changed or tags_changed):
|
||||||
|
try:
|
||||||
|
self._client.delete_async(self.member_path).get(timeout=1)
|
||||||
|
except NoNodeError:
|
||||||
|
pass
|
||||||
|
except Exception:
|
||||||
|
return False
|
||||||
|
member = None
|
||||||
|
|
||||||
encoded_data = json.dumps(data, separators=(',', ':')).encode('utf-8')
|
encoded_data = json.dumps(data, separators=(',', ':')).encode('utf-8')
|
||||||
if member and member_data:
|
if member and member_data:
|
||||||
|
|||||||
@@ -654,6 +654,14 @@ class Ha(object):
|
|||||||
current = CaseInsensitiveSet(sync.members)
|
current = CaseInsensitiveSet(sync.members)
|
||||||
picked, allow_promote = self.state_handler.sync_handler.current_state(self.cluster)
|
picked, allow_promote = self.state_handler.sync_handler.current_state(self.cluster)
|
||||||
|
|
||||||
|
if picked == current and current != allow_promote:
|
||||||
|
logger.warning('Inconsistent state between synchronous_standby_names = %s and /sync = %s key '
|
||||||
|
'detected, updating synchronous replication key...', list(allow_promote), list(current))
|
||||||
|
sync = self.dcs.write_sync_state(self.state_handler.name, allow_promote, version=sync.version)
|
||||||
|
if not sync:
|
||||||
|
return logger.warning("Updating sync state failed")
|
||||||
|
current = CaseInsensitiveSet(sync.members)
|
||||||
|
|
||||||
if picked != current:
|
if picked != current:
|
||||||
# update synchronous standby list in dcs temporarily to point to common nodes in current and picked
|
# update synchronous standby list in dcs temporarily to point to common nodes in current and picked
|
||||||
sync_common = current & allow_promote
|
sync_common = current & allow_promote
|
||||||
|
|||||||
@@ -118,17 +118,28 @@ class Postgresql(object):
|
|||||||
# Last known running process
|
# Last known running process
|
||||||
self._postmaster_proc = None
|
self._postmaster_proc = None
|
||||||
|
|
||||||
if self.is_running(): # we are "joining" already running postgres
|
if self.is_running():
|
||||||
self.set_state('running')
|
# If we found postmaster process we need to figure out whether postgres is accepting connections
|
||||||
|
self.set_state('starting')
|
||||||
|
self.check_startup_state_changed()
|
||||||
|
|
||||||
|
if self.state == 'running': # we are "joining" already running postgres
|
||||||
|
# we know that PostgreSQL is accepting connections and can read some GUC's from pg_settings
|
||||||
|
self.config.load_current_server_parameters()
|
||||||
|
|
||||||
self.set_role('master' if self.is_leader() else 'replica')
|
self.set_role('master' if self.is_leader() else 'replica')
|
||||||
# postpone writing postgresql.conf for 12+ because recovery parameters are not yet known
|
|
||||||
if self.major_version < 120000 or self.is_leader():
|
|
||||||
self.config.write_postgresql_conf()
|
|
||||||
hba_saved = self.config.replace_pg_hba()
|
hba_saved = self.config.replace_pg_hba()
|
||||||
ident_saved = self.config.replace_pg_ident()
|
ident_saved = self.config.replace_pg_ident()
|
||||||
if hba_saved or ident_saved:
|
|
||||||
|
if self.major_version < 120000 or self.role in ('master', 'primary'):
|
||||||
|
# If PostgreSQL is running as a primary or we run PostgreSQL that is older than 12 we can
|
||||||
|
# call reload_config() once again (the first call happened in the ConfigHandler constructor),
|
||||||
|
# so that it can figure out if config files should be updated and pg_ctl reload executed.
|
||||||
|
self.config.reload_config(config, sighup=bool(hba_saved or ident_saved))
|
||||||
|
elif hba_saved or ident_saved:
|
||||||
self.reload()
|
self.reload()
|
||||||
elif self.role in ('master', 'primary'):
|
elif not self.is_running() and self.role in ('master', 'primary'):
|
||||||
self.set_role('demoted')
|
self.set_role('demoted')
|
||||||
|
|
||||||
@property
|
@property
|
||||||
|
|||||||
@@ -326,14 +326,22 @@ class ConfigHandler(object):
|
|||||||
.format(self._pgpass))
|
.format(self._pgpass))
|
||||||
self._passfile = None
|
self._passfile = None
|
||||||
self._passfile_mtime = None
|
self._passfile_mtime = None
|
||||||
self._synchronous_standby_names = None
|
|
||||||
self._postmaster_ctime = None
|
self._postmaster_ctime = None
|
||||||
self._current_recovery_params: Optional[CaseInsensitiveDict] = None
|
self._current_recovery_params: Optional[CaseInsensitiveDict] = None
|
||||||
self._config = {}
|
self._config = {}
|
||||||
self._recovery_params = CaseInsensitiveDict()
|
self._recovery_params = CaseInsensitiveDict()
|
||||||
self._server_parameters: CaseInsensitiveDict
|
self._server_parameters: CaseInsensitiveDict = CaseInsensitiveDict()
|
||||||
self.reload_config(config)
|
self.reload_config(config)
|
||||||
|
|
||||||
|
def load_current_server_parameters(self) -> None:
|
||||||
|
"""Read GUC's values from ``pg_settings`` when Patroni is joining the the postgres that is already running."""
|
||||||
|
exclude = [name.lower() for name, value in self.CMDLINE_OPTIONS.items() if value[1] == _false_validator] \
|
||||||
|
+ [name.lower() for name in self._RECOVERY_PARAMETERS]
|
||||||
|
self._server_parameters = CaseInsensitiveDict({r[0]: r[1] for r in self._postgresql.query(
|
||||||
|
"SELECT name, pg_catalog.current_setting(name) FROM pg_catalog.pg_settings"
|
||||||
|
" WHERE (source IN ('command line', 'environment variable') OR sourcefile = %s)"
|
||||||
|
" AND pg_catalog.lower(name) != ALL(%s)", self._postgresql_conf, exclude)})
|
||||||
|
|
||||||
def setup_server_parameters(self) -> None:
|
def setup_server_parameters(self) -> None:
|
||||||
self._server_parameters = self.get_server_parameters(self._config)
|
self._server_parameters = self.get_server_parameters(self._config)
|
||||||
self._adjust_recovery_parameters()
|
self._adjust_recovery_parameters()
|
||||||
@@ -922,14 +930,15 @@ class ConfigHandler(object):
|
|||||||
listen_addresses, port = split_host_port(config['listen'], 5432)
|
listen_addresses, port = split_host_port(config['listen'], 5432)
|
||||||
parameters.update(cluster_name=self._postgresql.scope, listen_addresses=listen_addresses, port=str(port))
|
parameters.update(cluster_name=self._postgresql.scope, listen_addresses=listen_addresses, port=str(port))
|
||||||
if not self._postgresql.global_config or self._postgresql.global_config.is_synchronous_mode:
|
if not self._postgresql.global_config or self._postgresql.global_config.is_synchronous_mode:
|
||||||
if self._synchronous_standby_names is None:
|
synchronous_standby_names = self._server_parameters.get('synchronous_standby_names')
|
||||||
|
if synchronous_standby_names is None:
|
||||||
if self._postgresql.global_config and self._postgresql.global_config.is_synchronous_mode_strict\
|
if self._postgresql.global_config and self._postgresql.global_config.is_synchronous_mode_strict\
|
||||||
and self._postgresql.role in ('master', 'primary', 'promoted'):
|
and self._postgresql.role in ('master', 'primary', 'promoted'):
|
||||||
parameters['synchronous_standby_names'] = '*'
|
parameters['synchronous_standby_names'] = '*'
|
||||||
else:
|
else:
|
||||||
parameters.pop('synchronous_standby_names', None)
|
parameters.pop('synchronous_standby_names', None)
|
||||||
else:
|
else:
|
||||||
parameters['synchronous_standby_names'] = self._synchronous_standby_names
|
parameters['synchronous_standby_names'] = synchronous_standby_names
|
||||||
|
|
||||||
# Handle hot_standby <-> replica rename
|
# Handle hot_standby <-> replica rename
|
||||||
if parameters.get('wal_level') == ('hot_standby' if self._postgresql.major_version >= 90600 else 'replica'):
|
if parameters.get('wal_level') == ('hot_standby' if self._postgresql.major_version >= 90600 else 'replica'):
|
||||||
@@ -1129,12 +1138,11 @@ class ConfigHandler(object):
|
|||||||
def set_synchronous_standby_names(self, value: Optional[str]) -> Optional[bool]:
|
def set_synchronous_standby_names(self, value: Optional[str]) -> Optional[bool]:
|
||||||
"""Updates synchronous_standby_names and reloads if necessary.
|
"""Updates synchronous_standby_names and reloads if necessary.
|
||||||
:returns: True if value was updated."""
|
:returns: True if value was updated."""
|
||||||
if value != self._synchronous_standby_names:
|
if value != self._server_parameters.get('synchronous_standby_names'):
|
||||||
if value is None:
|
if value is None:
|
||||||
self._server_parameters.pop('synchronous_standby_names', None)
|
self._server_parameters.pop('synchronous_standby_names', None)
|
||||||
else:
|
else:
|
||||||
self._server_parameters['synchronous_standby_names'] = value
|
self._server_parameters['synchronous_standby_names'] = value
|
||||||
self._synchronous_standby_names = value
|
|
||||||
if self._postgresql.state == 'running':
|
if self._postgresql.state == 'running':
|
||||||
self.write_postgresql_conf()
|
self.write_postgresql_conf()
|
||||||
self._postgresql.reload()
|
self._postgresql.reload()
|
||||||
|
|||||||
+1
-1
@@ -2,4 +2,4 @@
|
|||||||
|
|
||||||
:var __version__: the current Patroni version.
|
:var __version__: the current Patroni version.
|
||||||
"""
|
"""
|
||||||
__version__ = '3.1.1'
|
__version__ = '3.1.2'
|
||||||
|
|||||||
@@ -118,11 +118,31 @@ class MockCursor(object):
|
|||||||
'"state":"streaming","sync_state":"async","sync_priority":0}]'
|
'"state":"streaming","sync_state":"async","sync_priority":0}]'
|
||||||
now = datetime.datetime.now(tzutc)
|
now = datetime.datetime.now(tzutc)
|
||||||
self.results = [(now, 0, '', 0, '', False, now, 'streaming', None, replication_info)]
|
self.results = [(now, 0, '', 0, '', False, now, 'streaming', None, replication_info)]
|
||||||
|
elif sql.startswith('SELECT name, pg_catalog.current_setting(name) FROM pg_catalog.pg_settings'):
|
||||||
|
self.results = [('data_directory', 'data'),
|
||||||
|
('hba_file', os.path.join('data', 'pg_hba.conf')),
|
||||||
|
('ident_file', os.path.join('data', 'pg_ident.conf')),
|
||||||
|
('max_connections', 42),
|
||||||
|
('max_locks_per_transaction', 73),
|
||||||
|
('max_prepared_transactions', 0),
|
||||||
|
('max_replication_slots', 21),
|
||||||
|
('max_wal_senders', 37),
|
||||||
|
('track_commit_timestamp', 'off'),
|
||||||
|
('wal_level', 'replica'),
|
||||||
|
('listen_addresses', '6.6.6.6'),
|
||||||
|
('port', 1984),
|
||||||
|
('archive_command', 'my archive command'),
|
||||||
|
('cluster_name', 'my_cluster')]
|
||||||
elif sql.startswith('SELECT name, setting'):
|
elif sql.startswith('SELECT name, setting'):
|
||||||
self.results = [('wal_segment_size', '2048', '8kB', 'integer', 'internal'),
|
self.results = [('wal_segment_size', '2048', '8kB', 'integer', 'internal'),
|
||||||
('wal_block_size', '8192', None, 'integer', 'internal'),
|
('wal_block_size', '8192', None, 'integer', 'internal'),
|
||||||
('shared_buffers', '16384', '8kB', 'integer', 'postmaster'),
|
('shared_buffers', '16384', '8kB', 'integer', 'postmaster'),
|
||||||
('wal_buffers', '-1', '8kB', 'integer', 'postmaster'),
|
('wal_buffers', '-1', '8kB', 'integer', 'postmaster'),
|
||||||
|
('max_connections', '100', None, 'integer', 'postmaster'),
|
||||||
|
('max_prepared_transactions', '0', None, 'integer', 'postmaster'),
|
||||||
|
('max_worker_processes', '8', None, 'integer', 'postmaster'),
|
||||||
|
('max_locks_per_transaction', '64', None, 'integer', 'postmaster'),
|
||||||
|
('max_wal_senders', '5', None, 'integer', 'postmaster'),
|
||||||
('search_path', 'public', None, 'string', 'user'),
|
('search_path', 'public', None, 'string', 'user'),
|
||||||
('port', '5433', None, 'integer', 'postmaster'),
|
('port', '5433', None, 'integer', 'postmaster'),
|
||||||
('listen_addresses', '*', None, 'string', 'postmaster'),
|
('listen_addresses', '*', None, 'string', 'postmaster'),
|
||||||
@@ -222,6 +242,7 @@ class PostgresInit(unittest.TestCase):
|
|||||||
|
|
||||||
class BaseTestPostgresql(PostgresInit):
|
class BaseTestPostgresql(PostgresInit):
|
||||||
|
|
||||||
|
@patch('time.sleep', Mock())
|
||||||
def setUp(self):
|
def setUp(self):
|
||||||
super(BaseTestPostgresql, self).setUp()
|
super(BaseTestPostgresql, self).setUp()
|
||||||
|
|
||||||
|
|||||||
@@ -155,6 +155,7 @@ class TestConfig(unittest.TestCase):
|
|||||||
expected_params = {
|
expected_params = {
|
||||||
'f.oo': 'bar', # not in ConfigHandler.CMDLINE_OPTIONS
|
'f.oo': 'bar', # not in ConfigHandler.CMDLINE_OPTIONS
|
||||||
'max_connections': 100, # IntValidator
|
'max_connections': 100, # IntValidator
|
||||||
|
'wal_keep_size': '128MB', # IntValidator
|
||||||
'wal_level': 'hot_standby', # EnumValidator
|
'wal_level': 'hot_standby', # EnumValidator
|
||||||
}
|
}
|
||||||
input_params = deepcopy(expected_params)
|
input_params = deepcopy(expected_params)
|
||||||
|
|||||||
@@ -1313,6 +1313,24 @@ class TestHa(PostgresInit):
|
|||||||
self.ha.run_cycle()
|
self.ha.run_cycle()
|
||||||
self.assertEqual(mock_logger.call_args[0][0], 'Updating sync state failed')
|
self.assertEqual(mock_logger.call_args[0][0], 'Updating sync state failed')
|
||||||
|
|
||||||
|
@patch.object(Cluster, 'is_unlocked', Mock(return_value=False))
|
||||||
|
def test_inconsistent_synchronous_state(self):
|
||||||
|
self.ha.is_synchronous_mode = true
|
||||||
|
self.ha.has_lock = true
|
||||||
|
self.p.name = 'leader'
|
||||||
|
self.ha.cluster = get_cluster_initialized_without_leader(sync=('leader', 'a'))
|
||||||
|
self.p.sync_handler.current_state = Mock(return_value=(CaseInsensitiveSet('a'), CaseInsensitiveSet()))
|
||||||
|
self.ha.dcs.write_sync_state = Mock(return_value=SyncState.empty())
|
||||||
|
mock_set_sync = self.p.sync_handler.set_synchronous_standby_names = Mock()
|
||||||
|
with patch('patroni.ha.logger.warning') as mock_logger:
|
||||||
|
self.ha.run_cycle()
|
||||||
|
mock_set_sync.assert_called_once()
|
||||||
|
self.assertTrue(mock_logger.call_args_list[0][0][0].startswith('Inconsistent state between '))
|
||||||
|
self.ha.dcs.write_sync_state = Mock(return_value=None)
|
||||||
|
with patch('patroni.ha.logger.warning') as mock_logger:
|
||||||
|
self.ha.run_cycle()
|
||||||
|
self.assertEqual(mock_logger.call_args[0][0], 'Updating sync state failed')
|
||||||
|
|
||||||
def test_effective_tags(self):
|
def test_effective_tags(self):
|
||||||
self.ha._disable_sync = True
|
self.ha._disable_sync = True
|
||||||
self.assertEqual(self.ha.get_effective_tags(), {'foo': 'bar', 'nosync': True})
|
self.assertEqual(self.ha.get_effective_tags(), {'foo': 'bar', 'nosync': True})
|
||||||
|
|||||||
@@ -40,7 +40,7 @@ class MockFrozenImporter(object):
|
|||||||
@patch('time.sleep', Mock())
|
@patch('time.sleep', Mock())
|
||||||
@patch('subprocess.call', Mock(return_value=0))
|
@patch('subprocess.call', Mock(return_value=0))
|
||||||
@patch('patroni.psycopg.connect', psycopg_connect)
|
@patch('patroni.psycopg.connect', psycopg_connect)
|
||||||
@patch('urllib3.PoolManager.request', Mock(side_effect=Exception))
|
@patch('urllib3.connection.HTTPConnection.connect', Mock(side_effect=Exception))
|
||||||
@patch.object(ConfigHandler, 'append_pg_hba', Mock())
|
@patch.object(ConfigHandler, 'append_pg_hba', Mock())
|
||||||
@patch.object(ConfigHandler, 'write_postgresql_conf', Mock())
|
@patch.object(ConfigHandler, 'write_postgresql_conf', Mock())
|
||||||
@patch.object(ConfigHandler, 'write_recovery_conf', Mock())
|
@patch.object(ConfigHandler, 'write_recovery_conf', Mock())
|
||||||
@@ -64,7 +64,7 @@ class TestPatroni(unittest.TestCase):
|
|||||||
self.assertRaises(SystemExit, _main)
|
self.assertRaises(SystemExit, _main)
|
||||||
|
|
||||||
@patch('pkgutil.iter_importers', Mock(return_value=[MockFrozenImporter()]))
|
@patch('pkgutil.iter_importers', Mock(return_value=[MockFrozenImporter()]))
|
||||||
@patch('urllib3.PoolManager.request', Mock(side_effect=Exception))
|
@patch('urllib3.connection.HTTPConnection.connect', Mock(side_effect=Exception))
|
||||||
@patch('sys.frozen', Mock(return_value=True), create=True)
|
@patch('sys.frozen', Mock(return_value=True), create=True)
|
||||||
@patch.object(HTTPServer, '__init__', Mock())
|
@patch.object(HTTPServer, '__init__', Mock())
|
||||||
@patch.object(etcd.Client, 'read', etcd_read)
|
@patch.object(etcd.Client, 'read', etcd_read)
|
||||||
@@ -108,6 +108,7 @@ class TestPatroni(unittest.TestCase):
|
|||||||
@patch('os.getpid')
|
@patch('os.getpid')
|
||||||
@patch('multiprocessing.Process')
|
@patch('multiprocessing.Process')
|
||||||
@patch('patroni.__main__.patroni_main', Mock())
|
@patch('patroni.__main__.patroni_main', Mock())
|
||||||
|
@patch('sys.argv', ['patroni.py', 'postgres0.yml'])
|
||||||
def test_patroni_main(self, mock_process, mock_getpid):
|
def test_patroni_main(self, mock_process, mock_getpid):
|
||||||
mock_getpid.return_value = 2
|
mock_getpid.return_value = 2
|
||||||
_main()
|
_main()
|
||||||
@@ -233,8 +234,8 @@ class TestPatroni(unittest.TestCase):
|
|||||||
)
|
)
|
||||||
with patch('patroni.dcs.AbstractDCS.get_cluster', Mock(return_value=bad_cluster)):
|
with patch('patroni.dcs.AbstractDCS.get_cluster', Mock(return_value=bad_cluster)):
|
||||||
# If the api of the running node cannot be reached, this implies unique name
|
# If the api of the running node cannot be reached, this implies unique name
|
||||||
with patch.object(self.p, 'request', Mock(side_effect=ConnectionError)):
|
with patch('urllib3.connection.HTTPConnection.connect', Mock(side_effect=ConnectionError)):
|
||||||
self.assertIsNone(self.p.ensure_unique_name())
|
self.assertIsNone(self.p.ensure_unique_name())
|
||||||
# Only if the api of the running node is reachable do we throw an error
|
# Only if the api of the running node is reachable do we throw an error
|
||||||
with patch.object(self.p, 'request', Mock()):
|
with patch('urllib3.connection.HTTPConnection.connect', Mock()):
|
||||||
self.assertRaises(SystemExit, self.p.ensure_unique_name)
|
self.assertRaises(SystemExit, self.p.ensure_unique_name)
|
||||||
|
|||||||
@@ -715,6 +715,7 @@ class TestPostgresql(BaseTestPostgresql):
|
|||||||
self.assertEqual(self.p.get_primary_timeline(), 1)
|
self.assertEqual(self.p.get_primary_timeline(), 1)
|
||||||
|
|
||||||
@patch.object(Postgresql, 'get_postgres_role_from_data_directory', Mock(return_value='replica'))
|
@patch.object(Postgresql, 'get_postgres_role_from_data_directory', Mock(return_value='replica'))
|
||||||
|
@patch.object(Postgresql, 'is_running', Mock(return_value=False))
|
||||||
@patch.object(Bootstrap, 'running_custom_bootstrap', PropertyMock(return_value=True))
|
@patch.object(Bootstrap, 'running_custom_bootstrap', PropertyMock(return_value=True))
|
||||||
@patch.object(Postgresql, 'controldata', Mock(return_value={'max_connections setting': '200',
|
@patch.object(Postgresql, 'controldata', Mock(return_value={'max_connections setting': '200',
|
||||||
'max_worker_processes setting': '20',
|
'max_worker_processes setting': '20',
|
||||||
@@ -958,6 +959,7 @@ class TestPostgresql2(BaseTestPostgresql):
|
|||||||
@patch('patroni.postgresql.CallbackExecutor', Mock())
|
@patch('patroni.postgresql.CallbackExecutor', Mock())
|
||||||
@patch.object(Postgresql, 'get_major_version', Mock(return_value=140000))
|
@patch.object(Postgresql, 'get_major_version', Mock(return_value=140000))
|
||||||
@patch.object(Postgresql, 'is_running', Mock(return_value=True))
|
@patch.object(Postgresql, 'is_running', Mock(return_value=True))
|
||||||
|
@patch.object(Postgresql, 'is_leader', Mock(return_value=False))
|
||||||
def setUp(self):
|
def setUp(self):
|
||||||
super(TestPostgresql2, self).setUp()
|
super(TestPostgresql2, self).setUp()
|
||||||
|
|
||||||
|
|||||||
@@ -276,6 +276,7 @@ class TestZooKeeper(unittest.TestCase):
|
|||||||
self.assertTrue(self.zk.delete_cluster())
|
self.assertTrue(self.zk.delete_cluster())
|
||||||
|
|
||||||
def test_watch(self):
|
def test_watch(self):
|
||||||
|
self.zk.event.wait = Mock()
|
||||||
self.zk.watch(None, 0)
|
self.zk.watch(None, 0)
|
||||||
self.zk.event.is_set = Mock(return_value=True)
|
self.zk.event.is_set = Mock(return_value=True)
|
||||||
self.zk._fetch_status = False
|
self.zk._fetch_status = False
|
||||||
|
|||||||
Reference in New Issue
Block a user