mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Inherit Etcd from abstract class
This commit is contained in:
+58
-15
@@ -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
|
||||
|
||||
+33
-14
@@ -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)
|
||||
|
||||
+4
-4
@@ -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')
|
||||
|
||||
+4
-5
@@ -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('', ''))
|
||||
|
||||
@@ -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()
|
||||
|
||||
+3
-3
@@ -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')
|
||||
|
||||
@@ -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:[email protected]:5435/postgres', None, None, 28)
|
||||
self.other = Member('test1', 'postgres://replicator:[email protected]:5433/postgres', None, None, 28)
|
||||
self.me = Member('test0', 'postgres://replicator:[email protected]:5434/postgres', None, None, 28)
|
||||
self.leader = Member(0, 'leader', 'postgres://replicator:[email protected]:5435/postgres', None, None, 28)
|
||||
self.other = Member(0, 'test1', 'postgres://replicator:[email protected]:5433/postgres', None, None, 28)
|
||||
self.me = Member(0, 'test0', 'postgres://replicator:[email protected]:5434/postgres', None, None, 28)
|
||||
|
||||
def tear_down(self):
|
||||
shutil.rmtree('data')
|
||||
|
||||
Reference in New Issue
Block a user