Compare commits

..
22 Commits
Author SHA1 Message Date
Oleksii Kliukin bf5737614d Merge pull request #30 from zalando/feature/cleanup_on_failed_initialization
Make sure initialize flag is reset on failure.
2015-09-14 12:57:39 +02:00
Oleksii Kliukin d69403ab6f Merge pull request #37 from zalando/feature/zookeeper-fetch-initialize
Build Cluster object for ZooKeeper the same way as for Etcd
2015-09-14 12:55:10 +02:00
Oleksii Kliukin 51eacc5042 Handle the case when initialize flag is not set and leader is present. 2015-09-14 12:36:28 +02:00
Alexander Kukushkin 98488a00a2 Remove unused import of KazooException 2015-09-14 12:00:24 +02:00
Alexander Kukushkin 4a081bcb71 Run cancel_initialization with retry 2015-09-14 11:58:10 +02:00
Alexander Kukushkin 209c985420 get_node and get_children should catch only NoNodeError exception.
All other exceptions are needed to have retry functionality working
correctly.
2015-09-14 11:45:00 +02:00
Alexander Kukushkin f494d2ce64 Build Cluster object for ZooKeeper the same way as for Etcd
Previous implementation was always setting Cluster.initialize to True.
Also it was throwing ZooKeeperError when there were no members in a
cluster.

Plus BUGFIX of a bug introduced with
https://github.com/zalando/patroni/pull/34 in a `load_members` method.
- data = self.get_node(self.member_path)
+ data = self.get_node(self.members_path + member)
It was always fetching the same node for all cluster members.
Fortunately Etcd doesn't have such problem because we are fetching the
whole cluster directory with one recursive API call.
2015-09-14 11:19:46 +02:00
Oleksii Kliukin be110c4ba0 Do not try to stop postgres twice if initialization had failed. 2015-09-14 09:20:45 +02:00
Oleksii Kliukin 15cd10669d Change an outdated comment. 2015-09-10 18:04:47 +02:00
Oleksii Kliukin cd312de252 Fix a flake8 warning.
Improve some unit tests by expecting specific exceptions.
2015-09-10 17:15:43 +02:00
Oleksii Kliukin 2377c417e4 Fix etcd and zookeper interactions with initialize key.
Fix unittests as well.
2015-09-10 16:05:10 +02:00
Oleksii Kliukin 938b946e55 Merge branch 'master' into feature/cleanup_on_failed_initialization 2015-09-10 15:43:31 +02:00
Oleksii Kliukin 30a9e0f7f5 Move PostgreSQL data directory if init had failed.
Prevent treating the incompletely-initialized PostgreSQL cluster
as a valid on restart by forcefully moving the data directory.
I don't want to remove it altogether, since a DBA might decide
to analyze the failed PG cluster in order to resolve the init
issue.
2015-09-10 15:34:29 +02:00
Alexander Kukushkin 30a7d50a56 Merge pull request #34 from zalando/feature/constants-for-key-names
Define initialize, leader, optime and members string constansts in Ab…
2015-09-10 15:33:08 +02:00
Alexander Kukushkin 36cbd34ffc Fix zookeeper test coverage 2015-09-09 15:59:02 +02:00
Alexander Kukushkin 5bdb18761b Define initialize, leader, optime and members string constansts in AbstractDCS
Also define following properties:
* initialize_path
* members_path
* member_path
* leader_path
* leader_optime_path

And replace any occurrences of these strings or client_path calls in a
etcd and zookeeper implementations with given constants and properties.
2015-09-09 15:10:45 +02:00
Alexander Kukushkin e90b14cd3b Merge pull request #31 from zalando/bugfix/script_paths
Fix path to scripts subdirectory in configuration files.
2015-09-08 16:25:52 +02:00
Oleksii Kliukin 1c61280d70 Fix path to scripts subdirectory in configuration files. 2015-09-08 16:10:42 +02:00
Oleksii Kliukin ff499604f0 Act on removal of initialization flag.
If initializer node suddenly dies before the initialization is complete,
other nodes should try to take over.

