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 185777e3..7db82dd4 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 @@ -30,9 +29,9 @@ class Postgresql: def __init__(self, config): self.name = config['name'] - self.host, self.port = config['listen'].split(':') + self.listen_addresses, self.port = config['listen'].split(':') self.libpq_parameters = { - 'host': self.host, + 'host': self.listen_addresses.split(',')[0].strip(), 'port': self.port, 'fallback_application_name': 'Governor', 'connect_timeout': 5, @@ -141,7 +140,7 @@ class Postgresql: return os.system(self._pg_ctl + ' restart -m fast') == 0 def server_options(self): - options = '--listen_addresses={} --port={}'.format(self.host, self.port) + options = "--listen_addresses='{}' --port={}".format(self.listen_addresses, self.port) for setting, value in self.config['parameters'].items(): options += " --{}='{}'".format(setting, value) return options @@ -152,11 +151,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: @@ -167,19 +166,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))