mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Track postgresql role in a Postgresql class
This commit is contained in:
+29
-47
@@ -55,6 +55,7 @@ class Postgresql:
|
||||
self.postmaster_pid = os.path.join(self.data_dir, 'postmaster.pid')
|
||||
self.trigger_file = config.get('recovery_conf', {}).get('trigger_file', None) or 'promote'
|
||||
self.trigger_file = os.path.abspath(os.path.join(self.data_dir, self.trigger_file))
|
||||
self._role = 'replica'
|
||||
self.is_promoted = False
|
||||
|
||||
self._pg_ctl = ['pg_ctl', '-w', '-D', self.data_dir]
|
||||
@@ -163,31 +164,26 @@ class Postgresql:
|
||||
def is_running(self):
|
||||
return subprocess.call(' '.join(self._pg_ctl) + ' status > /dev/null 2>&1', shell=True) == 0
|
||||
|
||||
def call_nowait(self, cb_name, is_leader=None):
|
||||
def call_nowait(self, cb_name):
|
||||
""" pick a callback command and call it without waiting for it to finish """
|
||||
if not self.callback or cb_name not in self.callback:
|
||||
return False
|
||||
cmd = self.callback[cb_name]
|
||||
if is_leader is None:
|
||||
try:
|
||||
is_leader = self.is_leader(check_only=True)
|
||||
except psycopg2.OperationalError as e:
|
||||
logger.warning("unable to perform {0} action, cannot obtain the cluster role: {1}".format(cb_name, e))
|
||||
return False
|
||||
try:
|
||||
role = "master" if is_leader else "replica"
|
||||
subprocess.Popen(shlex.split(cmd) + [cb_name, role, self.scope])
|
||||
subprocess.Popen(shlex.split(cmd) + [cb_name, self._role, self.scope])
|
||||
except:
|
||||
logger.exception('callback %s %s %s %s failed', cmd, cb_name, role, self.scope)
|
||||
logger.exception('callback %s %s %s %s failed', cmd, cb_name, self._role, self.scope)
|
||||
return False
|
||||
return True
|
||||
|
||||
def start(self, block_callbacks=False):
|
||||
def start(self, role='replica', block_callbacks=False):
|
||||
if self.is_running():
|
||||
self._role = 'master' if self.is_leader(check_only=True) else 'replica'
|
||||
self.load_replication_slots()
|
||||
logger.error('Cannot start PostgreSQL because one is already running.')
|
||||
return False
|
||||
|
||||
self._role = role
|
||||
if os.path.exists(self.postmaster_pid):
|
||||
os.remove(self.postmaster_pid)
|
||||
logger.info('Removed %s', self.postmaster_pid)
|
||||
@@ -197,47 +193,32 @@ class Postgresql:
|
||||
self.save_configuration_files()
|
||||
# block_callbacks is used during restart to avoid
|
||||
# running start/stop callbacks in addition to restart ones
|
||||
if not block_callbacks and ret and ACTION_ON_START in self.callback:
|
||||
self.call_nowait(ACTION_ON_START)
|
||||
ret and not block_callbacks and ret and self.call_nowait(ACTION_ON_START)
|
||||
return ret
|
||||
|
||||
def stop(self, block_callbacks=False):
|
||||
try:
|
||||
is_leader = self.is_leader(check_only=True)
|
||||
except:
|
||||
is_leader = None
|
||||
pass
|
||||
try:
|
||||
self.query("CHECKPOINT")
|
||||
except psycopg2.OperationalError:
|
||||
# likely PostgreSQL is already stopped.
|
||||
ret = 0
|
||||
else:
|
||||
ret = subprocess.call(self._pg_ctl + ['stop', '-m', 'fast'])
|
||||
if block_callbacks:
|
||||
try:
|
||||
self.query('SET statement_timeout TO 0')
|
||||
self.query('CHECKPOINT')
|
||||
except:
|
||||
logging.exception('Exception diring CHECKPOINT')
|
||||
|
||||
ret = subprocess.call(self._pg_ctl + ['stop', '-m', 'fast']) == 0
|
||||
# block_callbacks is used during restart to avoid
|
||||
# running start/stop callbacks in addition to restart ones
|
||||
if not block_callbacks and ret == 0 and ACTION_ON_STOP in self.callback:
|
||||
self.call_nowait(ACTION_ON_STOP, is_leader=is_leader)
|
||||
return ret == 0
|
||||
ret and not block_callbacks and self.call_nowait(ACTION_ON_STOP)
|
||||
return ret
|
||||
|
||||
def reload(self):
|
||||
ret = subprocess.call(self._pg_ctl + ['reload'])
|
||||
if ret == 0 and ACTION_ON_RELOAD in self.callback:
|
||||
self.call_nowait(ACTION_ON_RELOAD)
|
||||
return ret == 0
|
||||
ret = subprocess.call(self._pg_ctl + ['reload']) == 0
|
||||
ret and self.call_nowait(ACTION_ON_RELOAD)
|
||||
return ret
|
||||
|
||||
def restart(self):
|
||||
try:
|
||||
is_leader = self.is_leader(check_only=True)
|
||||
except:
|
||||
is_leader = None
|
||||
pass
|
||||
ret = self.stop(block_callbacks=True)
|
||||
if ret == 0:
|
||||
ret = self.start(block_callbacks=True)
|
||||
if ret == 0 and ACTION_ON_RESTART in self.callback:
|
||||
self.call_nowait(ACTION_ON_RESTART, is_leader=is_leader)
|
||||
return ret == 0
|
||||
ret = self.stop(block_callbacks=True) and self.start(block_callbacks=True)
|
||||
ret and self.call_nowait(ACTION_ON_RESTART)
|
||||
return ret
|
||||
|
||||
def server_options(self):
|
||||
options = "--listen_addresses='{}' --port={}".format(self.listen_addresses, self.port)
|
||||
@@ -325,9 +306,9 @@ recovery_target_timeline = 'latest'
|
||||
def follow_the_leader(self, leader):
|
||||
if not self.check_recovery_conf(leader):
|
||||
self.write_recovery_conf(leader)
|
||||
run_callback = self._role == 'master'
|
||||
self.restart()
|
||||
if ACTION_ON_ROLE_CHANGE in self.callback:
|
||||
self.call_nowait(ACTION_ON_ROLE_CHANGE)
|
||||
run_callback and self.call_nowait(ACTION_ON_ROLE_CHANGE)
|
||||
|
||||
def save_configuration_files(self):
|
||||
"""
|
||||
@@ -347,7 +328,8 @@ recovery_target_timeline = 'latest'
|
||||
|
||||
def promote(self):
|
||||
self.is_promoted = subprocess.call(self._pg_ctl + ['promote']) == 0
|
||||
if self.is_promoted and ACTION_ON_ROLE_CHANGE in self.callback:
|
||||
if self.is_promoted:
|
||||
self._role = 'master'
|
||||
self.call_nowait(ACTION_ON_ROLE_CHANGE)
|
||||
return self.is_promoted
|
||||
|
||||
@@ -418,7 +400,7 @@ recovery_target_timeline = 'latest'
|
||||
"""
|
||||
ret = False
|
||||
if not current_leader:
|
||||
ret = self.initialize() and self.start()
|
||||
ret = self.initialize() and self.start(role='master')
|
||||
if ret:
|
||||
self.create_replication_user()
|
||||
self.create_connection_users()
|
||||
|
||||
@@ -61,7 +61,7 @@ def get_cluster_not_initialized_with_leader():
|
||||
|
||||
|
||||
def get_cluster_initialized_with_leader():
|
||||
return get_cluster(True, Leader(0, 0, 0,
|
||||
return get_cluster(True, Leader(0, 0, 0,
|
||||
Member(0, 'leader', 'postgres://replicator:[email protected]:5435/postgres',
|
||||
None, None, 28)))
|
||||
|
||||
@@ -132,7 +132,7 @@ class TestPatroni(unittest.TestCase):
|
||||
self.p.ha.dcs.client.read = etcd_read
|
||||
self.p.ha.dcs.watch = time_sleep
|
||||
self.assertRaises(SleepException, self.p.run)
|
||||
self.p.ha.state_handler.is_leader = lambda: False
|
||||
self.p.ha.state_handler.is_leader = false
|
||||
self.p.api.start = nop
|
||||
self.assertRaises(SleepException, self.p.run)
|
||||
|
||||
|
||||
+18
-11
@@ -17,10 +17,6 @@ def subprocess_call(cmd, shell=False, env=None):
|
||||
return 0
|
||||
|
||||
|
||||
def false(*args, **kwargs):
|
||||
return False
|
||||
|
||||
|
||||
class MockCursor:
|
||||
|
||||
def __init__(self):
|
||||
@@ -28,7 +24,7 @@ class MockCursor:
|
||||
self.results = []
|
||||
|
||||
def execute(self, sql, *params):
|
||||
if sql.startswith('blabla'):
|
||||
if sql.startswith('blabla') or sql == 'CHECKPOINT':
|
||||
raise psycopg2.OperationalError()
|
||||
elif sql.startswith('InterfaceError'):
|
||||
raise psycopg2.InterfaceError()
|
||||
@@ -88,12 +84,11 @@ class MockConnect:
|
||||
|
||||
|
||||
def psycopg2_connect(*args, **kwargs):
|
||||
|
||||
return MockConnect()
|
||||
|
||||
|
||||
def is_running():
|
||||
return False
|
||||
def raise_exception(*args, **kwargs):
|
||||
raise Exception
|
||||
|
||||
|
||||
class TestPostgresql(unittest.TestCase):
|
||||
@@ -143,7 +138,7 @@ class TestPostgresql(unittest.TestCase):
|
||||
|
||||
def test_start_stop(self):
|
||||
self.assertFalse(self.p.start())
|
||||
self.p.is_running = is_running
|
||||
self.p.is_running = false
|
||||
with open(os.path.join(self.p.data_dir, 'postmaster.pid'), 'w'):
|
||||
pass
|
||||
self.assertTrue(self.p.start())
|
||||
@@ -159,6 +154,10 @@ class TestPostgresql(unittest.TestCase):
|
||||
self.p.follow_the_leader(self.leader)
|
||||
self.p.follow_the_leader(Leader(-1, None, 28, self.other))
|
||||
|
||||
def test_create_replica(self):
|
||||
self.p.delete_trigger_file = raise_exception
|
||||
self.assertEquals(self.p.create_replica({'host': '', 'port': '', 'user': ''}, ''), 1)
|
||||
|
||||
def test_create_connection_users(self):
|
||||
cfg = self.p.config
|
||||
cfg['superuser']['username'] = 'test'
|
||||
@@ -202,7 +201,7 @@ class TestPostgresql(unittest.TestCase):
|
||||
|
||||
def test_is_healthy(self):
|
||||
self.assertTrue(self.p.is_healthy())
|
||||
self.p.is_running = is_running
|
||||
self.p.is_running = false
|
||||
self.assertFalse(self.p.is_healthy())
|
||||
|
||||
def test_promote(self):
|
||||
@@ -211,6 +210,12 @@ class TestPostgresql(unittest.TestCase):
|
||||
def test_last_operation(self):
|
||||
self.assertEquals(self.p.last_operation(), '0')
|
||||
|
||||
def test_call_nowait(self):
|
||||
popen = subprocess.Popen
|
||||
subprocess.Popen = raise_exception
|
||||
self.assertFalse(self.p.call_nowait('on_start'))
|
||||
subprocess.Popen = popen
|
||||
|
||||
def test_non_existing_callback(self):
|
||||
self.assertFalse(self.p.call_nowait('foobar'))
|
||||
|
||||
@@ -220,7 +225,9 @@ class TestPostgresql(unittest.TestCase):
|
||||
self.assertTrue(self.p.stop())
|
||||
|
||||
def test_move_data_directory(self):
|
||||
self.p.is_running = is_running
|
||||
self.p.is_running = false
|
||||
os.rename = nop
|
||||
os.path.isdir = true
|
||||
self.p.move_data_directory()
|
||||
os.rename = raise_exception
|
||||
self.p.move_data_directory()
|
||||
|
||||
Reference in New Issue
Block a user