Introduce is_paused method in the Cluster

This commit is contained in:
Murat Kabilov
2016-08-29 09:29:49 +02:00
parent 89ef5da5ae
commit 3d1fe3fa49
6 changed files with 29 additions and 40 deletions
+11 -14
View File
@@ -194,13 +194,14 @@ class RestApiHandler(BaseHTTPRequestHandler):
status_code = 500
data = 'restart failed'
request = self._read_json_content(body_is_optional=True)
cluster = self.server.patroni.dcs.get_cluster()
if request is None:
# failed to parse the json
return
if request:
logger.debug("received restart request: {0}".format(request))
if self.server.patroni.ha.is_paused() and 'schedule' in request and self:
if cluster.is_paused() and 'schedule' in request:
self._write_response(status_code, "Can't schedule restart in the paused state")
return
@@ -244,16 +245,12 @@ class RestApiHandler(BaseHTTPRequestHandler):
@check_auth
def do_DELETE_restart(self):
if self.server.patroni.ha.is_paused():
data = "Can't delete scheduled restart in the paused state"
code = 500
if self.server.patroni.ha.delete_future_restart():
data = "scheduled restart deleted"
code = 200
else:
if self.server.patroni.ha.delete_future_restart():
data = "scheduled restart deleted"
code = 200
else:
data = "no restarts are scheduled"
code = 404
data = "no restarts are scheduled"
code = 404
self._write_response(code, data)
@check_auth
@@ -261,7 +258,7 @@ class RestApiHandler(BaseHTTPRequestHandler):
patroni = self.server.patroni
cluster = patroni.dcs.get_cluster()
status_code = 500
if self.server.patroni.ha.is_paused():
if cluster.is_paused():
self._write_response(status_code, "Can't do reinitialize in the paused state")
return
@@ -324,11 +321,11 @@ class RestApiHandler(BaseHTTPRequestHandler):
leader = request.get('leader')
candidate = request.get('candidate') or request.get('member')
scheduled_at = request.get('scheduled_at')
if scheduled_at and self.server.patroni.ha.is_paused():
self._write_response(status_code, "Can't schedule failover in the paused state")
cluster = self.server.patroni.dcs.get_cluster()
if scheduled_at and cluster.is_paused():
self._write_response(status_code, "Can't schedule failover in the paused state")
logger.info("received failover request with leader=%s candidate=%s scheduled_at=%s",
leader, candidate, scheduled_at)
+9 -16
View File
@@ -499,7 +499,7 @@ def restart(cluster_name, member_names, config_file, dcs, force, role, p_any, sc
scheduled_at = parse_scheduled(scheduled)
if scheduled_at:
if is_paused(cluster):
if cluster.is_paused():
raise PatroniCtlException("Can't schedule restart in the paused state")
content['schedule'] = scheduled_at.isoformat()
@@ -556,17 +556,17 @@ def failover(config_file, cluster_name, master, candidate, force, dcs, scheduled
config, dcs, cluster = ctl_load_config(cluster_name, config_file, dcs)
if cluster.leader is None and not is_paused(cluster):
if cluster.leader is None and not cluster.is_paused():
raise PatroniCtlException('This cluster has no master')
if master is None and (not is_paused(cluster) or cluster.leader):
if master is None and (not cluster.is_paused() or cluster.leader):
if force:
master = cluster.leader.member.name
else:
master = click.prompt('Master', type=str, default=cluster.leader.member.name)
if not is_paused(cluster) and cluster.leader.member.name != master:
raise PatroniCtlException('Member {0} is not the leader of cluster {1}'.format(master, cluster_name))
if not (master is not None and cluster.leader and cluster.leader.member.name == master):
raise PatroniCtlException('Member {0} is not the leader of cluster {1}'.format(master, cluster_name))
candidate_names = [str(m.name) for m in cluster.members if m.name != master]
# We sort the names for consistent output to the client
@@ -591,13 +591,11 @@ def failover(config_file, cluster_name, master, candidate, force, dcs, scheduled
scheduled_at = parse_scheduled(scheduled)
if scheduled_at:
if is_paused(cluster):
if cluster.is_paused():
raise PatroniCtlException("Can't schedule failover in the paused state")
scheduled_at = scheduled_at.isoformat()
failover_value = {'candidate': candidate, 'scheduled_at': scheduled_at}
if master:
failover_value['leader'] = master
failover_value = {'leader': master, 'candidate': candidate, 'scheduled_at': scheduled_at}
logging.debug(failover_value)
@@ -754,11 +752,6 @@ def touch_member(config, dcs):
return dcs.touch_member(json.dumps(data, separators=(',', ':')), permanent=True)
def is_paused(cluster):
"""Check if cluster management is paused"""
return cluster.config and 'pause' in cluster.config.data and cluster.config.data['pause']
def set_defaults(config, cluster_name):
"""fill-in some basic configuration parameters if config file is not set """
config['postgresql'].setdefault('name', cluster_name)
@@ -821,7 +814,7 @@ def flush(cluster_name, member_names, config_file, dcs, force, role, target):
def disable(config_file, cluster_name, dcs):
config, dcs, cluster = ctl_load_config(cluster_name, config_file, dcs)
if is_paused(cluster):
if cluster.is_paused():
raise PatroniCtlException("Cluster is already paused")
r = request_patroni(cluster.leader.member, 'patch', 'config', {'pause': True}, auth_header(config))
@@ -840,7 +833,7 @@ def disable(config_file, cluster_name, dcs):
def resume(config_file, cluster_name, dcs):
config, dcs, cluster = ctl_load_config(cluster_name, config_file, dcs)
if not is_paused(cluster):
if not cluster.is_paused():
raise PatroniCtlException("Cluster is not paused")
r = request_patroni(cluster.leader.member, 'patch', 'config', {'pause': False}, auth_header(config))
+3
View File
@@ -228,6 +228,9 @@ class Cluster(namedtuple('Cluster', 'initialize,config,leader,last_leader_operat
candidates = [m for m in self.members if m.clonefrom and (not self.leader or m.name != self.leader.name)]
return candidates[randint(0, len(candidates) - 1)] if candidates else self.leader
def is_paused(self):
return self.config and self.config.data.get('pause', False)
@six.add_metaclass(abc.ABCMeta)
class AbstractDCS(object):
+1 -1
View File
@@ -27,7 +27,7 @@ class Ha(object):
self._async_executor = AsyncExecutor()
def is_paused(self):
return self.cluster and self.cluster.config and self.cluster.config.data.get('pause', False)
return self.cluster and self.cluster.is_paused()
def load_cluster_from_dcs(self):
cluster = self.dcs.get_cluster()
-3
View File
@@ -234,9 +234,6 @@ class TestRestApiHandler(unittest.TestCase):
with patch.object(MockHa, 'delete_future_restart', Mock(return_value=retval)):
request = 'DELETE /restart HTTP/1.0' + self._authorization
self.assertIsNotNone(MockRestApiServer(RestApiHandler, request))
with patch.object(MockHa, 'is_paused', Mock(return_value=True)):
request = 'DELETE /restart HTTP/1.0' + self._authorization
self.assertIsNotNone(MockRestApiServer(RestApiHandler, request))
@patch.object(MockPatroni, 'dcs')
def test_do_POST_reinitialize(self, dcs):
+5 -6
View File
@@ -82,7 +82,7 @@ class TestCtl(unittest.TestCase):
result = self.runner.invoke(ctl, ['failover', 'dummy'], input='leader\nother\n2030-01-01T12:23:00\ny')
assert result.exit_code == 0
with patch('patroni.ctl.is_paused', Mock(return_value=True)):
with patch('patroni.dcs.Cluster.is_paused', Mock(return_value=True)):
result = self.runner.invoke(ctl,
['failover', 'dummy', '--force', '--scheduled', '2015-01-01T12:00:00+01:00'])
assert result.exit_code == 1
@@ -246,12 +246,11 @@ class TestCtl(unittest.TestCase):
'--scheduled', '2300-10-01T14:30'])
assert 'Failed: flush scheduled restart' in result.output
with patch('patroni.ctl.is_paused', Mock(return_value=True)):
with patch('patroni.dcs.Cluster.is_paused', Mock(return_value=True)):
result = self.runner.invoke(ctl,
['restart', 'alpha', 'other', '--force', '--scheduled', '2300-10-01T14:30'])
assert result.exit_code == 1
with patch('requests.post', Mock(return_value=MockResponse())):
# normal restart, the schedule is actually parsed, but not validated in patronictl
result = self.runner.invoke(ctl, ['restart', 'alpha', '--pg-version', '42.0.0',
@@ -419,7 +418,7 @@ class TestCtl(unittest.TestCase):
assert 'Failed' in result.output
with patch('requests.patch', Mock(return_value=MockResponse(200))),\
patch('patroni.ctl.is_paused', Mock(return_value=True)):
patch('patroni.dcs.Cluster.is_paused', Mock(return_value=True)):
result = self.runner.invoke(ctl, ['disable', 'dummy'])
assert 'Cluster is already paused' in result.output
@@ -428,7 +427,7 @@ class TestCtl(unittest.TestCase):
mock_get_dcs.return_value = self.e
mock_get_dcs.return_value.get_cluster = get_cluster_initialized_with_leader
with patch('patroni.ctl.is_paused', Mock(return_value=True)):
with patch('patroni.dcs.Cluster.is_paused', Mock(return_value=True)):
with patch('requests.patch', Mock(return_value=MockResponse(200))):
result = self.runner.invoke(ctl, ['resume', 'dummy'])
assert 'Success' in result.output
@@ -438,6 +437,6 @@ class TestCtl(unittest.TestCase):
assert 'Failed' in result.output
with patch('requests.patch', Mock(return_value=MockResponse(200))),\
patch('patroni.ctl.is_paused', Mock(return_value=False)):
patch('patroni.dcs.Cluster.is_paused', Mock(return_value=False)):
result = self.runner.invoke(ctl, ['resume', 'dummy'])
assert 'Cluster is not paused' in result.output