diff --git a/patroni/__init__.py b/patroni/__init__.py index 2da27729..cc008fda 100644 --- a/patroni/__init__.py +++ b/patroni/__init__.py @@ -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 diff --git a/patroni/dcs.py b/patroni/dcs.py index a33dea29..9017da7f 100644 --- a/patroni/dcs.py +++ b/patroni/dcs.py @@ -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) diff --git a/patroni/etcd.py b/patroni/etcd.py index 7a527f6e..e24fa66b 100644 --- a/patroni/etcd.py +++ b/patroni/etcd.py @@ -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) diff --git a/patroni/exceptions.py b/patroni/exceptions.py index e6159d47..507edfc3 100644 --- a/patroni/exceptions.py +++ b/patroni/exceptions.py @@ -13,5 +13,9 @@ class PatroniException(Exception): return repr(self.value) +class PostgresException(PatroniException): + pass + + class DCSError(PatroniException): pass diff --git a/patroni/postgresql.py b/patroni/postgresql.py index d729942c..b7f76fc3 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -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)) diff --git a/patroni/zookeeper.py b/patroni/zookeeper.py index 3901e89d..708cbafa 100644 --- a/patroni/zookeeper.py +++ b/patroni/zookeeper.py @@ -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) diff --git a/tests/test_etcd.py b/tests/test_etcd.py index 94b7f807..66261eeb 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -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 diff --git a/tests/test_patroni.py b/tests/test_patroni.py index a45ebf62..317e1e5c 100644 --- a/tests/test_patroni.py +++ b/tests/test_patroni.py @@ -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: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 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) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 273d077d..56dbc557 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -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() diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py index d123932e..4b31bad2 100644 --- a/tests/test_zookeeper.py +++ b/tests/test_zookeeper.py @@ -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')