From 3b1efff53e0db0e05b8bb7a2a4443994631293a2 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 27 Aug 2015 10:53:22 +0200 Subject: [PATCH] Refactor `Cluster` object `Cluster.leader` is not reference to `Member` anymore, but to `Leader` `Leader` class contains field `index` (update index). This field is very useful for watching for events which changing leader key. Also `Leader` contains `member` field, which should reference real member. --- helpers/dcs.py | 15 ++++++++++++--- helpers/etcd.py | 7 ++++--- helpers/ha.py | 2 +- helpers/postgresql.py | 8 ++++---- helpers/zookeeper.py | 27 ++++++++++++--------------- tests/test_etcd.py | 12 ++++++++---- tests/test_patroni.py | 12 ++++++++---- tests/test_postgresql.py | 11 ++++++----- tests/test_zookeeper.py | 12 +++++++++--- 9 files changed, 64 insertions(+), 42 deletions(-) diff --git a/helpers/dcs.py b/helpers/dcs.py index c7140c22..f5984fd9 100644 --- a/helpers/dcs.py +++ b/helpers/dcs.py @@ -39,7 +39,7 @@ class DCSError(Exception): class Member(namedtuple('Member', 'index,name,conn_url,api_url,expiration,ttl')): """Immutable object (namedtuple) which represents single member of PostgreSQL cluster. Consists of the following fields: - :param index: modification index of a given member key in DCS + :param index: modification index of a given member key in a Configuration Store :param name: name of PostgreSQL cluster member :param conn_url: connection string containing host, user and password which could be used to access this member. :param api_url: REST API url of patroni instance @@ -50,17 +50,26 @@ class Member(namedtuple('Member', 'index,name,conn_url,api_url,expiration,ttl')) return calculate_ttl(self.expiration) or -1 +class Leader(namedtuple('Leader', 'index,expiration,ttl,member')): + """Immutable object (namedtuple) which represents leader key. + Consists of the following fields: + :param index: modification index of a leader key in a Configuration Store + :param expiration: expiration time of the leader key + :param ttl: ttl of the leader key + :param member: reference to a `Member` object which represents current leader (see `Cluster.members`)""" + + class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members')): """Immutable object (namedtuple) which represents PostgreSQL cluster. Consists of the following fields: :param initialize: boolean, shows whether this cluster has initialization key stored in DC or not. - :param leader: `Member` object which represents current leader of the cluster + :param leader: `Leader` object which represents current leader of the cluster :param last_leader_operation: int or long object containing position of last known leader operation. This value is stored in `/optime/leader` key :param members: list of Member object, all PostgreSQL cluster members including leader""" def is_unlocked(self): - return not (self.leader and self.leader.name) + return not (self.leader and self.leader.member.name) class AbstractDCS: diff --git a/helpers/etcd.py b/helpers/etcd.py index add736c5..f767175b 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -8,7 +8,7 @@ import socket from dns.exception import DNSException from dns import resolver -from helpers.dcs import AbstractDCS, Cluster, DCSError, Member, parse_connection_string +from helpers.dcs import AbstractDCS, Cluster, DCSError, Leader, Member, parse_connection_string from helpers.utils import sleep from requests.exceptions import RequestException @@ -170,8 +170,9 @@ class Etcd(AbstractDCS): # get leader leader = nodes.get('leader', None) if leader: - leader = Member(-1, leader.value, None, None, None, None) - leader = ([m for m in members if m.name == leader.name] or [leader])[0] + member = Member(-1, leader.value, None, None, None, None) + member = ([m for m in members if m.name == leader.value] or [member])[0] + leader = Leader(leader.modifiedIndex, leader.expiration, leader.ttl, member) return Cluster(initialize, leader, last_leader_operation, members) except etcd.EtcdKeyNotFound: diff --git a/helpers/ha.py b/helpers/ha.py index 8e283c55..ff67f0e0 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -31,7 +31,7 @@ class Ha: return self.dcs.update_leader(self.state_handler) def has_lock(self): - lock_owner = self.cluster.leader and self.cluster.leader.name + lock_owner = self.cluster.leader and self.cluster.leader.member.name logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name) return lock_owner == self.state_handler.name diff --git a/helpers/postgresql.py b/helpers/postgresql.py index ce8ca18f..ca4ea9c0 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -128,7 +128,7 @@ class Postgresql: os.path.exists(self.trigger_file) and os.unlink(self.trigger_file) def sync_from_leader(self, leader): - r = parseurl(leader.conn_url) + r = parseurl(leader.member.conn_url) pgpass = 'pgpass' with open(pgpass, 'w') as f: @@ -288,7 +288,7 @@ class Postgresql: if not os.path.isfile(self.recovery_conf): return False - pattern = leader and leader.conn_url and self.primary_conninfo(leader.conn_url) + pattern = leader and leader.member.conn_url and self.primary_conninfo(leader.member.conn_url) with open(self.recovery_conf, 'r') as f: for line in f: @@ -304,11 +304,11 @@ class Postgresql: f.write("""standby_mode = 'on' recovery_target_timeline = 'latest' """) - if leader and leader.conn_url: + if leader and leader.member.conn_url: f.write(""" primary_slot_name = '{}' primary_conninfo = '{}' -""".format(self.name, self.primary_conninfo(leader.conn_url))) +""".format(self.name, self.primary_conninfo(leader.member.conn_url))) for name, value in self.config.get('recovery_conf', {}).items(): f.write("{} = '{}'\n".format(name, value)) diff --git a/helpers/zookeeper.py b/helpers/zookeeper.py index cb2918cd..56e755a0 100644 --- a/helpers/zookeeper.py +++ b/helpers/zookeeper.py @@ -3,7 +3,7 @@ import random import requests import time -from helpers.dcs import AbstractDCS, Cluster, DCSError, Member, parse_connection_string +from helpers.dcs import AbstractDCS, Cluster, DCSError, Leader, Member, parse_connection_string from helpers.utils import sleep from kazoo.client import KazooClient, KazooState from kazoo.exceptions import NoNodeError, NodeExistsError @@ -134,21 +134,18 @@ class ZooKeeper(AbstractDCS): 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 + client_id = self.client.client_id + 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') + 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) + member = Member(-1, leader[0], None, None, None, None) + member = ([m for m in self.members if m.name == leader[0]] or [member])[0] + leader = Leader(leader[1].mzxid, None, None, member) + self.fetch_cluster = member.index == -1 + self.leader = leader if self.fetch_cluster: last_leader_operation = self.get_node('/optime/leader') @@ -220,7 +217,7 @@ class ZooKeeper(AbstractDCS): return True def delete_leader(self): - if isinstance(self.leader, Member) and self.leader.name == self._name: + if isinstance(self.leader, Leader) and self.leader.member.name == self._name: self.client.delete(self.client_path('/leader')) def sleep(self, timeout): diff --git a/tests/test_etcd.py b/tests/test_etcd.py index 38692c52..52de036e 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -8,7 +8,7 @@ import time import unittest from dns.exception import DNSException -from helpers.dcs import Cluster, DCSError, Member +from helpers.dcs import Cluster, DCSError, Leader, Member from helpers.etcd import Client, Etcd from mock import Mock, patch @@ -107,8 +107,12 @@ def time_sleep(_): pass +class SleepException(Exception): + pass + + def time_sleep_exception(_): - raise Exception() + raise SleepException() class MockSRV: @@ -204,7 +208,7 @@ class TestEtcd(unittest.TestCase): time.sleep = time_sleep_exception with patch.object(etcd.Client, 'machines') as mock_machines: mock_machines.__get__ = Mock(side_effect=etcd.EtcdException) - self.assertRaises(Exception, self.etcd.get_etcd_client, {'discovery_srv': 'test'}) + self.assertRaises(SleepException, self.etcd.get_etcd_client, {'discovery_srv': 'test'}) def test_get_cluster(self): self.assertIsInstance(self.etcd.get_cluster(), Cluster) @@ -214,7 +218,7 @@ class TestEtcd(unittest.TestCase): self.assertIsNone(cluster.leader) def test_current_leader(self): - self.assertIsInstance(self.etcd.current_leader(), Member) + self.assertIsInstance(self.etcd.current_leader(), Leader) self.etcd._base_path = '/service/noleader' self.assertIsNone(self.etcd.current_leader()) diff --git a/tests/test_patroni.py b/tests/test_patroni.py index 67ac1ffd..575c7cf3 100644 --- a/tests/test_patroni.py +++ b/tests/test_patroni.py @@ -24,8 +24,12 @@ def nop(*args, **kwargs): pass +class SleepException(Exception): + pass + + def time_sleep(*args): - raise Exception() + raise SleepException() class Mock_BaseServer__is_shut_down: @@ -90,7 +94,7 @@ class TestPatroni(unittest.TestCase): Etcd.delete_leader = nop - self.assertRaises(Exception, main) + self.assertRaises(SleepException, main) Patroni.run = run Patroni.touch_member = touch_member @@ -100,10 +104,10 @@ class TestPatroni(unittest.TestCase): self.p.touch_member = self.touch_member self.p.ha.state_handler.sync_replication_slots = time_sleep self.p.ha.dcs.client.read = etcd_read - self.assertRaises(Exception, self.p.run) + self.assertRaises(SleepException, self.p.run) self.p.ha.state_handler.is_leader = lambda: False self.p.api.start = nop - self.assertRaises(Exception, self.p.run) + self.assertRaises(SleepException, self.p.run) def touch_member(self, ttl=None): if not self.touched: diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 040d1c66..08ba0b33 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -4,7 +4,7 @@ import shutil import subprocess import unittest -from helpers.dcs import Cluster, Member +from helpers.dcs import Cluster, Leader, Member from helpers.postgresql import Postgresql @@ -123,7 +123,8 @@ 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(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', None, None, 28) + self.leadermem = Member(0, 'leader', 'postgres://replicator:rep-pass@127.0.0.1:5435/postgres', None, None, 28) + self.leader = Leader(-1, None, 28, self.leadermem) 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) @@ -156,7 +157,7 @@ class TestPostgresql(unittest.TestCase): self.p.follow_the_leader(None) self.p.demote(self.leader) self.p.follow_the_leader(self.leader) - self.p.follow_the_leader(self.other) + self.p.follow_the_leader(Leader(-1, None, 28, self.other)) def test_create_connection_users(self): cfg = self.p.config @@ -166,7 +167,7 @@ class TestPostgresql(unittest.TestCase): def test_create_replication_slots(self): self.p.start() - cluster = Cluster(True, self.leader, 0, [self.me, self.other, self.leader]) + cluster = Cluster(True, self.leader, 0, [self.me, self.other, self.leadermem]) self.p.create_replication_slots(cluster) def test_query(self): @@ -180,7 +181,7 @@ class TestPostgresql(unittest.TestCase): self.assertRaises(psycopg2.OperationalError, self.p.query, 'blabla') def test_is_healthiest_node(self): - cluster = Cluster(True, self.leader, 0, [self.me, self.other, self.leader]) + cluster = Cluster(True, self.leader, 0, [self.me, self.other, self.leadermem]) self.assertTrue(self.p.is_healthiest_node(cluster)) self.p.is_leader = false self.assertFalse(self.p.is_healthiest_node(cluster)) diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py index 53f5ab82..54adea5e 100644 --- a/tests/test_zookeeper.py +++ b/tests/test_zookeeper.py @@ -2,6 +2,7 @@ import helpers.zookeeper import requests import unittest +from helpers.dcs import Leader from helpers.zookeeper import ExhibitorEnsembleProvider, ZooKeeper, ZooKeeperError from kazoo.client import KazooState from kazoo.exceptions import NoNodeError, NodeExistsError @@ -30,6 +31,10 @@ class MockEventHandler: return MockEvent() +class SleepException(Exception): + pass + + class MockKazooClient: def __init__(self, **kwargs): @@ -94,7 +99,7 @@ class MockKazooClient: def exhibitor_sleep(_): - raise Exception + raise SleepException class TestExhibitorEnsembleProvider(unittest.TestCase): @@ -108,7 +113,7 @@ class TestExhibitorEnsembleProvider(unittest.TestCase): helpers.zookeeper.sleep = exhibitor_sleep def test_init(self): - self.assertRaises(Exception, ExhibitorEnsembleProvider, ['localhost'], 8181) + self.assertRaises(SleepException, ExhibitorEnsembleProvider, ['localhost'], 8181) class TestZooKeeper(unittest.TestCase): @@ -136,7 +141,8 @@ class TestZooKeeper(unittest.TestCase): def test_get_cluster(self): self.assertRaises(ZooKeeperError, self.zk.get_cluster) self.zk.exhibitor.poll = lambda: True - self.zk.get_cluster() + cluster = self.zk.get_cluster() + self.assertIsInstance(cluster.leader, Leader) self.zk.touch_member('foo') self.zk.delete_leader()