mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-26 07:30:14 +00:00
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.
This commit is contained in:
+12
-3
@@ -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:
|
||||
|
||||
+4
-3
@@ -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:
|
||||
|
||||
+1
-1
@@ -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
|
||||
|
||||
|
||||
@@ -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))
|
||||
|
||||
|
||||
+12
-15
@@ -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):
|
||||
|
||||
+8
-4
@@ -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())
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:[email protected]:5435/postgres', None, None, 28)
|
||||
self.leadermem = Member(0, 'leader', 'postgres://replicator:[email protected]:5435/postgres', None, None, 28)
|
||||
self.leader = Leader(-1, None, 28, self.leadermem)
|
||||
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)
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user