From b24fb0488c777bebb4ce97a7e97c529cd852e673 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Mon, 18 May 2015 13:33:27 +0200 Subject: [PATCH] last_leader_operation is the propery of Cluster object and the value is set in get_cluster method --- helpers/etcd.py | 49 +++++++++++++++++++------------------------ helpers/ha.py | 8 +------ helpers/postgresql.py | 17 +++++---------- 3 files changed, 28 insertions(+), 46 deletions(-) diff --git a/helpers/etcd.py b/helpers/etcd.py index cbe44507..29552a2b 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -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) diff --git a/helpers/ha.py b/helpers/ha.py index a2788e4c..89467bf4 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -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) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 1c0e29b1..7b8a15c8 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -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))