Merge pull request #30 from zalando/feature/cleanup_on_failed_initialization

Make sure initialize flag is reset on failure.
This commit is contained in:
Oleksii Kliukin
2015-09-14 12:57:39 +02:00
10 changed files with 209 additions and 59 deletions
+30 -15
View File
@@ -6,6 +6,7 @@ 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
@@ -42,6 +43,14 @@ 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():
@@ -50,21 +59,27 @@ class Patroni:
# is data directory empty?
if self.postgresql.data_directory_empty():
# racing to initialize
if self.ha.dcs.race('/initialize'):
self.postgresql.initialize()
self.ha.dcs.take_leader()
self.postgresql.start()
self.postgresql.create_replication_user()
self.postgresql.create_connection_users()
else:
while True:
leader = self.ha.dcs.current_leader()
if leader and self.postgresql.sync_from_leader(leader):
self.postgresql.write_recovery_conf(leader)
self.postgresql.start()
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
sleep(5)
# 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.load_replication_slots()
@@ -110,8 +125,8 @@ def main():
config = yaml.load(f)
patroni = Patroni(config)
patroni.initialize()
try:
patroni.initialize()
patroni.run()
except KeyboardInterrupt:
pass
+6 -3
View File
@@ -165,12 +165,11 @@ class AbstractDCS:
overwriting the key if necessary."""
@abc.abstractmethod
def race(self, path):
def initialize(self):
"""Race for cluster initialization.
:param path: usually this is just '/initialize'
:returns: `!True` if key has been created successfully.
this method should create atomically `path` key and return `!True`
this method should create atomically initialize key and return `!True`
otherwise it should return `!False`"""
@abc.abstractmethod
@@ -178,5 +177,9 @@ class AbstractDCS:
"""Voluntarily remove leader key from DCS
This method should remove leader key if current instance is the leader"""
@abc.abstractmethod
def cancel_initialization(self):
""" Removes the initialize key for a cluster """
def watch(self, timeout):
sleep(timeout)
+6 -3
View File
@@ -232,19 +232,22 @@ class Etcd(AbstractDCS):
return ret
@catch_etcd_errors
def race(self, path):
return self.retry(self.client.write, self.client_path(path), self._name, prevExist=False)
def initialize(self):
return self.client.write(self.initialize_path, self._name, prevExist=False)
@catch_etcd_errors
def delete_leader(self):
return self.client.delete(self.leader_path, prevValue=self._name)
@catch_etcd_errors
def cancel_initialization(self):
return self.client.delete(self.initialize_path, prevValue=self._name)
def watch(self, timeout):
# watch on leader key changes if it is defined and current node is not lock owner
if self.cluster and self.cluster.leader and self.cluster.leader.name != self._name:
end_time = time.time() + timeout
index = self.cluster.leader.index
while index and timeout >= 1: # when timeout is too small urllib3 doesn't have enough time to connect
try:
res = self.client.watch(self.leader_path, index=index + 1, timeout=timeout)
+4
View File
@@ -13,5 +13,9 @@ class PatroniException(Exception):
return repr(self.value)
class PostgresException(PatroniException):
pass
class DCSError(PatroniException):
pass
+34 -1
View File
@@ -4,7 +4,9 @@ import psycopg2
import shlex
import shutil
import subprocess
import time
from patroni.exceptions import PostgresException
from patroni.utils import sleep
from six.moves.urllib_parse import urlparse
@@ -159,7 +161,7 @@ class Postgresql:
return ret
def is_running(self):
return subprocess.call(' '.join(self._pg_ctl) + ' status > /dev/null', shell=True) == 0
return subprocess.call(' '.join(self._pg_ctl) + ' status > /dev/null 2>&1', shell=True) == 0
def call_nowait(self, cb_name, is_leader=None):
""" pick a callback command and call it without waiting for it to finish """
@@ -391,3 +393,34 @@ recovery_target_timeline = 'latest'
def last_operation(self):
return str(self.xlog_position())
def bootstrap(self, current_leader=None):
"""
Initially bootstrap PostgreSQL, either by creating a data
directory with initdb, or by initalizing a replica from an
exiting leader. Failure in the first case always leads to
exception, since there is no point in continuing if initdb failed.
In the second case, however, a False is returned on failure, since
it is normal for the replica to retry a failed attempt to initialize
from the master.
"""
ret = False
if not current_leader:
ret = self.initialize() and self.start()
if ret:
self.create_replication_user()
self.create_connection_users()
else:
raise PostgresException("Could not bootstrap master PostgreSQL")
else:
if self.sync_from_leader(current_leader):
self.write_recovery_conf(current_leader)
ret = self.start()
return ret
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')))
except:
logger.exception("Could not rename data directory {0}".format(self.data_dir))
+43 -24
View File
@@ -92,9 +92,8 @@ class ZooKeeper(AbstractDCS):
self.client.add_listener(self.session_listener)
self.cluster_event = self.client.handler.event_object()
self.cluster = None
self.fetch_cluster = True
self.members = []
self.leader = None
self.last_leader_operation = 0
self.client.start(None)
@@ -111,28 +110,39 @@ class ZooKeeper(AbstractDCS):
try:
return self.client.get(key, watch)
except NoNodeError:
pass
except:
logger.exception('get_node')
return None
return None
@staticmethod
def member(name, value, znode):
conn_url, api_url = parse_connection_string(value)
return Member(znode.mzxid, name, conn_url, api_url, None, None)
def get_children(self, key, watch=None):
try:
return self.client.get_children(key, watch)
except NoNodeError:
return []
def load_members(self):
members = []
for member in self.client.get_children(self.members_path, self.cluster_watcher):
data = self.get_node(self.member_path)
for member in self.get_children(self.members_path, self.cluster_watcher):
data = self.get_node(self.members_path + member)
if data is not None:
members.append(self.member(member, *data))
return members
def _inner_load_cluster(self):
self.cluster_event.clear()
leader = self.get_node(self.leader_path, self.cluster_watcher)
self.members = self.load_members()
nodes = set(self.get_children(self.client_path('')))
# get initialize flag
initialize = self._INITIALIZE in nodes
# get list of members
members = self.load_members() if self._MEMBERS[:-1] in nodes else []
# get leader
leader = self.get_node(self.leader_path, self.cluster_watcher) if self._LEADER in nodes else None
if leader:
client_id = self.client.client_id
if leader[0] == self._name and client_id is not None and client_id[0] != leader[1].ephemeralOwner:
@@ -142,15 +152,14 @@ class ZooKeeper(AbstractDCS):
if leader:
member = Member(-1, leader[0], None, None, None, None)
member = ([m for m in self.members if m.name == leader[0]] or [member])[0]
member = ([m for m in members if m.name == leader[0]] or [member])[0]
leader = Leader(leader[1].mzxid, None, None, member)
self.fetch_cluster = member.index == -1
self.leader = leader
if self.fetch_cluster:
last_leader_operation = self.get_node(self.leader_optime_path)
if last_leader_operation:
self.last_leader_operation = int(last_leader_operation[0])
# get last leader operation
self.last_leader_operation = self.get_node(self.leader_optime_path) if self.fetch_cluster else None
self.last_leader_operation = 0 if self.last_leader_operation is None else int(self.last_leader_operation[0])
self.cluster = Cluster(initialize, leader, self.last_leader_operation, members)
def get_cluster(self):
if self.exhibitor and self.exhibitor.poll():
@@ -163,7 +172,7 @@ class ZooKeeper(AbstractDCS):
logger.exception('get_cluster')
self.session_listener(KazooState.LOST)
raise ZooKeeperError('ZooKeeper in not responding properly')
return Cluster(True, self.leader, self.last_leader_operation, self.members)
return self.cluster
def _create(self, path, value, **kwargs):
try:
@@ -177,13 +186,12 @@ class ZooKeeper(AbstractDCS):
ret or logger.info('Could not take out TTL lock')
return ret
def race(self, path):
return self._create(self.client_path(path), self._name, makepath=True)
def initialize(self):
return self._create(self.initialize_path, self._name, makepath=True)
def touch_member(self, connection_string, ttl=None):
for m in self.members:
if m.name == self._name:
return True
if self.cluster and any(m.name == self._name for m in self.cluster.members):
return True
path = self.member_path
try:
self.client.retry(self.client.create, path, connection_string, makepath=True, ephemeral=True)
@@ -217,8 +225,19 @@ class ZooKeeper(AbstractDCS):
return True
def delete_leader(self):
if isinstance(self.leader, Leader) and self.leader.name == self._name:
self.client.delete(self.leader_path)
if isinstance(self.cluster, Cluster) and self.cluster.leader.name == self._name:
self.client.delete(self.leader_path, version=self.cluster.leader.index)
def _cancel_initialization(self):
node = self.get_node(self.initialize_path)
if node and node[0] == self._name:
self.client.delete(self.initialize_path, version=node[1].mzxid)
def cancel_initialization(self):
try:
self.client.retry(self._cancel_initialization)
except:
logger.exception("Unable to delete initialize key")
def watch(self, timeout):
self.cluster_event.wait(timeout)
+6 -2
View File
@@ -266,8 +266,12 @@ class TestEtcd(unittest.TestCase):
def test_update_leader(self):
self.assertTrue(self.etcd.update_leader(MockPostgresql()))
def test_race(self):
self.assertFalse(self.etcd.race(''))
def test_initialize(self):
self.assertFalse(self.etcd.initialize())
def test_cancel_initializion(self):
self.etcd.client.delete = etcd_delete
self.assertFalse(self.etcd.cancel_initialization())
def test_delete_leader(self):
self.etcd.client.delete = etcd_delete
+57 -6
View File
@@ -9,8 +9,9 @@ import yaml
from mock import Mock, patch
from patroni.api import RestApiServer
from patroni.dcs import Cluster, Member
from patroni.dcs import Cluster, Member, Leader
from patroni.etcd import Etcd
from patroni.exceptions import PostgresException
from patroni import Patroni, main
from patroni.zookeeper import ZooKeeper
from six.moves import BaseHTTPServer
@@ -41,6 +42,30 @@ class Mock_BaseServer__is_shut_down:
pass
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)))
class TestPatroni(unittest.TestCase):
def __init__(self, method_name='runTest'):
@@ -50,6 +75,7 @@ class TestPatroni(unittest.TestCase):
def set_up(self):
self.touched = False
self.init_cancelled = False
subprocess.call = subprocess_call
psycopg2.connect = psycopg2_connect
self.time_sleep = time.sleep
@@ -126,24 +152,49 @@ class TestPatroni(unittest.TestCase):
self.p.touch_member()
def test_patroni_initialize(self):
self.p.postgresql.should_use_s3_to_create_replica = false
self.p.ha.dcs.client.write = etcd_write
self.p.ha.dcs.client.read = etcd_read
self.p.touch_member = self.touch_member
self.p.postgresql.data_directory_empty = true
self.p.ha.dcs.race = 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.race = false
self.p.ha.dcs.initialize = false
self.p.ha.dcs.get_cluster = get_cluster_initialized_with_leader
time.sleep = time_sleep
self.p.ha.dcs.client.read = etcd_read
self.p.initialize()
self.p.ha.dcs.current_leader = nop
self.assertRaises(Exception, 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()
def test_schedule_next_run(self):
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.client.write = etcd_write
self.p.ha.dcs.client.read = etcd_read
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)
+7
View File
@@ -6,6 +6,7 @@ import unittest
from patroni.dcs import Cluster, Leader, Member
from patroni.postgresql import Postgresql
from test_ha import true, false
def nop(*args, **kwargs):
@@ -217,3 +218,9 @@ class TestPostgresql(unittest.TestCase):
self.p.start()
self.p.query = self.mock_query
self.assertTrue(self.p.stop())
def test_move_data_directory(self):
self.p.is_running = is_running
os.rename = nop
os.path.isdir = true
self.p.move_data_directory()
+16 -5
View File
@@ -58,8 +58,6 @@ class MockKazooClient:
def get(self, path, watch=None):
if path == '/no_node':
raise NoNodeError
elif path == '/other_exception':
raise Exception()
elif '/members/' in path:
return (
'postgres://repuser:rep-pass@localhost:5434/postgres?application_name=http://127.0.0.1:8009/patroni',
@@ -71,8 +69,14 @@ class MockKazooClient:
if self.leader:
return ('foo', ZnodeStat(0, 0, 0, 0, 0, 0, 0, -1, 0, 0, 0))
return ('foo', ZnodeStat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0))
elif path.endswith('/initialize'):
return ('foo', ZnodeStat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0))
def get_children(self, path, watch=None, include_data=False):
if path == '/no_node':
raise NoNodeError
elif path in ['/service/bla/', '/service/test/']:
return ['initialize', 'leader', 'members', 'optime']
return ['foo', 'bar', 'buzz']
def create(self, path, value="", acl=None, ephemeral=False, sequence=False, makepath=False):
@@ -93,6 +97,8 @@ class MockKazooClient:
return
self.leader = True
raise Exception
elif path.endswith('/initialize'):
raise NoNodeError
def set_hosts(self, hosts, randomize_hosts=None):
pass
@@ -132,7 +138,9 @@ class TestZooKeeper(unittest.TestCase):
def test_get_node(self):
self.assertIsNone(self.zk.get_node('/no_node'))
self.assertIsNone(self.zk.get_node('/other_exception'))
def test_get_children(self):
self.assertListEqual(self.zk.get_children('/no_node'), [])
def test__inner_load_cluster(self):
self.zk._base_path = self.zk._base_path.replace('test', 'bla')
@@ -146,8 +154,11 @@ class TestZooKeeper(unittest.TestCase):
self.zk.touch_member('foo')
self.zk.delete_leader()
def test_race(self):
self.assertFalse(self.zk.race('/initialize'))
def test_initialize(self):
self.assertFalse(self.zk.initialize())
def test_cancel_initialization(self):
self.zk.cancel_initialization()
def test_touch_member(self):
self.zk.touch_member('new')