Fix some unittests for etcd and zookeeper and add couple of new ones.
2015-09-08 16:04:54 +02:00
Oleksii Kliukin 92647b7aad Merge branch 'master' of https://github.com/zalando/patroni into feature/cleanup_on_failed_initialization 2015-09-08 14:54:52 +02:00
Feike Steenbergen dd8472f639 Tag on github is prefixed with v. 2015-09-08 13:17:36 +02:00
Oleksii Kliukin b842ed478b Make sure initialize flag is reset on failure.
Cleanup the initialize flag if the initializing node fails
to bootstrap its PostgreSQL database.

Rename dcs.race to initialize, since we only call it for the
initialize flag. Factored out PostgreSQL bootstrapping code
into a separate function.
2015-09-08 12:03:34 +02:00
13 changed files with 260 additions and 86 deletions
+30 -15
View File
@@ -6,6 +6,7 @@ import yaml
from patroni.api import RestApiServer from patroni.api import RestApiServer
from patroni.etcd import Etcd from patroni.etcd import Etcd
from patroni.exceptions import DCSError
from patroni.ha import Ha from patroni.ha import Ha
from patroni.postgresql import Postgresql from patroni.postgresql import Postgresql
from patroni.utils import setup_signal_handlers, sleep, reap_children from patroni.utils import setup_signal_handlers, sleep, reap_children
@@ -42,6 +43,14 @@ class Patroni:
return True return True
return self.ha.dcs.touch_member(connection_string, ttl) 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): def initialize(self):
# wait for etcd to be available # wait for etcd to be available
while not self.touch_member(): while not self.touch_member():
@@ -50,21 +59,27 @@ class Patroni:
# is data directory empty? # is data directory empty?
if self.postgresql.data_directory_empty(): if self.postgresql.data_directory_empty():
# racing to initialize while True:
if self.ha.dcs.race('/initialize'): try:
self.postgresql.initialize() cluster = self.ha.dcs.get_cluster()
self.ha.dcs.take_leader() if not cluster.is_unlocked(): # the leader already exists
self.postgresql.start() if not cluster.initialize:
self.postgresql.create_replication_user() self.ha.dcs.initialize()
self.postgresql.create_connection_users() self.postgresql.bootstrap(cluster.leader)
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()
break 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(): elif self.postgresql.is_running():
self.postgresql.load_replication_slots() self.postgresql.load_replication_slots()
@@ -110,8 +125,8 @@ def main():
config = yaml.load(f) config = yaml.load(f)
patroni = Patroni(config) patroni = Patroni(config)
patroni.initialize()
try: try:
patroni.initialize()
patroni.run() patroni.run()
except KeyboardInterrupt: except KeyboardInterrupt:
pass pass
+33 -4
View File
@@ -74,6 +74,12 @@ class AbstractDCS:
__metaclass__ = abc.ABCMeta __metaclass__ = abc.ABCMeta
_INITIALIZE = 'initialize'
_LEADER = 'leader'
_MEMBERS = 'members/'
_OPTIME = 'optime'
_LEADER_OPTIME = _OPTIME + '/' + _LEADER
def __init__(self, name, config): def __init__(self, name, config):
""" """
:param name: name of current instance (the same value as `~Postgresql.name`) :param name: name of current instance (the same value as `~Postgresql.name`)
@@ -85,7 +91,27 @@ class AbstractDCS:
self._base_path = '/service/' + self._scope self._base_path = '/service/' + self._scope
def client_path(self, path): def client_path(self, path):
return self._base_path + path return '/'.join([self._base_path, path.lstrip('/')])
@property
def initialize_path(self):
return self.client_path(self._INITIALIZE)
@property
def members_path(self):
return self.client_path(self._MEMBERS)
@property
def member_path(self):
return self.client_path(self._MEMBERS + self._name)
@property
def leader_path(self):
return self.client_path(self._LEADER)
@property
def leader_optime_path(self):
return self.client_path(self._LEADER_OPTIME)
@abc.abstractmethod @abc.abstractmethod
def get_cluster(self): def get_cluster(self):
@@ -139,12 +165,11 @@ class AbstractDCS:
overwriting the key if necessary.""" overwriting the key if necessary."""
@abc.abstractmethod @abc.abstractmethod
def race(self, path): def initialize(self):
"""Race for cluster initialization. """Race for cluster initialization.
:param path: usually this is just '/initialize'
:returns: `!True` if key has been created successfully. :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`""" otherwise it should return `!False`"""
@abc.abstractmethod @abc.abstractmethod
@@ -152,5 +177,9 @@ class AbstractDCS:
"""Voluntarily remove leader key from DCS """Voluntarily remove leader key from DCS
This method should remove leader key if current instance is the leader""" 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): def watch(self, timeout):
sleep(timeout) sleep(timeout)
+17 -16
View File
@@ -179,17 +179,17 @@ class Etcd(AbstractDCS):
nodes = {os.path.relpath(node.key, result.key): node for node in result.leaves} nodes = {os.path.relpath(node.key, result.key): node for node in result.leaves}
# get initialize flag # get initialize flag
initialize = bool(nodes.get('initialize', False)) initialize = bool(nodes.get(self._INITIALIZE, False))
# get last leader operation # get last leader operation
last_leader_operation = nodes.get('optime/leader', None) last_leader_operation = nodes.get(self._LEADER_OPTIME, None)
last_leader_operation = 0 if last_leader_operation is None else int(last_leader_operation.value) last_leader_operation = 0 if last_leader_operation is None else int(last_leader_operation.value)
# get list of members # get list of members
members = [self.member(n) for k, n in nodes.items() if k.startswith('members/') and len(k.split('/')) == 2] members = [self.member(n) for k, n in nodes.items() if k.startswith(self._MEMBERS) and k.count('/') == 1]
# get leader # get leader
leader = nodes.get('leader', None) leader = nodes.get(self._LEADER, None)
if leader: if leader:
member = Member(-1, leader.value, None, None, None, None) member = Member(-1, leader.value, None, None, None, None)
member = ([m for m in members if m.name == leader.value] or [member])[0] member = ([m for m in members if m.name == leader.value] or [member])[0]
@@ -206,17 +206,15 @@ class Etcd(AbstractDCS):
@catch_etcd_errors @catch_etcd_errors
def touch_member(self, connection_string, ttl=None): def touch_member(self, connection_string, ttl=None):
return self.retry(self.client.set, self.client_path('/members/' + self._name), return self.retry(self.client.set, self.member_path, connection_string, ttl or self.member_ttl)
connection_string, ttl or self.member_ttl)
@catch_etcd_errors @catch_etcd_errors
def take_leader(self): def take_leader(self):
return self.retry(self.client.set, self.client_path('/leader'), self._name, self.ttl) return self.retry(self.client.set, self.leader_path, self._name, self.ttl)
def attempt_to_acquire_leader(self): def attempt_to_acquire_leader(self):
try: try:
return not self.retry(self.client.write, self.client_path('/leader'), return bool(self.retry(self.client.write, self.leader_path, self._name, ttl=self.ttl, prevExist=False))
self._name, ttl=self.ttl, prevExist=False) is None
except etcd.EtcdAlreadyExist: except etcd.EtcdAlreadyExist:
logger.info('Could not take out TTL lock') logger.info('Could not take out TTL lock')
except (RetryFailedError, etcd.EtcdException): except (RetryFailedError, etcd.EtcdException):
@@ -225,31 +223,34 @@ class Etcd(AbstractDCS):
@catch_etcd_errors @catch_etcd_errors
def write_leader_optime(self, state_handler): def write_leader_optime(self, state_handler):
return self.client.set(self.client_path('/optime/leader'), state_handler.last_operation()) return self.client.set(self.leader_optime_path, state_handler.last_operation())
@catch_etcd_errors @catch_etcd_errors
def update_leader(self, state_handler): def update_leader(self, state_handler):
ret = self.retry(self.client.test_and_set, self.client_path('/leader'), self._name, self._name, self.ttl) ret = self.retry(self.client.test_and_set, self.leader_path, self._name, self._name, self.ttl)
ret and self.write_leader_optime(state_handler) ret and self.write_leader_optime(state_handler)
return ret return ret
@catch_etcd_errors @catch_etcd_errors
def race(self, path): def initialize(self):
return self.retry(self.client.write, self.client_path(path), self._name, prevExist=False) return self.client.write(self.initialize_path, self._name, prevExist=False)
@catch_etcd_errors @catch_etcd_errors
def delete_leader(self): def delete_leader(self):
return self.client.delete(self.client_path('/leader'), prevValue=self._name) 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): def watch(self, timeout):
# watch on leader key changes if it is defined and current node is not lock owner # 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: if self.cluster and self.cluster.leader and self.cluster.leader.name != self._name:
end_time = time.time() + timeout end_time = time.time() + timeout
index = self.cluster.leader.index index = self.cluster.leader.index
while index and timeout >= 1: # when timeout is too small urllib3 doesn't have enough time to connect while index and timeout >= 1: # when timeout is too small urllib3 doesn't have enough time to connect
try: try:
res = self.client.watch(self.client_path('/leader'), index=index + 1, timeout=timeout) res = self.client.watch(self.leader_path, index=index + 1, timeout=timeout)
if res.action not in ['set', 'compareAndSwap'] or res.value != self.cluster.leader.name: if res.action not in ['set', 'compareAndSwap'] or res.value != self.cluster.leader.name:
return return
index = res.modifiedIndex index = res.modifiedIndex
+4
View File
@@ -13,5 +13,9 @@ class PatroniException(Exception):
return repr(self.value) return repr(self.value)
class PostgresException(PatroniException):
pass
class DCSError(PatroniException): class DCSError(PatroniException):
pass pass
+34 -1
View File
@@ -4,7 +4,9 @@ import psycopg2
import shlex import shlex
import shutil import shutil
import subprocess import subprocess
import time
from patroni.exceptions import PostgresException
from patroni.utils import sleep from patroni.utils import sleep
from six.moves.urllib_parse import urlparse from six.moves.urllib_parse import urlparse
@@ -159,7 +161,7 @@ class Postgresql:
return ret return ret
def is_running(self): 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): def call_nowait(self, cb_name, is_leader=None):
""" pick a callback command and call it without waiting for it to finish """ """ 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): def last_operation(self):
return str(self.xlog_position()) 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))
+50 -31
View File
@@ -92,9 +92,8 @@ class ZooKeeper(AbstractDCS):
self.client.add_listener(self.session_listener) self.client.add_listener(self.session_listener)
self.cluster_event = self.client.handler.event_object() self.cluster_event = self.client.handler.event_object()
self.cluster = None
self.fetch_cluster = True self.fetch_cluster = True
self.members = []
self.leader = None
self.last_leader_operation = 0 self.last_leader_operation = 0
self.client.start(None) self.client.start(None)
@@ -107,50 +106,60 @@ class ZooKeeper(AbstractDCS):
self.fetch_cluster = True self.fetch_cluster = True
self.cluster_event.set() self.cluster_event.set()
def get_node(self, name, watch=None): def get_node(self, key, watch=None):
try: try:
return self.client.get(self.client_path(name), watch) return self.client.get(key, watch)
except NoNodeError: except NoNodeError:
pass return None
except:
logger.exception('get_node')
return None
@staticmethod @staticmethod
def member(name, value, znode): def member(name, value, znode):
conn_url, api_url = parse_connection_string(value) conn_url, api_url = parse_connection_string(value)
return Member(znode.mzxid, name, conn_url, api_url, None, None) 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): def load_members(self):
members = [] members = []
for member in self.client.get_children(self.client_path('/members'), self.cluster_watcher): for member in self.get_children(self.members_path, self.cluster_watcher):
data = self.get_node('/members/' + member) data = self.get_node(self.members_path + member)
if data is not None: if data is not None:
members.append(self.member(member, *data)) members.append(self.member(member, *data))
return members return members
def _inner_load_cluster(self): def _inner_load_cluster(self):
self.cluster_event.clear() self.cluster_event.clear()
leader = self.get_node('/leader', self.cluster_watcher) nodes = set(self.get_children(self.client_path('')))
self.members = self.load_members()
# 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: if leader:
client_id = self.client.client_id client_id = self.client.client_id
if leader[0] == self._name and client_id is not None and client_id[0] != leader[1].ephemeralOwner: if leader[0] == self._name and client_id is not None and client_id[0] != leader[1].ephemeralOwner:
logger.info('I am leader but not owner of the session. Removing leader node') logger.info('I am leader but not owner of the session. Removing leader node')
self.client.delete(self.client_path('/leader')) self.client.delete(self.leader_path)
leader = None leader = None
if leader: if leader:
member = Member(-1, leader[0], None, None, None, None) 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) leader = Leader(leader[1].mzxid, None, None, member)
self.fetch_cluster = member.index == -1 self.fetch_cluster = member.index == -1
self.leader = leader # get last leader operation
if self.fetch_cluster: self.last_leader_operation = self.get_node(self.leader_optime_path) if self.fetch_cluster else None
last_leader_operation = self.get_node('/optime/leader') self.last_leader_operation = 0 if self.last_leader_operation is None else int(self.last_leader_operation[0])
if last_leader_operation: self.cluster = Cluster(initialize, leader, self.last_leader_operation, members)
self.last_leader_operation = int(last_leader_operation[0])
def get_cluster(self): def get_cluster(self):
if self.exhibitor and self.exhibitor.poll(): if self.exhibitor and self.exhibitor.poll():
@@ -163,28 +172,27 @@ class ZooKeeper(AbstractDCS):
logger.exception('get_cluster') logger.exception('get_cluster')
self.session_listener(KazooState.LOST) self.session_listener(KazooState.LOST)
raise ZooKeeperError('ZooKeeper in not responding properly') 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): def _create(self, path, value, **kwargs):
try: try:
self.client.retry(self.client.create, self.client_path(path), value, **kwargs) self.client.retry(self.client.create, path, value, **kwargs)
return True return True
except: except:
return False return False
def attempt_to_acquire_leader(self): def attempt_to_acquire_leader(self):
ret = self._create('/leader', self._name, makepath=True, ephemeral=True) ret = self._create(self.leader_path, self._name, makepath=True, ephemeral=True)
ret or logger.info('Could not take out TTL lock') ret or logger.info('Could not take out TTL lock')
return ret return ret
def race(self, path): def initialize(self):
return self._create(path, self._name, makepath=True) return self._create(self.initialize_path, self._name, makepath=True)
def touch_member(self, connection_string, ttl=None): def touch_member(self, connection_string, ttl=None):
for m in self.members: if self.cluster and any(m.name == self._name for m in self.cluster.members):
if m.name == self._name: return True
return True path = self.member_path
path = self.client_path('/members/' + self._name)
try: try:
self.client.retry(self.client.create, path, connection_string, makepath=True, ephemeral=True) self.client.retry(self.client.create, path, connection_string, makepath=True, ephemeral=True)
return True return True
@@ -204,7 +212,7 @@ class ZooKeeper(AbstractDCS):
last_operation = state_handler.last_operation() last_operation = state_handler.last_operation()
if last_operation != self.last_leader_operation: if last_operation != self.last_leader_operation:
self.last_leader_operation = last_operation self.last_leader_operation = last_operation
path = self.client_path('/optime/leader') path = self.leader_optime_path
try: try:
self.client.retry(self.client.set, path, last_operation) self.client.retry(self.client.set, path, last_operation)
except NoNodeError: except NoNodeError:
@@ -217,8 +225,19 @@ class ZooKeeper(AbstractDCS):
return True return True
def delete_leader(self): def delete_leader(self):
if isinstance(self.leader, Leader) and self.leader.name == self._name: if isinstance(self.cluster, Cluster) and self.cluster.leader.name == self._name:
self.client.delete(self.client_path('/leader')) 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): def watch(self, timeout):
self.cluster_event.wait(timeout) self.cluster_event.wait(timeout)
+1 -1
View File
@@ -47,7 +47,7 @@ postgresql:
env_dir: /home/postgres/etc/wal-e.d/env env_dir: /home/postgres/etc/wal-e.d/env
threshold_megabytes: 10240 threshold_megabytes: 10240
threshold_backup_size_percentage: 30 threshold_backup_size_percentage: 30
restore: scripts/restore.py restore: patroni/scripts/restore.py
#recovery_conf: #recovery_conf:
#restore_command: cp ../wal_archive/%f %p #restore_command: cp ../wal_archive/%f %p
parameters: parameters:
+1 -1
View File
@@ -49,7 +49,7 @@ postgresql:
env_dir: /home/postgres/etc/wal-e.d/env env_dir: /home/postgres/etc/wal-e.d/env
threshold_megabytes: 10240 threshold_megabytes: 10240
threshold_backup_size_percentage: 30 threshold_backup_size_percentage: 30
restore: scripts/restore.py restore: patroni/scripts/restore.py
parameters: parameters:
archive_mode: "on" archive_mode: "on"
wal_level: hot_standby wal_level: hot_standby
+1 -1
View File
@@ -27,5 +27,5 @@ git push
python3 setup.py sdist bdist_wheel upload python3 setup.py sdist bdist_wheel upload
git tag ${version} git tag v${version}
git push --tags git push --tags
+8 -4
View File
@@ -93,9 +93,9 @@ def etcd_delete(key, **kwargs):
def etcd_read(key, **kwargs): def etcd_read(key, **kwargs):
if key == '/service/noleader': if key == '/service/noleader/':
raise DCSError('noleader') raise DCSError('noleader')
elif key == '/service/nocluster': elif key == '/service/nocluster/':
raise etcd.EtcdKeyNotFound raise etcd.EtcdKeyNotFound
response = {"action": "get", "node": {"key": "/service/batman5", "dir": True, "nodes": [ response = {"action": "get", "node": {"key": "/service/batman5", "dir": True, "nodes": [
@@ -266,8 +266,12 @@ class TestEtcd(unittest.TestCase):
def test_update_leader(self): def test_update_leader(self):
self.assertTrue(self.etcd.update_leader(MockPostgresql())) self.assertTrue(self.etcd.update_leader(MockPostgresql()))
def test_race(self): def test_initialize(self):
self.assertFalse(self.etcd.race('')) 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): def test_delete_leader(self):
self.etcd.client.delete = etcd_delete self.etcd.client.delete = etcd_delete
+57 -6
View File
@@ -9,8 +9,9 @@ import yaml
from mock import Mock, patch from mock import Mock, patch
from patroni.api import RestApiServer 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.etcd import Etcd
from patroni.exceptions import PostgresException
from patroni import Patroni, main from patroni import Patroni, main
from patroni.zookeeper import ZooKeeper from patroni.zookeeper import ZooKeeper
from six.moves import BaseHTTPServer from six.moves import BaseHTTPServer
@@ -41,6 +42,30 @@ class Mock_BaseServer__is_shut_down:
pass 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): class TestPatroni(unittest.TestCase):
def __init__(self, method_name='runTest'): def __init__(self, method_name='runTest'):
@@ -50,6 +75,7 @@ class TestPatroni(unittest.TestCase):
def set_up(self): def set_up(self):
self.touched = False self.touched = False
self.init_cancelled = False
subprocess.call = subprocess_call subprocess.call = subprocess_call
psycopg2.connect = psycopg2_connect psycopg2.connect = psycopg2_connect
self.time_sleep = time.sleep self.time_sleep = time.sleep
@@ -126,24 +152,49 @@ class TestPatroni(unittest.TestCase):
self.p.touch_member() self.p.touch_member()
def test_patroni_initialize(self): 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.write = etcd_write
self.p.ha.dcs.client.read = etcd_read
self.p.touch_member = self.touch_member self.p.touch_member = self.touch_member
self.p.postgresql.data_directory_empty = true 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.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 time.sleep = time_sleep
self.p.ha.dcs.client.read = etcd_read self.p.ha.dcs.client.read = etcd_read
self.p.initialize() self.p.initialize()
self.p.ha.dcs.current_leader = nop self.p.ha.dcs.get_cluster = get_cluster_initialized_without_leader
self.assertRaises(Exception, self.p.initialize) self.assertRaises(SleepException, self.p.initialize)
self.p.postgresql.data_directory_empty = false self.p.postgresql.data_directory_empty = false
self.p.initialize() 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): def test_schedule_next_run(self):
self.p.next_run = time.time() - self.p.nap_time - 1 self.p.next_run = time.time() - self.p.nap_time - 1
self.p.schedule_next_run() 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.dcs import Cluster, Leader, Member
from patroni.postgresql import Postgresql from patroni.postgresql import Postgresql
from test_ha import true, false
def nop(*args, **kwargs): def nop(*args, **kwargs):
@@ -217,3 +218,9 @@ class TestPostgresql(unittest.TestCase):
self.p.start() self.p.start()
self.p.query = self.mock_query self.p.query = self.mock_query
self.assertTrue(self.p.stop()) 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()
+17 -6
View File
@@ -56,10 +56,8 @@ class MockKazooClient:
func(*args, **kwargs) func(*args, **kwargs)
def get(self, path, watch=None): def get(self, path, watch=None):
if path == '/service/test/no_node': if path == '/no_node':
raise NoNodeError raise NoNodeError
elif path == '/service/test/other_exception':
raise Exception()
elif '/members/' in path: elif '/members/' in path:
return ( return (
'postgres://repuser:rep-pass@localhost:5434/postgres?application_name=http://127.0.0.1:8009/patroni', '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: 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, -1, 0, 0, 0))
return ('foo', ZnodeStat(0, 0, 0, 0, 0, 0, 0, 0, 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): 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'] return ['foo', 'bar', 'buzz']
def create(self, path, value="", acl=None, ephemeral=False, sequence=False, makepath=False): def create(self, path, value="", acl=None, ephemeral=False, sequence=False, makepath=False):
@@ -93,6 +97,8 @@ class MockKazooClient:
return return
self.leader = True self.leader = True
raise Exception raise Exception
elif path.endswith('/initialize'):
raise NoNodeError
def set_hosts(self, hosts, randomize_hosts=None): def set_hosts(self, hosts, randomize_hosts=None):
pass pass
@@ -132,7 +138,9 @@ class TestZooKeeper(unittest.TestCase):
def test_get_node(self): def test_get_node(self):
self.assertIsNone(self.zk.get_node('/no_node')) 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): def test__inner_load_cluster(self):
self.zk._base_path = self.zk._base_path.replace('test', 'bla') 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.touch_member('foo')
self.zk.delete_leader() self.zk.delete_leader()
def test_race(self): def test_initialize(self):
self.assertFalse(self.zk.race('/initialize')) self.assertFalse(self.zk.initialize())
def test_cancel_initialization(self):
self.zk.cancel_initialization()
def test_touch_member(self): def test_touch_member(self):
self.zk.touch_member('new') self.zk.touch_member('new')