From 6079c8359fbaa91d977ac9e124c34dd0a9f76f1e Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 2 Jul 2015 12:05:38 +0200 Subject: [PATCH 1/9] refactoring: preparing to support ZooKeeper move some common classes info separate file in preparation to support distributed configuration store other then Etcd --- helpers/dcs.py | 32 ++++++++++++++++++++++++++++++++ helpers/etcd.py | 31 ++----------------------------- tests/test_etcd.py | 3 ++- 3 files changed, 36 insertions(+), 30 deletions(-) create mode 100644 helpers/dcs.py diff --git a/helpers/dcs.py b/helpers/dcs.py new file mode 100644 index 00000000..e4a6e895 --- /dev/null +++ b/helpers/dcs.py @@ -0,0 +1,32 @@ +import sys + +from collections import namedtuple +from helpers.utils import calculate_ttl + +if sys.hexversion >= 0x03000000: + from urllib.parse import urlparse, urlunparse, parse_qsl +else: + from urlparse import urlparse, urlunparse, parse_qsl + + +class Member(namedtuple('Member', 'name,conn_url,api_url,expiration,ttl')): + + @staticmethod + def fromNode(node): + scheme, netloc, path, params, query, fragment = urlparse(node['value']) + conn_url = urlunparse((scheme, netloc, path, params, '', fragment)) + api_url = ([v for n, v in parse_qsl(query) if n == 'application_name'] or [None])[0] + expiration = node.get('expiration', None) + ttl = node.get('ttl', None) + return Member(node['key'].split('/')[-1], conn_url, api_url, expiration, ttl) + + def real_ttl(self): + return calculate_ttl(self.expiration) or -1 + + +class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members')): + + def is_unlocked(self): + return not (self.leader and self.leader.name) + + diff --git a/helpers/etcd.py b/helpers/etcd.py index 3385fbe6..390b8591 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -2,44 +2,17 @@ import logging import random import requests import socket -import sys -from collections import namedtuple from dns.exception import DNSException from dns import resolver +from helpers.dcs import Cluster, Member from helpers.errors import CurrentLeaderError, EtcdError, EtcdConnectionFailed -from helpers.utils import calculate_ttl, sleep +from helpers.utils import sleep from requests.exceptions import RequestException -if sys.hexversion >= 0x03000000: - from urllib.parse import urlparse, urlunparse, parse_qsl -else: - from urlparse import urlparse, urlunparse, parse_qsl - logger = logging.getLogger(__name__) -class Member(namedtuple('Member', 'name,conn_url,api_url,expiration,ttl')): - - @staticmethod - def fromNode(node): - scheme, netloc, path, params, query, fragment = urlparse(node['value']) - conn_url = urlunparse((scheme, netloc, path, params, '', fragment)) - api_url = ([v for n, v in parse_qsl(query) if n == 'application_name'] or [None])[0] - expiration = node.get('expiration', None) - ttl = node.get('ttl', None) - return Member(node['key'].split('/')[-1], conn_url, api_url, expiration, ttl) - - def real_ttl(self): - return calculate_ttl(self.expiration) or -1 - - -class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members')): - - def is_unlocked(self): - return not (self.leader and self.leader.name) - - class Client: API_VERSION = 'v2' diff --git a/tests/test_etcd.py b/tests/test_etcd.py index 11cfbdb5..e37486ac 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -8,7 +8,8 @@ import unittest from dns.exception import DNSException from helpers.errors import EtcdError, CurrentLeaderError, EtcdConnectionFailed -from helpers.etcd import Client, Cluster, Etcd, Member +from helpers.dcs import Cluster, Member +from helpers.etcd import Client, Etcd class MockResponse: From cc71906009496f2796ebfb57cb6e4d9066786a7e Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 2 Jul 2015 15:15:07 +0200 Subject: [PATCH 2/9] Inherit Etcd from abstract class --- helpers/dcs.py | 73 +++++++++++++++++++++++++++++++--------- helpers/etcd.py | 47 ++++++++++++++++++-------- helpers/ha.py | 8 ++--- tests/test_etcd.py | 9 +++-- tests/test_governor.py | 4 +-- tests/test_ha.py | 6 ++-- tests/test_postgresql.py | 8 ++--- 7 files changed, 108 insertions(+), 47 deletions(-) diff --git a/helpers/dcs.py b/helpers/dcs.py index e4a6e895..0b6c9332 100644 --- a/helpers/dcs.py +++ b/helpers/dcs.py @@ -1,24 +1,19 @@ -import sys +import abc from collections import namedtuple from helpers.utils import calculate_ttl -if sys.hexversion >= 0x03000000: - from urllib.parse import urlparse, urlunparse, parse_qsl -else: - from urlparse import urlparse, urlunparse, parse_qsl + +class DCSError(Exception): + + def __init__(self, value): + self.value = value + + def __str__(self): + return repr(self.value) -class Member(namedtuple('Member', 'name,conn_url,api_url,expiration,ttl')): - - @staticmethod - def fromNode(node): - scheme, netloc, path, params, query, fragment = urlparse(node['value']) - conn_url = urlunparse((scheme, netloc, path, params, '', fragment)) - api_url = ([v for n, v in parse_qsl(query) if n == 'application_name'] or [None])[0] - expiration = node.get('expiration', None) - ttl = node.get('ttl', None) - return Member(node['key'].split('/')[-1], conn_url, api_url, expiration, ttl) +class Member(namedtuple('Member', 'index,name,conn_url,api_url,expiration,ttl')): def real_ttl(self): return calculate_ttl(self.expiration) or -1 @@ -30,3 +25,51 @@ class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,mem return not (self.leader and self.leader.name) +class AbstractDCS: + + __metaclass__ = abc.ABCMeta + + def __init__(self, config): + self._base_path = '/service/' + config['scope'] + + def client_path(self, path): + return self._base_path + path + + @abc.abstractmethod + def get_cluster(self): + raise NotImplementedError + + @abc.abstractmethod + def update_leader(self, state_handler): + raise NotImplementedError + + @abc.abstractmethod + def attempt_to_acquire_leader(self, value): + raise NotImplementedError + + def current_leader(self): + try: + cluster = self.get_cluster() + return None if cluster.is_unlocked() else cluster.leader + except DCSError: + return None + + @abc.abstractmethod + def touch_member(self, member, connection_string, ttl=None): + raise NotImplementedError + + @abc.abstractmethod + def take_leader(self, value): + raise NotImplementedError + + @abc.abstractmethod + def race(self, path, value): + raise NotImplementedError + + @abc.abstractmethod + def delete_member(self, member): + raise NotImplementedError + + @abc.abstractmethod + def delete_leader(self, value): + raise NotImplementedError diff --git a/helpers/etcd.py b/helpers/etcd.py index 390b8591..72c18e1c 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -2,17 +2,34 @@ import logging import random import requests import socket +import sys from dns.exception import DNSException from dns import resolver -from helpers.dcs import Cluster, Member -from helpers.errors import CurrentLeaderError, EtcdError, EtcdConnectionFailed +from helpers.dcs import AbstractDCS, Cluster, DCSError, Member from helpers.utils import sleep from requests.exceptions import RequestException +if sys.hexversion >= 0x03000000: + from urllib.parse import urlparse, urlunparse, parse_qsl +else: + from urlparse import urlparse, urlunparse, parse_qsl + logger = logging.getLogger(__name__) +class EtcdError(DCSError): + pass + + +class CurrentLeaderError(EtcdError): + pass + + +class EtcdConnectionFailed(EtcdError): + pass + + class Client: API_VERSION = 'v2' @@ -177,12 +194,12 @@ class Client: return False -class Etcd: +class Etcd(AbstractDCS): def __init__(self, config): + super(Etcd, self).__init__(config) self.ttl = config['ttl'] self.member_ttl = config.get('member_ttl', 3600) - self._base_path = '/keys/service/' + config['scope'] self.client = self.get_etcd_client(config) def get_etcd_client(self, config): @@ -196,7 +213,7 @@ class Etcd: return client def client_path(self, path): - return self._base_path + path + return '/keys' + super(Etcd, self).client_path(path) def get_client_path(self, path): return self.client.get(self.client_path(path)) @@ -224,6 +241,15 @@ class Etcd: return n return None + @staticmethod + def member(node): + scheme, netloc, path, params, query, fragment = urlparse(node['value']) + conn_url = urlunparse((scheme, netloc, path, params, '', fragment)) + api_url = ([v for n, v in parse_qsl(query) if n == 'application_name'] or [None])[0] + expiration = node.get('expiration', None) + ttl = node.get('ttl', None) + return Member(node['modifiedIndex'], node['key'].split('/')[-1], conn_url, api_url, expiration, ttl) + def get_cluster(self): try: response, status_code = self.get_client_path('?recursive=true') @@ -232,7 +258,7 @@ class Etcd: initialize = True if node else False # get list of members node = self.find_node(response['node'], '/members') or {'nodes': []} - members = [Member.fromNode(n) for n in node['nodes']] + members = [self.member(n) for n in node['nodes']] # get last leader operation last_leader_operation = 0 @@ -251,7 +277,7 @@ class Etcd: leader = m break if not leader: - leader = Member(node['value'], None, None, None, None) + leader = Member(-1, node['value'], None, None, None, None) return Cluster(initialize, leader, last_leader_operation, members) elif status_code == 404: @@ -261,13 +287,6 @@ class Etcd: raise EtcdError('Etcd is not responding properly') - def current_leader(self): - try: - cluster = self.get_cluster() - return None if cluster.is_unlocked() else cluster.leader - except EtcdError: - raise CurrentLeaderError('Etcd is not responding properly') - def touch_member(self, member, connection_string, ttl=None): try: return self.put_client_path('/members/' + member, value=connection_string, ttl=ttl or self.member_ttl) diff --git a/helpers/ha.py b/helpers/ha.py index 4f0e61bc..52462b92 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -1,6 +1,6 @@ import logging -from helpers.errors import EtcdError +from helpers.dcs import DCSError from psycopg2 import InterfaceError, OperationalError logger = logging.getLogger(__name__) @@ -84,10 +84,10 @@ class Ha: else: self.follow_the_leader() return 'no action. i am a secondary and i am following a leader' - except EtcdError: - logger.error('Error communicating with Etcd') + except DCSError as e: + logger.error('Error communicating with DCS') if self.state_handler.is_leader(): self.state_handler.demote(None) - return 'demoted self because etcd is not accessible and i was a leader' + return 'demoted self because DCS is not accessible and i was a leader' except (InterfaceError, OperationalError): logger.error('Error communicating with Postgresql. Will try again') diff --git a/tests/test_etcd.py b/tests/test_etcd.py index e37486ac..a5416c45 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -7,9 +7,8 @@ import time import unittest from dns.exception import DNSException -from helpers.errors import EtcdError, CurrentLeaderError, EtcdConnectionFailed from helpers.dcs import Cluster, Member -from helpers.etcd import Client, Etcd +from helpers.etcd import Client, CurrentLeaderError, Etcd, EtcdConnectionFailed, EtcdError class MockResponse: @@ -106,9 +105,9 @@ class TestMember(unittest.TestCase): def test_real_ttl(self): now = datetime.datetime.utcnow() - member = Member('a', 'b', 'c', (now + datetime.timedelta(seconds=2)).strftime('%Y-%m-%dT%H:%M:%S.%fZ'), None) + member = Member(0, 'a', 'b', 'c', (now + datetime.timedelta(seconds=2)).strftime('%Y-%m-%dT%H:%M:%S.%fZ'), None) self.assertLess(member.real_ttl(), 2) - self.assertEquals(Member('a', 'b', 'c', '', None).real_ttl(), -1) + self.assertEquals(Member(0, 'a', 'b', 'c', '', None).real_ttl(), -1) class TestClient(unittest.TestCase): @@ -203,7 +202,7 @@ class TestEtcd(unittest.TestCase): self.etcd.get_cluster() def test_current_leader(self): - self.assertRaises(CurrentLeaderError, self.etcd.current_leader) + self.assertIsNone(self.etcd.current_leader()) def test_touch_member(self): self.assertFalse(self.etcd.touch_member('', '')) diff --git a/tests/test_governor.py b/tests/test_governor.py index b850deda..0a8d4ab8 100644 --- a/tests/test_governor.py +++ b/tests/test_governor.py @@ -8,7 +8,7 @@ import unittest import yaml from governor import Governor, main -from helpers.etcd import Cluster, Member +from helpers.dcs import Cluster, Member from test_ha import true, false from test_postgresql import Postgresql, subprocess_call, psycopg2_connect from test_etcd import requests_get, requests_put, requests_delete @@ -77,7 +77,7 @@ class TestGovernor(unittest.TestCase): def test_touch_member(self): now = datetime.datetime.utcnow() - member = Member(self.g.postgresql.name, 'b', 'c', (now + datetime.timedelta( + member = Member(0, self.g.postgresql.name, 'b', 'c', (now + datetime.timedelta( seconds=self.g.shutdown_member_ttl + 10)).strftime('%Y-%m-%dT%H:%M:%S.%fZ'), None) self.g.ha.cluster = Cluster(True, member, 0, [member]) self.g.touch_member() diff --git a/tests/test_ha.py b/tests/test_ha.py index 574f949f..e45885e8 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -1,7 +1,7 @@ import unittest import requests -from helpers.errors import EtcdError +from helpers.dcs import DCSError from helpers.etcd import Cluster, Etcd from helpers.ha import Ha from test_etcd import requests_get, requests_put, requests_delete @@ -57,7 +57,7 @@ def nop(*args, **kwargs): def dead_etcd(): - raise EtcdError('Etcd is not responding properly') + raise DCSError('Etcd is not responding properly') class TestHa(unittest.TestCase): @@ -134,4 +134,4 @@ class TestHa(unittest.TestCase): def test_no_etcd_connection_master_demote(self): self.ha.load_cluster_from_etcd = dead_etcd - self.assertEquals(self.ha.run_cycle(), 'demoted self because etcd is not accessible and i was a leader') + self.assertEquals(self.ha.run_cycle(), 'demoted self because DCS is not accessible and i was a leader') diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 7fae838a..6948b9ff 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -4,7 +4,7 @@ import shutil import subprocess import unittest -from helpers.etcd import Cluster, Member +from helpers.dcs import Cluster, Member from helpers.postgresql import Postgresql @@ -101,9 +101,9 @@ class TestPostgresql(unittest.TestCase): psycopg2.connect = psycopg2_connect if not os.path.exists(self.p.data_dir): os.makedirs(self.p.data_dir) - self.leader = Member('leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', None, None, 28) - self.other = Member('test1', 'postgres://replicator:rep-pass@127.0.0.1:5433/postgres', None, None, 28) - self.me = Member('test0', 'postgres://replicator:rep-pass@127.0.0.1:5434/postgres', None, None, 28) + self.leader = Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', None, None, 28) + self.other = Member(0, 'test1', 'postgres://replicator:rep-pass@127.0.0.1:5433/postgres', None, None, 28) + self.me = Member(0, 'test0', 'postgres://replicator:rep-pass@127.0.0.1:5434/postgres', None, None, 28) def tear_down(self): shutil.rmtree('data') From ae689a2f83692638001f8535f3d72c03faca0f11 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 2 Jul 2015 15:33:21 +0200 Subject: [PATCH 3/9] Rename ha.etcd into ha.dcs --- governor.py | 19 ++++++++++++------- helpers/errors.py | 15 --------------- helpers/ha.py | 20 ++++++++++---------- tests/test_governor.py | 6 +++--- tests/test_ha.py | 6 +++--- 5 files changed, 28 insertions(+), 38 deletions(-) delete mode 100644 helpers/errors.py diff --git a/governor.py b/governor.py index eb029b22..cfe1c546 100755 --- a/governor.py +++ b/governor.py @@ -16,14 +16,19 @@ class Governor: def __init__(self, config): self.nap_time = config['loop_wait'] - self.etcd = Etcd(config['etcd']) self.postgresql = Postgresql(config['postgresql']) - self.ha = Ha(self.postgresql, self.etcd) + self.ha = Ha(self.postgresql, self.get_dcs(config)) host, port = config['restapi']['listen'].split(':') self.api = RestApiServer(self, config['restapi']) self.next_run = time.time() self.shutdown_member_ttl = 300 + @staticmethod + def get_dcs(config): + if 'etcd' in config: + return Etcd(config['etcd']) + raise Exception('Can not find sutable configuration of distributed configuration store') + def touch_member(self, ttl=None): connection_string = self.postgresql.connection_string + '?application_name=' + self.api.connection_string if self.ha.cluster: @@ -31,7 +36,7 @@ class Governor: # Do not update member TTL when it is far from being expired if m.name == self.postgresql.name and m.real_ttl() > self.shutdown_member_ttl: return True - return self.etcd.touch_member(self.postgresql.name, connection_string, ttl) + return self.ha.dcs.touch_member(self.postgresql.name, connection_string, ttl) def initialize(self): # wait for etcd to be available @@ -42,14 +47,14 @@ class Governor: # is data directory empty? if self.postgresql.data_directory_empty(): # racing to initialize - if self.etcd.race('/initialize', self.postgresql.name): + if self.ha.dcs.race('/initialize', self.postgresql.name): self.postgresql.initialize() - self.etcd.take_leader(self.postgresql.name) + self.ha.dcs.take_leader(self.postgresql.name) self.postgresql.start() self.postgresql.create_replication_user() else: while True: - leader = self.etcd.current_leader() + leader = self.ha.dcs.current_leader() if leader and self.postgresql.sync_from_leader(leader): self.postgresql.write_recovery_conf(leader) self.postgresql.start() @@ -105,7 +110,7 @@ def main(): finally: governor.touch_member(governor.shutdown_member_ttl) # schedule member removal governor.postgresql.stop() - governor.etcd.delete_leader(governor.postgresql.name) + governor.ha.dcs.delete_leader(governor.postgresql.name) if __name__ == '__main__': diff --git a/helpers/errors.py b/helpers/errors.py deleted file mode 100644 index 3a56a1e0..00000000 --- a/helpers/errors.py +++ /dev/null @@ -1,15 +0,0 @@ -class EtcdError(Exception): - - def __init__(self, value): - self.value = value - - def __str__(self): - return repr(self.value) - - -class CurrentLeaderError(EtcdError): - pass - - -class EtcdConnectionFailed(EtcdError): - pass diff --git a/helpers/ha.py b/helpers/ha.py index 52462b92..1ffa06ae 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -10,17 +10,17 @@ class Ha: def __init__(self, state_handler, etcd): self.state_handler = state_handler - self.etcd = etcd + self.dcs = etcd self.cluster = None - def load_cluster_from_etcd(self): - self.cluster = self.etcd.get_cluster() + def load_cluster_from_dcs(self): + self.cluster = self.dcs.get_cluster() def acquire_lock(self): - return self.etcd.attempt_to_acquire_leader(self.state_handler.name) + return self.dcs.attempt_to_acquire_leader(self.state_handler.name) def update_lock(self): - return self.etcd.update_leader(self.state_handler) + return self.dcs.update_leader(self.state_handler) def has_lock(self): lock_owner = self.cluster.leader and self.cluster.leader.name @@ -35,7 +35,7 @@ class Ha: def run_cycle(self): try: - self.load_cluster_from_etcd() + 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) @@ -43,7 +43,7 @@ class Ha: if not has_lock: return 'started as a secondary' logger.info('started as readonly because i had the session lock') - self.load_cluster_from_etcd() + self.load_cluster_from_dcs() if self.cluster.is_unlocked(): if self.state_handler.is_healthiest_node(self.cluster): @@ -54,7 +54,7 @@ class Ha: self.state_handler.promote() return 'promoted self to leader by acquiring session lock' else: - self.load_cluster_from_etcd() + self.load_cluster_from_dcs() if self.state_handler.is_leader(): self.demote() return 'demoted self due after trying and failing to obtain lock' @@ -62,7 +62,7 @@ class Ha: self.follow_the_leader() return 'following new leader after trying and failing to obtain lock' else: - self.load_cluster_from_etcd() + self.load_cluster_from_dcs() if self.state_handler.is_leader(): self.demote() return 'demoting self because i am not the healthiest node' @@ -84,7 +84,7 @@ class Ha: else: self.follow_the_leader() return 'no action. i am a secondary and i am following a leader' - except DCSError as e: + except DCSError: logger.error('Error communicating with DCS') if self.state_handler.is_leader(): self.state_handler.demote(None) diff --git a/tests/test_governor.py b/tests/test_governor.py index 0a8d4ab8..daced98b 100644 --- a/tests/test_governor.py +++ b/tests/test_governor.py @@ -83,11 +83,11 @@ class TestGovernor(unittest.TestCase): self.g.touch_member() def test_governor_initialize(self): - self.g.etcd.client._base_uri = 'http://remote' + self.g.ha.dcs.client._base_uri = 'http://remote' self.g.postgresql.data_directory_empty = true - self.g.etcd.race = true + self.g.ha.dcs.race = true self.g.initialize() - self.g.etcd.race = false + self.g.ha.dcs.race = false self.g.initialize() self.g.postgresql.data_directory_empty = false self.g.touch_member = self.touch_member diff --git a/tests/test_ha.py b/tests/test_ha.py index e45885e8..299bc4b6 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -73,9 +73,9 @@ class TestHa(unittest.TestCase): self.p = MockPostgresql() self.e = Etcd({'ttl': 30, 'host': 'remotehost:2379', 'scope': 'test'}) self.ha = Ha(self.p, self.e) - self.ha.load_cluster_from_etcd() + self.ha.load_cluster_from_dcs() self.ha.cluster = Cluster(False, None, None, []) - self.ha.load_cluster_from_etcd = nop + self.ha.load_cluster_from_dcs = nop def test_start_as_slave(self): self.p.is_healthy = false @@ -133,5 +133,5 @@ class TestHa(unittest.TestCase): self.assertEquals(self.ha.run_cycle(), 'no action. i am a secondary and i am following a leader') def test_no_etcd_connection_master_demote(self): - self.ha.load_cluster_from_etcd = dead_etcd + self.ha.load_cluster_from_dcs = dead_etcd self.assertEquals(self.ha.run_cycle(), 'demoted self because DCS is not accessible and i was a leader') From 1531304baad444f92d843a8e0e5e3d9fed2f20e2 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 7 Jul 2015 12:23:06 +0200 Subject: [PATCH 4/9] return result of last_operation as a string --- helpers/postgresql.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 8d615a0c..7e1ba9ca 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -289,4 +289,4 @@ primary_conninfo = '{}' self.sync_replication_slots([]) def last_operation(self): - return self.xlog_position() + return str(self.xlog_position()) From c5edcb589b13c193bf2e9c764e1582ad0bdda841 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 7 Jul 2015 12:24:02 +0200 Subject: [PATCH 5/9] return result of last_operation as a string --- tests/test_postgresql.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 6948b9ff..f6019f13 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -174,4 +174,4 @@ class TestPostgresql(unittest.TestCase): self.assertTrue(self.p.promote()) def test_last_operation(self): - self.assertEquals(self.p.last_operation(), 0) + self.assertEquals(self.p.last_operation(), '0') From 74455905466377d79de5d6262a410edb8a6deea1 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 7 Jul 2015 12:26:52 +0200 Subject: [PATCH 6/9] Keep node name in AbstractDCS class It eliminates need to pass this name into most of the methods --- governor.py | 18 ++++++++-------- helpers/dcs.py | 51 ++++++++++++++++++++++++++++++---------------- helpers/etcd.py | 39 +++++++++++++---------------------- helpers/ha.py | 2 +- tests/test_etcd.py | 13 +++++------- tests/test_ha.py | 2 +- 6 files changed, 64 insertions(+), 61 deletions(-) diff --git a/governor.py b/governor.py index cfe1c546..ec4b39ad 100755 --- a/governor.py +++ b/governor.py @@ -17,16 +17,16 @@ class Governor: def __init__(self, config): self.nap_time = config['loop_wait'] self.postgresql = Postgresql(config['postgresql']) - self.ha = Ha(self.postgresql, self.get_dcs(config)) + self.ha = Ha(self.postgresql, self.get_dcs(self.postgresql.name, config)) host, port = config['restapi']['listen'].split(':') self.api = RestApiServer(self, config['restapi']) self.next_run = time.time() self.shutdown_member_ttl = 300 @staticmethod - def get_dcs(config): + def get_dcs(name, config): if 'etcd' in config: - return Etcd(config['etcd']) + return Etcd(name, config['etcd']) raise Exception('Can not find sutable configuration of distributed configuration store') def touch_member(self, ttl=None): @@ -36,20 +36,20 @@ class Governor: # Do not update member TTL when it is far from being expired if m.name == self.postgresql.name and m.real_ttl() > self.shutdown_member_ttl: return True - return self.ha.dcs.touch_member(self.postgresql.name, connection_string, ttl) + return self.ha.dcs.touch_member(connection_string, ttl) def initialize(self): # wait for etcd to be available while not self.touch_member(): - logging.info('waiting on etcd') + logging.info('waiting on DCS') sleep(5) # is data directory empty? if self.postgresql.data_directory_empty(): # racing to initialize - if self.ha.dcs.race('/initialize', self.postgresql.name): + if self.ha.dcs.race('/initialize'): self.postgresql.initialize() - self.ha.dcs.take_leader(self.postgresql.name) + self.ha.dcs.take_leader() self.postgresql.start() self.postgresql.create_replication_user() else: @@ -70,7 +70,7 @@ class Governor: if nap_time <= 0: self.next_run = current_time else: - sleep(nap_time) + self.ha.dcs.sleep(nap_time) def run(self): self.api.start() @@ -110,7 +110,7 @@ def main(): finally: governor.touch_member(governor.shutdown_member_ttl) # schedule member removal governor.postgresql.stop() - governor.ha.dcs.delete_leader(governor.postgresql.name) + governor.ha.dcs.delete_leader() if __name__ == '__main__': diff --git a/helpers/dcs.py b/helpers/dcs.py index 0b6c9332..a5e43488 100644 --- a/helpers/dcs.py +++ b/helpers/dcs.py @@ -1,7 +1,20 @@ import abc +import sys from collections import namedtuple -from helpers.utils import calculate_ttl +from helpers.utils import calculate_ttl, sleep + +if sys.hexversion >= 0x03000000: + from urllib.parse import urlparse, urlunparse, parse_qsl +else: + from urlparse import urlparse, urlunparse, parse_qsl + + +def parse_connection_string(value): + scheme, netloc, path, params, query, fragment = urlparse(value) + conn_url = urlunparse((scheme, netloc, path, params, '', fragment)) + api_url = ([v for n, v in parse_qsl(query) if n == 'application_name'] or [None])[0] + return conn_url, api_url class DCSError(Exception): @@ -10,6 +23,10 @@ class DCSError(Exception): self.value = value def __str__(self): + """ + >>> str(DCSError('foo')) + "'foo'" + """ return repr(self.value) @@ -29,7 +46,8 @@ class AbstractDCS: __metaclass__ = abc.ABCMeta - def __init__(self, config): + def __init__(self, name, config): + self._name = name self._base_path = '/service/' + config['scope'] def client_path(self, path): @@ -37,15 +55,15 @@ class AbstractDCS: @abc.abstractmethod def get_cluster(self): - raise NotImplementedError + """get_cluster""" @abc.abstractmethod def update_leader(self, state_handler): - raise NotImplementedError + """update_leader""" @abc.abstractmethod - def attempt_to_acquire_leader(self, value): - raise NotImplementedError + def attempt_to_acquire_leader(self): + """attempt_to_acquire_leader""" def current_leader(self): try: @@ -55,21 +73,20 @@ class AbstractDCS: return None @abc.abstractmethod - def touch_member(self, member, connection_string, ttl=None): - raise NotImplementedError + def touch_member(self, connection_string, ttl=None): + """touch_member""" @abc.abstractmethod - def take_leader(self, value): - raise NotImplementedError + def take_leader(self): + """take_leader""" @abc.abstractmethod - def race(self, path, value): - raise NotImplementedError + def race(self, path): + """race""" @abc.abstractmethod - def delete_member(self, member): - raise NotImplementedError + def delete_leader(self): + """delete_leader""" - @abc.abstractmethod - def delete_leader(self, value): - raise NotImplementedError + def sleep(self, timeout): + sleep(timeout) diff --git a/helpers/etcd.py b/helpers/etcd.py index 72c18e1c..1e971f30 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -2,19 +2,13 @@ import logging import random import requests import socket -import sys from dns.exception import DNSException from dns import resolver -from helpers.dcs import AbstractDCS, Cluster, DCSError, Member +from helpers.dcs import AbstractDCS, Cluster, DCSError, Member, parse_connection_string from helpers.utils import sleep from requests.exceptions import RequestException -if sys.hexversion >= 0x03000000: - from urllib.parse import urlparse, urlunparse, parse_qsl -else: - from urlparse import urlparse, urlunparse, parse_qsl - logger = logging.getLogger(__name__) @@ -196,8 +190,8 @@ class Client: class Etcd(AbstractDCS): - def __init__(self, config): - super(Etcd, self).__init__(config) + def __init__(self, name, config): + super(Etcd, self).__init__(name, config) self.ttl = config['ttl'] self.member_ttl = config.get('member_ttl', 3600) self.client = self.get_etcd_client(config) @@ -243,9 +237,7 @@ class Etcd(AbstractDCS): @staticmethod def member(node): - scheme, netloc, path, params, query, fragment = urlparse(node['value']) - conn_url = urlunparse((scheme, netloc, path, params, '', fragment)) - api_url = ([v for n, v in parse_qsl(query) if n == 'application_name'] or [None])[0] + conn_url, api_url = parse_connection_string(node['value']) expiration = node.get('expiration', None) ttl = node.get('ttl', None) return Member(node['modifiedIndex'], node['key'].split('/')[-1], conn_url, api_url, expiration, ttl) @@ -287,21 +279,21 @@ class Etcd(AbstractDCS): raise EtcdError('Etcd is not responding properly') - def touch_member(self, member, connection_string, ttl=None): + def touch_member(self, connection_string, ttl=None): try: - return self.put_client_path('/members/' + member, value=connection_string, ttl=ttl or self.member_ttl) + return self.put_client_path('/members/' + self._name, value=connection_string, ttl=ttl or self.member_ttl) except EtcdError: return False - def take_leader(self, value): + def take_leader(self): try: - return self.put_client_path('/leader', value=value, ttl=self.ttl) + return self.put_client_path('/leader', value=self._name, ttl=self.ttl) except EtcdError: return False - def attempt_to_acquire_leader(self, value): + def attempt_to_acquire_leader(self): try: - ret = self.put_client_path('/leader', value=value, ttl=self.ttl, prevExist=False) + ret = self.put_client_path('/leader', value=self._name, ttl=self.ttl, prevExist=False) ret or logger.info('Could not take out TTL lock') return ret except EtcdError: @@ -316,14 +308,11 @@ class Etcd(AbstractDCS): return True return False - def race(self, path, value): + def race(self, path): try: - return self.put_client_path(path, value=value, prevExist=False) + return self.put_client_path(path, value=self._name, prevExist=False) except EtcdError: return False - def delete_member(self, member): - return self.delete_client_path('/members/' + member) - - def delete_leader(self, value): - return self.delete_client_path('/leader?prevValue=' + value) + def delete_leader(self): + return self.delete_client_path('/leader?prevValue=' + self._name) diff --git a/helpers/ha.py b/helpers/ha.py index 1ffa06ae..3d48505c 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -17,7 +17,7 @@ class Ha: self.cluster = self.dcs.get_cluster() def acquire_lock(self): - return self.dcs.attempt_to_acquire_leader(self.state_handler.name) + return self.dcs.attempt_to_acquire_leader() def update_lock(self): return self.dcs.update_leader(self.state_handler) diff --git a/tests/test_etcd.py b/tests/test_etcd.py index a5416c45..c0361a6f 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -8,7 +8,7 @@ import unittest from dns.exception import DNSException from helpers.dcs import Cluster, Member -from helpers.etcd import Client, CurrentLeaderError, Etcd, EtcdConnectionFailed, EtcdError +from helpers.etcd import Client, Etcd, EtcdConnectionFailed, EtcdError class MockResponse: @@ -176,7 +176,7 @@ class TestEtcd(unittest.TestCase): requests.put = requests_put requests.delete = requests_delete time.sleep = time_sleep - self.etcd = Etcd({'ttl': 30, 'host': 'localhost:2379', 'scope': 'test'}) + self.etcd = Etcd('foo', {'ttl': 30, 'host': 'localhost:2379', 'scope': 'test'}) def test_get_etcd_client(self): time.sleep = time_sleep_exception @@ -208,10 +208,10 @@ class TestEtcd(unittest.TestCase): self.assertFalse(self.etcd.touch_member('', '')) def test_take_leader(self): - self.assertFalse(self.etcd.take_leader('')) + self.assertFalse(self.etcd.take_leader()) def test_attempt_to_acquire_leader(self): - self.assertFalse(self.etcd.attempt_to_acquire_leader('')) + self.assertFalse(self.etcd.attempt_to_acquire_leader()) def test_update_leader(self): url = self.etcd.client._base_uri = self.etcd.client._base_uri.replace('local', 'remote') @@ -220,7 +220,4 @@ class TestEtcd(unittest.TestCase): self.assertFalse(self.etcd.update_leader(MockPostgresql())) def test_race(self): - self.assertFalse(self.etcd.race('', '')) - - def test_delete_member(self): - self.assertFalse(self.etcd.delete_member('')) + self.assertFalse(self.etcd.race('')) diff --git a/tests/test_ha.py b/tests/test_ha.py index 299bc4b6..fcb010c4 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -71,7 +71,7 @@ class TestHa(unittest.TestCase): requests.put = requests_put requests.delete = requests_delete self.p = MockPostgresql() - self.e = Etcd({'ttl': 30, 'host': 'remotehost:2379', 'scope': 'test'}) + self.e = Etcd('foo', {'ttl': 30, 'host': 'remotehost:2379', 'scope': 'test'}) self.ha = Ha(self.p, self.e) self.ha.load_cluster_from_dcs() self.ha.cluster = Cluster(False, None, None, []) From 43b12af3a7ea5307b06273a26bbd0fd8cc712e59 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Tue, 7 Jul 2015 12:45:14 +0200 Subject: [PATCH 7/9] Implement possibility to work against ZooKeeper This implementation is using the same interface (AbstractDCS) as Etcd class. It means that there should be no problem to implement another plugin to work agains Consul for example. --- governor.py | 5 +- helpers/zookeeper.py | 157 ++++++++++++++++++++++++++++++++++++++++ postgres0.yml | 10 ++- postgres1.yml | 14 +++- requirements-py2.txt | 1 + requirements-py3.txt | 1 + tests/test_governor.py | 10 ++- tests/test_zookeeper.py | 138 +++++++++++++++++++++++++++++++++++ 8 files changed, 328 insertions(+), 8 deletions(-) create mode 100644 helpers/zookeeper.py create mode 100644 tests/test_zookeeper.py diff --git a/governor.py b/governor.py index ec4b39ad..b3c8a48c 100755 --- a/governor.py +++ b/governor.py @@ -7,9 +7,10 @@ import yaml from helpers.api import RestApiServer from helpers.etcd import Etcd -from helpers.postgresql import Postgresql from helpers.ha import Ha +from helpers.postgresql import Postgresql from helpers.utils import setup_signal_handlers, sleep +from helpers.zookeeper import ZooKeeper class Governor: @@ -27,6 +28,8 @@ class Governor: def get_dcs(name, config): if 'etcd' in config: return Etcd(name, config['etcd']) + if 'zookeeper' in config: + return ZooKeeper(name, config['zookeeper']) raise Exception('Can not find sutable configuration of distributed configuration store') def touch_member(self, ttl=None): diff --git a/helpers/zookeeper.py b/helpers/zookeeper.py new file mode 100644 index 00000000..4e7a22c6 --- /dev/null +++ b/helpers/zookeeper.py @@ -0,0 +1,157 @@ +import logging + +from helpers.dcs import AbstractDCS, Cluster, DCSError, Member, parse_connection_string +from kazoo.client import KazooClient, KazooState +from kazoo.exceptions import NoNodeError, NodeExistsError + +logger = logging.getLogger(__name__) + + +class ZooKeeperError(DCSError): + pass + + +class ZooKeeper(AbstractDCS): + + def __init__(self, name, config): + super(ZooKeeper, self).__init__(name, config) + self.fetch_cluster = True + self.members = [] + self.leader = None + self.last_leader_operation = 0 + self.client = KazooClient(hosts=config['hosts'], + timeout=(config.get('session_timeout', None) or 30), + command_retry={ + 'deadline': (config.get('reconnect_timeout', None) or 10), + 'max_delay': 1, + 'max_tries': -1}, + connection_retry={'max_delay': 1, 'max_tries': -1}) + self.client.add_listener(self.session_listener) + self.cluster_event = self.client.handler.event_object() + self.client.start(None) + + def session_listener(self, state): + if state in [KazooState.SUSPENDED, KazooState.LOST]: + self.cluster_watcher(None) + + def cluster_watcher(self, event): + self.fetch_cluster = True + self.cluster_event.set() + + def get_node(self, name, watch=None): + try: + return self.client.get(self.client_path(name), watch) + except NoNodeError: + pass + except: + logger.exception('get_node') + 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 load_members(self): + members = [] + for member in self.client.get_children(self.client_path('/members'), self.cluster_watcher): + data = self.get_node('/members/' + 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('/leader', self.cluster_watcher) + self.members = self.load_members() + if leader: + if leader[0] == self._name: + client_id = self.client.client_id + if 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') + self.client.delete(self.client_path('/leader')) + leader = None + + if leader: + for member in self.members: + if member.name == leader[0]: + leader = member + self.fetch_cluster = False + break + if not isinstance(leader, Member): + leader = Member(-1, leader, None, None, None, None) + self.leader = leader + if self.fetch_cluster: + last_leader_operation = self.get_node('/optime/leader') + if last_leader_operation: + self.last_leader_operation = int(last_leader_operation[0]) + + def get_cluster(self): + if self.fetch_cluster: + try: + self.client.retry(self._inner_load_cluster) + except: + 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) + + def _create(self, path, value, **kwargs): + try: + self.client.retry(self.client.create, self.client_path(path), value, **kwargs) + return True + except: + return False + + def attempt_to_acquire_leader(self): + ret = self._create('/leader', self._name, makepath=True, ephemeral=True) + ret or logger.info('Could not take out TTL lock') + return ret + + def race(self, path): + return self._create(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 + path = self.client_path('/members/' + self._name) + try: + self.client.retry(self.client.create, path, connection_string, makepath=True, ephemeral=True) + return True + except NodeExistsError: + try: + self.client.retry(self.client.delete, path) + self.client.retry(self.client.create, path, connection_string, makepath=True, ephemeral=True) + return True + except: + logger.exception('touch_member') + return False + + def take_leader(self): + return self.attempt_to_acquire_leader() + + def update_leader(self, state_handler): + last_operation = state_handler.last_operation() + if last_operation != self.last_leader_operation: + self.last_leader_operation = last_operation + path = self.client_path('/optime/leader') + try: + self.client.retry(self.client.set, path, last_operation) + except NoNodeError: + try: + self.client.retry(self.client.create, path, last_operation, makepath=True) + except: + logger.exception('Failed to create %s', path) + except: + logger.exception('Failed to update %s', path) + return True + + def delete_leader(self): + if isinstance(self.leader, Member) and self.leader.name == self._name: + self.client.delete(self.client_path('/leader')) + + def sleep(self, timeout): + self.cluster_event.wait(timeout) + if self.cluster_event.isSet(): + self.fetch_cluster = True diff --git a/postgres0.yml b/postgres0.yml index 0b5910f9..7dc453b4 100644 --- a/postgres0.yml +++ b/postgres0.yml @@ -1,12 +1,18 @@ -loop_wait: 10 +ttl: &ttl 30 +loop_wait: &loop_wait 10 restapi: listen: 127.0.0.1:8008 connect_address: 127.0.0.1:8008 etcd: scope: batman - ttl: 30 + ttl: *ttl host: 127.0.0.1:4001 #discovery_srv: my-etcd.domain +#zookeeper: +# scope: batman +# session_timeout: *ttl +# reconnect_timeout: *loop_wait +# hosts: 127.0.0.1:2181 postgresql: name: postgresql0 listen: 127.0.0.1:5432 diff --git a/postgres1.yml b/postgres1.yml index 5d5687d3..1f85542d 100644 --- a/postgres1.yml +++ b/postgres1.yml @@ -1,12 +1,18 @@ -loop_wait: 10 +ttl: &ttl 30 +loop_wait: &loop_wait 10 restapi: - listen: 127.0.0.1:8009 - connect_address: 127.0.0.1:8009 + listen: 127.0.0.1:8010 + connect_address: 127.0.0.1:8010 etcd: scope: batman - ttl: 30 + ttl: *ttl host: 127.0.0.1:4001 #discovery_srv: my-etcd.domain +#zookeeper: +# scope: batman +# session_timeout: *ttl +# reconnect_timeout: *loop_wait +# hosts: 127.0.0.1:2181 postgresql: name: postgresql1 listen: 127.0.0.1:5433 diff --git a/requirements-py2.txt b/requirements-py2.txt index 436f88ab..f2703a77 100644 --- a/requirements-py2.txt +++ b/requirements-py2.txt @@ -2,3 +2,4 @@ dnspython psycopg2 PyYAML requests +kazoo>=2.2.1 diff --git a/requirements-py3.txt b/requirements-py3.txt index 13f53d0e..93428c08 100644 --- a/requirements-py3.txt +++ b/requirements-py3.txt @@ -2,3 +2,4 @@ dnspython3 psycopg2 PyYAML requests +kazoo>=2.2.1 diff --git a/tests/test_governor.py b/tests/test_governor.py index daced98b..ed9d6a51 100644 --- a/tests/test_governor.py +++ b/tests/test_governor.py @@ -1,4 +1,5 @@ import datetime +import helpers.zookeeper import psycopg2 import requests import subprocess @@ -9,9 +10,11 @@ import yaml from governor import Governor, main from helpers.dcs import Cluster, Member +from helpers.zookeeper import ZooKeeper +from test_etcd import requests_get, requests_put, requests_delete from test_ha import true, false from test_postgresql import Postgresql, subprocess_call, psycopg2_connect -from test_etcd import requests_get, requests_put, requests_delete +from test_zookeeper import MockKazooClient if sys.hexversion >= 0x03000000: import http.server as BaseHTTPServer @@ -57,6 +60,11 @@ class TestGovernor(unittest.TestCase): Postgresql.write_pg_hba = self.write_pg_hba Postgresql.write_recovery_conf = self.write_recovery_conf + def test_get_dcs(self): + helpers.zookeeper.KazooClient = MockKazooClient + self.assertIsInstance(self.g.get_dcs('', {'zookeeper': {'scope': '', 'hosts': ''}}), ZooKeeper) + self.assertRaises(Exception, self.g.get_dcs, '', {}) + def test_governor_main(self): main() sys.argv = ['governor.py', 'postgres0.yml'] diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py new file mode 100644 index 00000000..3b0bb4a2 --- /dev/null +++ b/tests/test_zookeeper.py @@ -0,0 +1,138 @@ +import helpers.zookeeper +import unittest + +from helpers.zookeeper import ZooKeeper, ZooKeeperError +from kazoo.client import KazooState +from kazoo.exceptions import NoNodeError, NodeExistsError +from kazoo.protocol.states import ZnodeStat +from test_etcd import MockPostgresql + + +class MockEvent: + + def clear(self): + pass + + def set(self): + pass + + def wait(self, timeout): + pass + + def isSet(self): + return True + + +class MockEventHandler: + + def event_object(self): + return MockEvent() + + +class MockKazooClient: + + def __init__(self, **kwargs): + self.handler = MockEventHandler() + self.leader = False + self.exists = True + + def start(self, timeout): + pass + + @property + def client_id(self): + return (-1, '') + + def add_listener(self, cb): + pass + + def retry(self, func, *args, **kwargs): + func(*args, **kwargs) + + def get(self, path, watch=None): + if path == '/service/test/no_node': + raise NoNodeError + elif path == '/service/test/other_exception': + raise Exception() + elif '/members/' in path: + return ( + 'postgres://repuser:rep-pass@localhost:5434/postgres?application_name=http://127.0.0.1:8009/governor', + ZnodeStat(0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0) + ) + elif path.endswith('/optime/leader'): + return '1' + elif path.endswith('/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, 0, 0, 0, 0)) + + def get_children(self, path, watch=None, include_data=False): + return ['foo', 'bar', 'buzz'] + + def create(self, path, value="", acl=None, ephemeral=False, sequence=False, makepath=False): + if path.endswith('/initialize') or path == '/service/test/optime/leader': + raise Exception + elif value == 'retry' or (value == 'exists' and self.exists): + raise NodeExistsError + + def set(self, path, value, version=-1): + if path == '/service/bla/optime/leader': + raise Exception + raise NoNodeError + + def delete(self, path, version=-1, recursive=False): + self.exists = False + if path == '/service/test/leader': + if self.leader: + return + self.leader = True + raise Exception + + +class TestZooKeeper(unittest.TestCase): + + def __init__(self, method_name='runTest'): + self.setUp = self.set_up + super(TestZooKeeper, self).__init__(method_name) + + def set_up(self): + helpers.zookeeper.KazooClient = MockKazooClient + self.zk = ZooKeeper('foo', {'hosts': 'localhost:2181', 'scope': 'test'}) + + def test_session_listener(self): + self.zk.session_listener(KazooState.SUSPENDED) + + def test_get_node(self): + self.assertIsNone(self.zk.get_node('/no_node')) + self.assertIsNone(self.zk.get_node('/other_exception')) + + def test__inner_load_cluster(self): + self.zk._base_path = self.zk._base_path.replace('test', 'bla') + self.zk._inner_load_cluster() + + def test_get_cluster(self): + self.assertRaises(ZooKeeperError, self.zk.get_cluster) + self.zk.get_cluster() + self.zk.touch_member('foo') + self.zk.delete_leader() + + def test_race(self): + self.assertFalse(self.zk.race('/initialize')) + + def test_touch_member(self): + self.zk.touch_member('new') + self.zk.touch_member('exists') + self.zk.touch_member('retry') + + def test_take_leader(self): + self.zk.take_leader() + + def test_update_leader(self): + self.zk.last_leader_operation = -1 + self.assertTrue(self.zk.update_leader(MockPostgresql())) + self.zk._base_path = self.zk._base_path.replace('test', 'bla') + self.zk.last_leader_operation = -1 + self.assertTrue(self.zk.update_leader(MockPostgresql())) + + def test_sleep(self): + self.zk.sleep(0) From a1f11fe2fe3fb5f062e5bac25e458bcee18be05e Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 8 Jul 2015 08:53:08 +0200 Subject: [PATCH 8/9] Remove unneeded class CurrentLeaderError --- helpers/etcd.py | 4 ---- 1 file changed, 4 deletions(-) diff --git a/helpers/etcd.py b/helpers/etcd.py index 1e971f30..83f9e0a9 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -16,10 +16,6 @@ class EtcdError(DCSError): pass -class CurrentLeaderError(EtcdError): - pass - - class EtcdConnectionFailed(EtcdError): pass From 00aa99b38fdf0b2cc8daf09d5f58cbfa2f6c8d45 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Wed, 8 Jul 2015 08:53:51 +0200 Subject: [PATCH 9/9] Revert restapi port to 8009 --- postgres1.yml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/postgres1.yml b/postgres1.yml index 1f85542d..a2caa196 100644 --- a/postgres1.yml +++ b/postgres1.yml @@ -1,8 +1,8 @@ ttl: &ttl 30 loop_wait: &loop_wait 10 restapi: - listen: 127.0.0.1:8010 - connect_address: 127.0.0.1:8010 + listen: 127.0.0.1:8009 + connect_address: 127.0.0.1:8009 etcd: scope: batman ttl: *ttl