From e83651b57b9d8a3672e9526d81c6a8d5aa05f33a Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 23 Sep 2015 10:55:38 +0200 Subject: [PATCH 1/2] Run initial cluster bootstrap from the main loop --- patroni/__init__.py | 35 +--------- patroni/ha.py | 142 +++++++++++++++++++++++++-------------- patroni/postgresql.py | 25 +++++-- tests/test_ha.py | 111 +++++++++++++++++++++--------- tests/test_patroni.py | 102 +++++----------------------- tests/test_postgresql.py | 19 +++++- 6 files changed, 228 insertions(+), 206 deletions(-) diff --git a/patroni/__init__.py b/patroni/__init__.py index 5a8588ce..10e70956 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -6,7 +6,6 @@ import yaml from patroni.api import RestApiServer from patroni.etcd import Etcd -from patroni.exceptions import DCSError from patroni.ha import Ha from patroni.postgresql import Postgresql from patroni.utils import setup_signal_handlers, sleep, reap_children @@ -43,45 +42,13 @@ class Patroni: return True return self.ha.dcs.touch_member(connection_string, ttl) - def cleanup_on_failed_initialization(self): - """ cleanup the DCS if initialization was not successfull """ - logger.info("removing initialize key after failed attempt to initialize the cluster") - self.ha.dcs.cancel_initialization() - self.touch_member(self.shutdown_member_ttl) - self.postgresql.stop() - self.postgresql.move_data_directory() - def initialize(self): # wait for etcd to be available while not self.touch_member(): logger.info('waiting on DCS') sleep(5) - # is data directory empty? - if self.postgresql.data_directory_empty(): - while True: - try: - cluster = self.ha.dcs.get_cluster() - if not cluster.is_unlocked(): # the leader already exists - if not cluster.initialize: - self.ha.dcs.initialize() - self.postgresql.bootstrap(cluster.leader) - break - # racing to initialize - elif not cluster.initialize and self.ha.dcs.initialize(): - try: - self.postgresql.bootstrap() - except: - # bail out and clean the initialize flag. - self.cleanup_on_failed_initialization() - raise - self.ha.dcs.take_leader() - break - except DCSError: - logger.info('waiting on DCS') - sleep(5) - elif self.postgresql.is_running(): - self.postgresql.schedule_load_slots = True + self.postgresql.schedule_load_slots = self.postgresql.is_running() and self.postgresql.use_slots def schedule_next_run(self): self.next_run += self.nap_time diff --git a/patroni/ha.py b/patroni/ha.py index 9ed8faf3..0d95d372 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -35,63 +35,107 @@ class Ha: logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name) return lock_owner == self.state_handler.name - def demote(self): - return self.state_handler.demote(self.cluster.leader) + def bootstrap(self): + if not self.cluster.is_unlocked(): # cluster already has leader + logger.info('trying to bootstrap from leader', ) + if self.state_handler.bootstrap(self.cluster.leader): + return 'bootstrapped from leader' + else: + self.state_handler.stop('immediate') + self.state_handler.remove_data_directory() + return 'failed to bootstrap from leader' + elif not self.cluster.initialize: # no initialize key + if self.dcs.initialize(): # race for initialization + try: + self.state_handler.bootstrap() + except: # initdb or start failed + # remove initialization key and give a chance to other members + logger.info("removing initialize key after failed attempt to initialize the cluster") + self.dcs.cancel_initialization() + self.state_handler.stop('immediate') + self.state_handler.move_data_directory() + raise + self.dcs.take_leader() + return 'initialized a new cluster' + else: + return 'failed to acquire initialize lock' + else: + return 'waiting for leader to bootstrap' - def follow_the_leader(self): - return self.state_handler.follow_the_leader(self.cluster.leader) + def recover(self): + if self.state_handler.is_healthy(): + return False + has_lock = self.has_lock() + self.state_handler.write_recovery_conf(None if has_lock else self.cluster.leader) + self.state_handler.start() + if has_lock: + logger.info('started as readonly because i had the session lock') + self.load_cluster_from_dcs() + return True + + def follow_the_leader(self, demote_reason, follow_reason, refresh=True): + refresh and self.load_cluster_from_dcs() + ret = demote_reason if self.state_handler.is_leader() else follow_reason + self.state_handler.follow_the_leader(self.cluster.leader) + return ret + + def enforce_master_role(self, message, promote_message): + if self.state_handler.is_leader() or self.state_handler.role == 'master': + return message + else: + self.state_handler.promote() + return promote_message + + def process_unhealthy_cluster(self): + if self.state_handler.is_healthiest_node(self.old_cluster): + if self.acquire_lock(): + return self.enforce_master_role('acquired session lock as a leader', + 'promoted self to leader by acquiring session lock') + else: + return self.follow_the_leader('demoted self due after trying and failing to obtain lock', + 'following new leader after trying and failing to obtain lock') + else: + return self.follow_the_leader('demoting self because i am not the healthiest node', + 'following a different leader because i am not the healthiest node') + + def process_healthy_cluster(self): + if self.has_lock(): + if self.update_lock(): + return self.enforce_master_role('no action. i am the leader with the lock', + 'promoted self to leader because i had the session lock') + else: + # Either there is no connection to DCS or someone else acquired the lock + logger.error('failed to update leader lock') + self.load_cluster_from_dcs() + else: + logger.info('does not have lock') + return self.follow_the_leader('demoting self because i do not have the lock and i was a leader', + 'no action. i am a secondary and i am following a leader', False) def run_cycle(self): try: self.load_cluster_from_dcs() - if not self.state_handler.is_healthy(): - has_lock = self.has_lock() - self.state_handler.write_recovery_conf(None if has_lock else self.cluster.leader) - self.state_handler.start() - if not has_lock: - return 'started as a secondary' - logger.info('started as readonly because i had the session lock') - self.load_cluster_from_dcs() + + # cluster has leader key but not initialize key + if not self.cluster.is_unlocked() and not self.cluster.initialize: + self.dcs.initialize() # fix it + + # is data directory empty? + if self.state_handler.data_directory_empty(): + return self.bootstrap() # new node + # "bootstrap", but data directory is not empty + elif not self.cluster.initialize and self.cluster.is_unlocked(): + self.dcs.initialize() + + # try to start dead postgres + if self.recover() and not self.has_lock(): + # no lock, do not try to promote immediately + return 'started as a secondary' if self.cluster.is_unlocked(): - if self.state_handler.is_healthiest_node(self.old_cluster): - if self.acquire_lock(): - if self.state_handler.is_leader() or self.state_handler.role == 'master': - return 'acquired session lock as a leader' - else: - self.state_handler.promote() - return 'promoted self to leader by acquiring session lock' - else: - self.load_cluster_from_dcs() - if self.state_handler.is_leader(): - self.demote() - return 'demoted self due after trying and failing to obtain lock' - else: - self.follow_the_leader() - return 'following new leader after trying and failing to obtain lock' - else: - self.load_cluster_from_dcs() - if self.state_handler.is_leader(): - self.demote() - return 'demoting self because i am not the healthiest node' - else: - self.follow_the_leader() - return 'following a different leader because i am not the healthiest node' + return self.process_unhealthy_cluster() else: - if self.has_lock() and self.update_lock(): - if self.state_handler.is_leader() or self.state_handler.role == 'master': - return 'no action. i am the leader with the lock' - else: - self.state_handler.promote() - return 'promoted self to leader because i had the session lock' - else: - logger.info('does not have lock') - if self.state_handler.is_leader(): - self.demote() - return 'demoting self because i do not have the lock and i was a leader' - else: - self.follow_the_leader() - return 'no action. i am a secondary and i am following a leader' + return self.process_healthy_cluster() except DCSError: logger.error('Error communicating with DCS') if self.state_handler.is_leader(): diff --git a/patroni/postgresql.py b/patroni/postgresql.py index 532e65c8..81f5c02e 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -189,7 +189,7 @@ class Postgresql: ret and not block_callbacks and ret and self.call_nowait(ACTION_ON_START) return ret - def stop(self, block_callbacks=False): + def stop(self, mode='fast', block_callbacks=False): if block_callbacks: try: self.query('SET statement_timeout TO 0') @@ -197,7 +197,7 @@ class Postgresql: except: logging.exception('Exception diring CHECKPOINT') - ret = subprocess.call(self._pg_ctl + ['stop', '-m', 'fast']) == 0 + ret = subprocess.call(self._pg_ctl + ['stop', '-m', mode]) == 0 # block_callbacks is used during restart to avoid # running start/stop callbacks in addition to restart ones ret and not block_callbacks and self.call_nowait(ACTION_ON_STOP) @@ -407,6 +407,23 @@ recovery_target_timeline = 'latest' def move_data_directory(self): if os.path.isdir(self.data_dir) and not self.is_running(): try: - os.rename(self.data_dir, '{0}_{1}'.format(self.data_dir, time.strftime('%Y-%m-%d-%H-%M-%S'))) + new_name = '{0}_{1}'.format(self.data_dir, time.strftime('%Y-%m-%d-%H-%M-%S')) + logger.info('renaming data directory to %s', new_name) + os.rename(self.data_dir, new_name) except: - logger.exception("Could not rename data directory {0}".format(self.data_dir)) + logger.exception("Could not rename data directory %s", self.data_dir) + + def remove_data_directory(self): + logger.info('Removing data directory: %s', self.data_dir) + try: + if os.path.islink(self.data_dir): + os.unlink(self.data_dir) + elif not os.path.exists(self.data_dir): + return + elif os.path.isfile(self.data_dir): + os.remove(self.data_dir) + elif os.path.isdir(self.data_dir): + shutil.rmtree(self.data_dir) + except: + logger.exception('Could not remove data directory %s', self.data_dir) + self.move_data_directory() diff --git a/tests/test_ha.py b/tests/test_ha.py index 3fa50b41..6cb6e441 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -1,8 +1,9 @@ import unittest from mock import Mock, patch -from patroni.dcs import Cluster, DCSError +from patroni.dcs import Cluster, DCSError, Leader, Member from patroni.etcd import Client, Etcd +from patroni.exceptions import PostgresException from patroni.ha import Ha from test_etcd import socket_getaddrinfo, etcd_read, etcd_write @@ -15,18 +16,38 @@ def false(*args, **kwargs): return False -class MockPostgresql: +def get_cluster(initialize, leader): + return Cluster(initialize, leader, None, None) - def __init__(self): - self.name = 'postgresql0' - self.role = 'replica' + +def get_cluster_not_initialized_without_leader(): + return get_cluster(None, None) + + +def get_cluster_initialized_without_leader(): + return get_cluster(True, None) + + +def get_cluster_not_initialized_with_leader(): + return get_cluster(False, Leader(0, 0, 0, + Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', + None, None, 28))) + + +def get_cluster_initialized_with_leader(): + return get_cluster(True, Leader(0, 0, 0, + Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', + None, None, 28))) + + +class MockPostgresql(Mock): + + name = 'postgresql0' + role = 'replica' def is_healthy(self): return True - def write_recovery_conf(self, _): - return True - def start(self): return True @@ -36,45 +57,35 @@ class MockPostgresql: def is_leader(self): return True - def promote(self): - return True - - def demote(self, _): - return True - - def follow_the_leader(self, _): - return True - - def create_replication_slots(self, _): - return True - def last_operation(self): return 0 + def data_directory_empty(self): + return False -def get_unlocked_cluster(): - return Cluster(False, None, None, []) + def bootstrap(self, *args, **kwargs): + return True class TestHa(unittest.TestCase): @patch('socket.getaddrinfo', socket_getaddrinfo) - def setUp(self): + @patch.object(Client, 'machines') + def setUp(self, mock_machines): + mock_machines.__get__ = Mock(return_value=['http://remotehost:2379']) self.p = MockPostgresql() - with patch.object(Client, 'machines') as mock_machines: - mock_machines.__get__ = Mock(return_value=['http://remotehost:2379']) - self.e = Etcd('foo', {'ttl': 30, 'host': 'ok:2379', 'scope': 'test'}) - self.e.client.read = etcd_read - self.e.client.write = etcd_write - self.ha = Ha(self.p, self.e) - self.ha.load_cluster_from_dcs() - self.ha.cluster = get_unlocked_cluster() - self.ha.load_cluster_from_dcs = Mock() + self.e = Etcd('foo', {'ttl': 30, 'host': 'ok:2379', 'scope': 'test'}) + self.e.client.read = etcd_read + self.e.client.write = etcd_write + self.ha = Ha(self.p, self.e) + self.ha.load_cluster_from_dcs() + self.ha.cluster = get_cluster_not_initialized_without_leader() + self.ha.load_cluster_from_dcs = Mock() def test_load_cluster_from_dcs(self): ha = Ha(self.p, self.e) ha.load_cluster_from_dcs() - self.e.get_cluster = get_unlocked_cluster + self.e.get_cluster = get_cluster_not_initialized_without_leader ha.load_cluster_from_dcs() def test_start_as_slave(self): @@ -127,6 +138,12 @@ class TestHa(unittest.TestCase): self.ha.cluster.is_unlocked = false self.assertEquals(self.ha.run_cycle(), 'demoting self because i do not have the lock and i was a leader') + def test_demote_because_update_lock_failed(self): + self.ha.cluster.is_unlocked = false + self.ha.has_lock = true + self.ha.update_lock = false + self.assertEquals(self.ha.run_cycle(), 'demoting self because i do not have the lock and i was a leader') + def test_follow_the_leader(self): self.ha.cluster.is_unlocked = false self.p.is_leader = false @@ -135,3 +152,31 @@ class TestHa(unittest.TestCase): def test_no_etcd_connection_master_demote(self): self.ha.load_cluster_from_dcs = Mock(side_effect=DCSError('Etcd is not responding properly')) self.assertEquals(self.ha.run_cycle(), 'demoted self because DCS is not accessible and i was a leader') + + def test_bootstrap_from_leader(self): + self.ha.cluster = get_cluster_initialized_with_leader() + self.assertEquals(self.ha.bootstrap(), 'bootstrapped from leader') + + def test_bootstrap_from_leader_failed(self): + self.ha.cluster = get_cluster_initialized_with_leader() + self.p.bootstrap = false + self.assertEquals(self.ha.bootstrap(), 'failed to bootstrap from leader') + + def test_bootstrap_waiting_for_leader(self): + self.ha.cluster = get_cluster_initialized_without_leader() + self.assertEquals(self.ha.bootstrap(), 'waiting for leader to bootstrap') + + def test_bootstrap_initialize_lock_failed(self): + self.ha.cluster = get_cluster_not_initialized_without_leader() + self.assertEquals(self.ha.bootstrap(), 'failed to acquire initialize lock') + + def test_bootstrap_initialized_new_cluster(self): + self.ha.cluster = get_cluster_not_initialized_without_leader() + self.e.initialize = true + self.assertEquals(self.ha.bootstrap(), 'initialized a new cluster') + + def test_bootstrap_release_initialize_key_on_failure(self): + self.ha.cluster = get_cluster_not_initialized_without_leader() + self.e.initialize = true + self.p.bootstrap = Mock(side_effect=PostgresException("Could not bootstrap master PostgreSQL")) + self.assertRaises(PostgresException, self.ha.bootstrap) diff --git a/tests/test_patroni.py b/tests/test_patroni.py index 13621bf9..7cac2eba 100644 --- a/tests/test_patroni.py +++ b/tests/test_patroni.py @@ -6,14 +6,12 @@ import yaml from mock import Mock, patch from patroni.api import RestApiServer -from patroni.dcs import Cluster, Member, Leader +from patroni.dcs import Cluster, Member from patroni.etcd import Etcd -from patroni.exceptions import DCSError, PostgresException from patroni import Patroni, main from patroni.zookeeper import ZooKeeper from six.moves import BaseHTTPServer from test_etcd import Client, SleepException, etcd_read, etcd_write -from test_ha import true, false from test_postgresql import Postgresql, psycopg2_connect from test_zookeeper import MockKazooClient @@ -22,34 +20,6 @@ def time_sleep(*args): raise SleepException() -def get_cluster(initialize, leader): - return Cluster(initialize, leader, None, None) - - -def get_cluster_not_initialized_without_leader(): - return get_cluster(None, None) - - -def get_cluster_initialized_without_leader(): - return get_cluster(True, None) - - -def get_cluster_not_initialized_with_leader(): - return get_cluster(False, Leader(0, 0, 0, - Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', - None, None, 28))) - - -def get_cluster_initialized_with_leader(): - return get_cluster(True, Leader(0, 0, 0, - Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', - None, None, 28))) - - -def get_cluster_dcs_error(): - raise DCSError('') - - @patch('time.sleep', Mock()) @patch('subprocess.call', Mock(return_value=0)) @patch('psycopg2.connect', psycopg2_connect) @@ -58,7 +28,9 @@ def get_cluster_dcs_error(): @patch.object(BaseHTTPServer.HTTPServer, '__init__', Mock()) class TestPatroni(unittest.TestCase): - def setUp(self): + @patch.object(Client, 'machines') + def setUp(self, mock_machines): + mock_machines.__get__ = Mock(return_value=['http://remotehost:2379']) self.touched = False self.init_cancelled = False RestApiServer._BaseServer__is_shut_down = Mock() @@ -66,11 +38,9 @@ class TestPatroni(unittest.TestCase): RestApiServer.socket = 0 with open('postgres0.yml', 'r') as f: config = yaml.load(f) - with patch.object(Client, 'machines') as mock_machines: - mock_machines.__get__ = Mock(return_value=['http://remotehost:2379']) - self.p = Patroni(config) - self.p.ha.dcs.client.write = etcd_write - self.p.ha.dcs.client.read = etcd_read + self.p = Patroni(config) + self.p.ha.dcs.client.write = etcd_write + self.p.ha.dcs.client.read = etcd_read @patch('patroni.zookeeper.KazooClient', MockKazooClient()) def test_get_dcs(self): @@ -80,26 +50,26 @@ class TestPatroni(unittest.TestCase): @patch('time.sleep', Mock(side_effect=SleepException())) @patch.object(Patroni, 'initialize', Mock()) @patch.object(Etcd, 'delete_leader', Mock()) - def test_patroni_main(self): + @patch.object(Client, 'machines') + def test_patroni_main(self, mock_machines): main() sys.argv = ['patroni.py', 'postgres0.yml'] - with patch.object(Client, 'machines') as mock_machines: - mock_machines.__get__ = Mock(return_value=['http://remotehost:2379']) - with patch.object(Patroni, 'touch_member', self.touch_member): - with patch.object(Patroni, 'run', Mock(side_effect=SleepException())): - self.assertRaises(SleepException, main) - with patch.object(Patroni, 'run', Mock(side_effect=KeyboardInterrupt())): - main() + mock_machines.__get__ = Mock(return_value=['http://remotehost:2379']) + with patch.object(Patroni, 'touch_member', self.touch_member): + with patch.object(Patroni, 'run', Mock(side_effect=SleepException())): + self.assertRaises(SleepException, main) + with patch.object(Patroni, 'run', Mock(side_effect=KeyboardInterrupt())): + main() @patch('time.sleep', Mock(side_effect=SleepException())) - def test_patroni_run(self): + def test_run(self): self.p.touch_member = self.touch_member self.p.ha.state_handler.sync_replication_slots = time_sleep self.p.ha.dcs.watch = time_sleep self.assertRaises(SleepException, self.p.run) - self.p.ha.state_handler.is_leader = false + self.p.ha.state_handler.is_leader = Mock(return_value=False) self.p.api.start = Mock() self.assertRaises(SleepException, self.p.run) @@ -119,48 +89,10 @@ class TestPatroni(unittest.TestCase): def test_patroni_initialize(self): self.p.touch_member = self.touch_member - self.p.postgresql.data_directory_empty = true - self.p.ha.dcs.initialize = true - self.p.postgresql.initialize = true - self.p.postgresql.start = true - self.p.ha.dcs.get_cluster = get_cluster_not_initialized_without_leader self.p.initialize() - self.p.ha.dcs.initialize = false - self.p.ha.dcs.get_cluster = get_cluster_initialized_with_leader - with patch('time.sleep', time_sleep): - self.p.initialize() - - self.p.ha.dcs.get_cluster = get_cluster_initialized_without_leader - self.assertRaises(SleepException, self.p.initialize) - - self.p.postgresql.data_directory_empty = false - self.p.initialize() - - self.p.ha.dcs.get_cluster = get_cluster_not_initialized_with_leader - self.p.postgresql.data_directory_empty = true - self.p.initialize() - - self.p.ha.dcs.get_cluster = get_cluster_dcs_error - self.assertRaises(SleepException, self.p.initialize) - def test_schedule_next_run(self): self.p.ha.dcs.watch = Mock(return_value=True) self.p.schedule_next_run() self.p.next_run = time.time() - self.p.nap_time - 1 self.p.schedule_next_run() - - def cancel_initialization(self): - self.init_cancelled = True - - def test_cleanup_on_initialization(self): - self.p.ha.dcs.get_cluster = get_cluster_not_initialized_without_leader - self.p.touch_member = self.touch_member - self.p.postgresql.data_directory_empty = true - self.p.ha.dcs.initialize = true - self.p.postgresql.initialize = true - self.p.postgresql.start = false - - self.p.ha.dcs.cancel_initialization = self.cancel_initialization - self.assertRaises(PostgresException, self.p.initialize) - self.assertTrue(self.init_cancelled) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index eaf84943..b4988c25 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -5,7 +5,7 @@ import unittest from mock import Mock, patch from patroni.dcs import Cluster, Leader, Member -from patroni.exceptions import PostgresConnectionException +from patroni.exceptions import PostgresException, PostgresConnectionException from patroni.postgresql import Postgresql from patroni.utils import RetryFailedError from test_ha import false @@ -209,3 +209,20 @@ class TestPostgresql(unittest.TestCase): self.p.move_data_directory() with patch('os.rename', Mock(side_effect=OSError())): self.p.move_data_directory() + + def test_bootstrap(self): + self.assertRaises(PostgresException, self.p.bootstrap) + self.p.start = Mock(return_value=True) + self.p.bootstrap() + + def test_remove_data_directory(self): + self.p.data_dir = 'data_dir' + self.p.remove_data_directory() + os.mkdir(self.p.data_dir) + self.p.remove_data_directory() + open(self.p.data_dir, 'w').close() + self.p.remove_data_directory() + os.symlink('unexisting', self.p.data_dir) + with patch('os.unlink', Mock(side_effect=Exception)): + self.p.remove_data_directory() + self.p.remove_data_directory() From a4266be3da8321dad2a924f89ae7a5cfc8e8006d Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 23 Sep 2015 10:59:55 +0200 Subject: [PATCH 2/2] remove unused function --- tests/test_ha.py | 6 ------ 1 file changed, 6 deletions(-) diff --git a/tests/test_ha.py b/tests/test_ha.py index 6cb6e441..925beb33 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -28,12 +28,6 @@ def get_cluster_initialized_without_leader(): return get_cluster(True, None) -def get_cluster_not_initialized_with_leader(): - return get_cluster(False, Leader(0, 0, 0, - Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', - None, None, 28))) - - def get_cluster_initialized_with_leader(): return get_cluster(True, Leader(0, 0, 0, Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres',