mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Merge branch 'master' of github.com:CyberDem0n/governor into features/refactoring
Conflicts: helpers/postgresql.py
This commit is contained in:
+22
-27
@@ -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
@@ -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)
|
||||
|
||||
+8
-15
@@ -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))
|
||||
|
||||
Reference in New Issue
Block a user