Run initial cluster bootstrap from the main loop

This commit is contained in:
Alexander Kukushkin
2015-09-23 10:55:38 +02:00
parent 83c5416c82
commit e83651b57b
6 changed files with 228 additions and 206 deletions
+1 -34
View File
@@ -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
+93 -49
View File
@@ -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():
+21 -4
View File
@@ -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()
+78 -33
View File
@@ -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:[email protected]:5435/postgres',
None, None, 28)))
def get_cluster_initialized_with_leader():
return get_cluster(True, Leader(0, 0, 0,
Member(0, 'leader', 'postgres://replicator:[email protected]: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)
+17 -85
View File
@@ -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:[email protected]:5435/postgres',
None, None, 28)))
def get_cluster_initialized_with_leader():
return get_cluster(True, Leader(0, 0, 0,
Member(0, 'leader', 'postgres://replicator:[email protected]: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)
+18 -1
View File
@@ -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()