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.
This commit is contained in:
Alexander Kukushkin
2015-07-07 12:45:14 +02:00
parent 7445590546
commit 43b12af3a7
8 changed files with 328 additions and 8 deletions
+4 -1
View File
@@ -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):
+157
View File
@@ -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
+8 -2
View File
@@ -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
+10 -4
View File
@@ -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
+1
View File
@@ -2,3 +2,4 @@ dnspython
psycopg2
PyYAML
requests
kazoo>=2.2.1
+1
View File
@@ -2,3 +2,4 @@ dnspython3
psycopg2
PyYAML
requests
kazoo>=2.2.1
+9 -1
View File
@@ -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']
+138
View File
@@ -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)