last_leader_operation is the propery of Cluster object and the value is set in get_cluster method

This commit is contained in:
Alexander Kukushkin
2015-05-18 13:33:27 +02:00
parent d221d1de1c
commit b24fb0488c
3 changed files with 28 additions and 46 deletions
+22 -27
View File
@@ -8,14 +8,8 @@ from helpers.errors import CurrentLeaderError, EtcdError
logger = logging.getLogger(__name__)
class Member(namedtuple('Member', 'hostname,address')):
pass
class Cluster(namedtuple('Cluster', 'leader,members')):
pass
Member = namedtuple('Member', 'hostname,address,ttl')
Cluster = namedtuple('Cluster', 'leader,last_leader_operation,members')
class Etcd:
@@ -87,21 +81,32 @@ class Etcd:
try:
response, status_code = self.get_client_path('?recursive=true')
if status_code == 200:
leader = None
members = self.find_node(response['node'], '/members')
members = [Member(n['key'].split('/')[-1], n['value']) for n in members['nodes']] if members else []
# get list of members
node = self.find_node(response['node'], '/members') or {'nodes': []}
members = [Member(n['key'].split('/')[-1], n['value'], n.get('ttl', None)) for n in node['nodes']]
leader_node = self.find_node(response['node'], '/leader')
if leader_node:
# get last leader operation
last_leader_operation = 0
node = self.find_node(response['node'], '/optime')
if node:
node = self.find_node(node, '/leader')
if node:
last_leader_operation = int(node['value'])
# get leader
leader = None
node = self.find_node(response['node'], '/leader')
if node:
for m in members:
if m.hostname == leader_node['value']:
if m.hostname == node['value']:
leader = m
break
if not leader:
leader = Member(leader['value'], None)
return Cluster(leader, members)
leader = Member(leader['value'], None, None)
return Cluster(leader, last_leader_operation, members)
elif status_code == 404:
return Cluster(None, [])
return Cluster(None, None, [])
except:
logger.exception('get_cluster')
@@ -135,16 +140,6 @@ class Etcd:
def race(self, path, value):
return self.put_client_path(path, value=value, prevExist=False)
def last_leader_operation(self):
try:
response, status_code = self.get_client_path('/optime/leader')
if status_code == 404:
return None
return int(response['node']['value'])
except:
logger.exception('last_leader_operation')
raise EtcdError('Etcd is not responding properly')
def delete_member(self, member):
return self.delete_client_path('/members/' + member)
+1 -7
View File
@@ -1,5 +1,4 @@
import logging
import time
from helpers.errors import EtcdError, HealthiestMemberError
from psycopg2 import OperationalError
@@ -50,7 +49,7 @@ class Ha:
self.load_cluster_from_etcd()
if self.is_unlocked():
if self.state_handler.is_healthiest_node(self.etcd.last_leader_operation(), self.cluster.members):
if self.state_handler.is_healthiest_node(self.cluster):
if self.acquire_lock():
if not self.state_handler.is_leader():
self.state_handler.promote()
@@ -97,8 +96,3 @@ class Ha:
logger.error('Error communicating with Postgresql. Will try again')
except HealthiestMemberError:
logger.error('failed to determine healthiest member fromt etcd')
def run(self):
while True:
self.run_cycle()
time.sleep(10)
+5 -12
View File
@@ -1,7 +1,6 @@
import logging
import os
import psycopg2
import re
import sys
import time
@@ -144,11 +143,11 @@ class Postgresql:
return False
return True
def is_healthiest_node(self, last_leader_operation, members):
if (last_leader_operation or 0) - self.xlog_position() > self.config.get('maximum_lag_on_failover', 0):
def is_healthiest_node(self, cluster):
if cluster.last_leader_operation - self.xlog_position() > self.config.get('maximum_lag_on_failover', 0):
return False
for member in members:
for member in cluster.members:
if member.hostname == self.name:
continue
try:
@@ -159,19 +158,13 @@ class Postgresql:
"SELECT %s - (pg_last_xlog_replay_location() - '0/0000000'::pg_lsn)", (self.xlog_position(), ))
xlog_diff = member_cursor.fetchone()[0]
logger.info([self.name, member.hostname, xlog_diff])
if xlog_diff < 0:
member_cursor.close()
return False
member_cursor.close()
if xlog_diff < 0:
return False
except psycopg2.OperationalError:
continue
return True
def replication_slot_name(self):
member = os.environ.get("MEMBER")
(member, _) = re.subn(r'[^a-z0-9]+', r'_', member)
return member
def write_pg_hba(self):
with open(os.path.join(self.data_dir, 'pg_hba.conf'), 'a') as f:
f.write('host replication {username} {network} md5'.format(**self.replication))