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
This commit is contained in:
+6
-21
@@ -1,11 +1,10 @@
|
||||
#!/usr/bin/env python
|
||||
|
||||
import sys
|
||||
import yaml
|
||||
import time
|
||||
import urllib2
|
||||
import atexit
|
||||
import logging
|
||||
import sys
|
||||
import time
|
||||
import yaml
|
||||
|
||||
from helpers.etcd import Etcd
|
||||
from helpers.postgresql import Postgresql
|
||||
@@ -42,14 +41,9 @@ def stop_postgresql():
|
||||
atexit.register(stop_postgresql)
|
||||
|
||||
# wait for etcd to be available
|
||||
etcd_ready = False
|
||||
while not etcd_ready:
|
||||
try:
|
||||
etcd.touch_member(postgresql.name, postgresql.connection_string)
|
||||
etcd_ready = True
|
||||
except urllib2.URLError:
|
||||
logging.info("waiting on etcd")
|
||||
time.sleep(5)
|
||||
while not etcd.touch_member(postgresql.name, postgresql.connection_string):
|
||||
logging.info("waiting on etcd")
|
||||
time.sleep(5)
|
||||
|
||||
# is data directory empty?
|
||||
if postgresql.data_directory_empty():
|
||||
@@ -73,16 +67,7 @@ if postgresql.data_directory_empty():
|
||||
synced_from_leader = True
|
||||
else:
|
||||
time.sleep(5)
|
||||
else:
|
||||
postgresql.write_recovery_conf(None)
|
||||
postgresql.start()
|
||||
|
||||
while True:
|
||||
logging.info(ha.run_cycle())
|
||||
|
||||
# create replication slots
|
||||
if postgresql.is_leader():
|
||||
members = [m['hostname'] for m in etcd.members() if m['hostname'] != postgresql.name]
|
||||
postgresql.create_replication_slots(members)
|
||||
|
||||
time.sleep(config["loop_wait"])
|
||||
|
||||
+83
-80
@@ -1,117 +1,120 @@
|
||||
import urllib2
|
||||
import json
|
||||
import time
|
||||
import logging
|
||||
import requests
|
||||
import time
|
||||
|
||||
from helpers.errors import CurrentLeaderError
|
||||
from urllib import urlencode
|
||||
from collections import namedtuple
|
||||
from helpers.errors import CurrentLeaderError, EtcdError
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class Member(namedtuple('Member', 'hostname,address')):
|
||||
|
||||
pass
|
||||
|
||||
|
||||
class Cluster(namedtuple('Cluster', 'leader,members')):
|
||||
|
||||
pass
|
||||
|
||||
|
||||
class Etcd:
|
||||
|
||||
def __init__(self, config):
|
||||
self.ttl = config['ttl']
|
||||
self.base_client_url = 'http://{host}/v2/keys/service/{scope}'.format(**config)
|
||||
self.postgres_cluster = None
|
||||
|
||||
def get_client_path(self, path, max_attempts=1):
|
||||
attempts = 0
|
||||
response = None
|
||||
|
||||
while True:
|
||||
ex = None
|
||||
try:
|
||||
response = urllib2.urlopen(self.client_url(path)).read()
|
||||
break
|
||||
except (urllib2.HTTPError, urllib2.URLError) as e:
|
||||
attempts += 1
|
||||
if attempts < max_attempts:
|
||||
logger.info('Failed to return %s, trying again. (%s of %s)', path, attempts, max_attempts)
|
||||
time.sleep(3)
|
||||
else:
|
||||
raise e
|
||||
try:
|
||||
return json.loads(response)
|
||||
except ValueError:
|
||||
return response
|
||||
response = requests.get(self.client_url(path))
|
||||
if response.status_code == 200:
|
||||
break
|
||||
except Exception, e:
|
||||
logger.exception('get_client_path')
|
||||
ex = e
|
||||
|
||||
def put_client_path(self, path, data):
|
||||
opener = urllib2.build_opener(urllib2.HTTPHandler)
|
||||
request = urllib2.Request(self.client_url(path), data=urlencode(data).replace("false", "False"))
|
||||
request.get_method = lambda: 'PUT'
|
||||
opener.open(request)
|
||||
attempts += 1
|
||||
if attempts < max_attempts:
|
||||
logger.info('Failed to return %s, trying again. (%s of %s)', path, attempts, max_attempts)
|
||||
time.sleep(3)
|
||||
elif ex:
|
||||
raise ex
|
||||
|
||||
return response.json(), response.status_code
|
||||
|
||||
def put_client_path(self, path, **data):
|
||||
try:
|
||||
response = requests.put(self.client_url(path), data=data)
|
||||
return response.status_code in [200, 201]
|
||||
except:
|
||||
logger.exception('PUT %s data=%s', path, data)
|
||||
return False
|
||||
|
||||
def client_url(self, path):
|
||||
return self.base_client_url + path
|
||||
|
||||
@staticmethod
|
||||
def find_node(node, key):
|
||||
if not node['dir']:
|
||||
return None
|
||||
key = node['key'] + key
|
||||
for n in node['nodes']:
|
||||
if n['key'] == key:
|
||||
return n
|
||||
return None
|
||||
|
||||
def get_cluster(self):
|
||||
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 []
|
||||
|
||||
leader_node = self.find_node(response['node'], '/leader')
|
||||
if leader_node:
|
||||
for m in members:
|
||||
if m.hostname == leader_node['value']:
|
||||
leader = m
|
||||
break
|
||||
if not leader:
|
||||
leader = Member(leader['value'], None)
|
||||
return Cluster(leader, members)
|
||||
elif status_code == 404:
|
||||
return Cluster(None, [])
|
||||
except:
|
||||
logger.exception('get_cluster')
|
||||
|
||||
raise EtcdError('Etcd is not responding properly')
|
||||
|
||||
def current_leader(self):
|
||||
try:
|
||||
hostname = self.get_client_path('/leader')['node']['value']
|
||||
address = self.get_client_path('/members/' + hostname)['node']['value']
|
||||
|
||||
return {'hostname': hostname, 'address': address}
|
||||
except urllib2.HTTPError as e:
|
||||
if e.code == 404:
|
||||
return None
|
||||
raise CurrentLeaderError("Etcd is not responding properly")
|
||||
|
||||
def members(self):
|
||||
try:
|
||||
members = []
|
||||
|
||||
r = self.get_client_path("/members?recursive=true")
|
||||
for node in r["node"]["nodes"]:
|
||||
members.append({"hostname": node["key"].split('/')[-1], "address": node["value"]})
|
||||
|
||||
return members
|
||||
except urllib2.HTTPError as e:
|
||||
if e.code == 404:
|
||||
cluster = self.get_cluster()
|
||||
if not cluster['leader'] or not cluster['leader'].address:
|
||||
return None
|
||||
return cluster['leader']
|
||||
except:
|
||||
raise CurrentLeaderError("Etcd is not responding properly")
|
||||
|
||||
def touch_member(self, member, connection_string):
|
||||
self.put_client_path('/members/' + member, {"value": connection_string})
|
||||
return self.put_client_path('/members/' + member, value=connection_string)
|
||||
|
||||
def take_leader(self, value):
|
||||
return self.put_client_path("/leader", {"value": value, "ttl": self.ttl}) is None
|
||||
return self.put_client_path('/leader', value=value, ttl=self.ttl)
|
||||
|
||||
def attempt_to_acquire_leader(self, value):
|
||||
try:
|
||||
return self.put_client_path("/leader", {"value": value, "ttl": self.ttl, "prevExist": False}) is None
|
||||
except urllib2.HTTPError as e:
|
||||
if e.code == 412:
|
||||
logger.info('Could not take out TTL lock: %s', e)
|
||||
return False
|
||||
ret = self.put_client_path('/leader', value=value, ttl=self.ttl, prevExist=False)
|
||||
ret or logger.info('Could not take out TTL lock')
|
||||
return ret
|
||||
|
||||
def update_leader(self, value):
|
||||
try:
|
||||
self.put_client_path("/leader", {"value": value, "ttl": self.ttl, "prevValue": value})
|
||||
return True
|
||||
except urllib2.HTTPError:
|
||||
logger.error("Error updating TTL on ETCD for primary.")
|
||||
return False
|
||||
|
||||
def leader_unlocked(self):
|
||||
try:
|
||||
self.get_client_path("/leader")
|
||||
return False
|
||||
except urllib2.HTTPError as e:
|
||||
if e.code == 404:
|
||||
return True
|
||||
return False
|
||||
except ValueError as e:
|
||||
return False
|
||||
|
||||
def am_i_leader(self, value):
|
||||
# try:
|
||||
reponse = self.get_client_path("/leader")
|
||||
logger.info('Lock owner: %s; I am %s', reponse["node"]["value"], value)
|
||||
return reponse["node"]["value"] == value
|
||||
# except Exception as e:
|
||||
# return False
|
||||
return self.put_client_path('/leader', value=value, ttl=self.ttl, prevValue=value)
|
||||
|
||||
def race(self, path, value):
|
||||
try:
|
||||
return self.put_client_path(path, {"prevExist": False, "value": value}) is None
|
||||
except urllib2.HTTPError:
|
||||
return False
|
||||
return self.put_client_path(path, value=value, prevExist=False)
|
||||
|
||||
+59
-39
@@ -2,7 +2,7 @@ import inspect
|
||||
import logging
|
||||
import time
|
||||
|
||||
from helpers.errors import CurrentLeaderError, HealthiestMemberError
|
||||
from helpers.errors import EtcdError, HealthiestMemberError
|
||||
from psycopg2 import OperationalError
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -18,6 +18,10 @@ class Ha:
|
||||
def __init__(self, state_handler, etcd):
|
||||
self.state_handler = state_handler
|
||||
self.etcd = etcd
|
||||
self.cluster = None
|
||||
|
||||
def load_cluster_from_etcd(self):
|
||||
self.cluster = self.etcd.get_cluster()
|
||||
|
||||
def acquire_lock(self):
|
||||
return self.etcd.attempt_to_acquire_leader(self.state_handler.name)
|
||||
@@ -26,60 +30,76 @@ class Ha:
|
||||
return self.etcd.update_leader(self.state_handler.name)
|
||||
|
||||
def is_unlocked(self):
|
||||
return self.etcd.leader_unlocked()
|
||||
return not (self.cluster.leader and self.cluster.leader.hostname)
|
||||
|
||||
def has_lock(self):
|
||||
return self.etcd.am_i_leader(self.state_handler.name)
|
||||
logger.info('Lock owner: %s; I am %s', self.cluster.leader.hostname, self.state_handler.name)
|
||||
return self.cluster.leader.hostname == self.state_handler.name
|
||||
|
||||
def fetch_current_leader(self):
|
||||
return self.etcd.current_leader()
|
||||
def demote(self):
|
||||
return self.state_handler.demote(self.cluster.leader)
|
||||
|
||||
def follow_the_leader(self):
|
||||
return self.state_handler.follow_the_leader(self.cluster.leader)
|
||||
|
||||
def run_cycle(self):
|
||||
try:
|
||||
if self.state_handler.is_healthy():
|
||||
if self.is_unlocked():
|
||||
if self.state_handler.is_healthiest_node(self.etcd.members()):
|
||||
if self.acquire_lock():
|
||||
if not self.state_handler.is_leader():
|
||||
self.state_handler.promote()
|
||||
return "promoted self to leader by acquiring session lock"
|
||||
return "acquired session lock as a leader"
|
||||
else:
|
||||
if self.state_handler.is_leader():
|
||||
self.state_handler.demote(self.fetch_current_leader())
|
||||
return "demoted self due after trying and failing to obtain lock"
|
||||
else:
|
||||
self.state_handler.follow_the_leader(self.fetch_current_leader())
|
||||
return "following new leader after trying and failing to obtain lock"
|
||||
self.load_cluster_from_etcd()
|
||||
if self.is_unlocked():
|
||||
if not self.state_handler.is_healthy():
|
||||
return 'no action. not healthy enough to do anything.'
|
||||
elif self.state_handler.is_healthiest_node(self.cluster.members):
|
||||
if self.acquire_lock():
|
||||
if not self.state_handler.is_leader():
|
||||
self.state_handler.promote()
|
||||
return "promoted self to leader by acquiring session lock"
|
||||
return "acquired session lock as a leader"
|
||||
else:
|
||||
self.load_cluster_from_etcd()
|
||||
if self.state_handler.is_leader():
|
||||
self.state_handler.demote(self.fetch_current_leader())
|
||||
return "demoting self because i am not the healthiest node"
|
||||
self.demote()
|
||||
return "demoted self due after trying and failing to obtain lock"
|
||||
else:
|
||||
self.state_handler.follow_the_leader(self.fetch_current_leader())
|
||||
return "following a different leader because i am not the healthiest node"
|
||||
self.follow_the_leader()
|
||||
return "following new leader after trying and failing to obtain lock"
|
||||
else:
|
||||
if self.has_lock() and self.update_lock():
|
||||
self.load_cluster_from_etcd()
|
||||
if self.state_handler.is_leader():
|
||||
self.demote()
|
||||
return "demoting self because i am not the healthiest node"
|
||||
else:
|
||||
self.follow_the_leader()
|
||||
return "following a different leader because i am not the healthiest node"
|
||||
else:
|
||||
if self.has_lock() and not self.state_handler.is_healthy():
|
||||
self.state_handler.write_recovery_conf(None)
|
||||
self.state_handler.start()
|
||||
self.load_cluster_from_etcd()
|
||||
|
||||
if self.has_lock() and self.update_lock():
|
||||
try:
|
||||
if not self.state_handler.is_leader():
|
||||
self.state_handler.promote()
|
||||
return "promoted self to leader because i had the session lock"
|
||||
else:
|
||||
return "no action. i am the leader with the lock"
|
||||
finally:
|
||||
# create replication slots
|
||||
self.state_handler.create_replication_slots([m.hostname for m in self.cluster.members])
|
||||
else:
|
||||
logger.info("does not have lock")
|
||||
if not self.state_handler.is_healthy():
|
||||
self.state_handler.write_recovery_conf(self.cluster.leader)
|
||||
self.state_handler.start()
|
||||
return 'starting as a secondary'
|
||||
elif self.state_handler.is_leader():
|
||||
self.demote()
|
||||
return "demoting self because i do not have the lock and i was a leader"
|
||||
else:
|
||||
logger.info("does not have lock")
|
||||
if self.state_handler.is_leader():
|
||||
self.state_handler.demote(self.fetch_current_leader())
|
||||
return "demoting self because i do not have the lock and i was a leader"
|
||||
else:
|
||||
self.state_handler.follow_the_leader(self.fetch_current_leader())
|
||||
return "no action. i am a secondary and i am following a leader"
|
||||
else:
|
||||
if not self.state_handler.is_running(): # XXX is_running == is_healthy
|
||||
self.state_handler.start()
|
||||
return "postgresql was stopped. starting again."
|
||||
return "no action. not healthy enough to do anything."
|
||||
except CurrentLeaderError:
|
||||
logger.error("failed to fetch current leader from etcd")
|
||||
self.follow_the_leader()
|
||||
return "no action. i am a secondary and i am following a leader"
|
||||
except EtcdError:
|
||||
logger.error("Error communicating with Etcd")
|
||||
except OperationalError:
|
||||
logger.error("Error communicating with Postgresql. Will try again.")
|
||||
except HealthiestMemberError:
|
||||
|
||||
+10
-10
@@ -142,10 +142,10 @@ class Postgresql:
|
||||
|
||||
def is_healthiest_node(self, members):
|
||||
for member in members:
|
||||
if member['hostname'] == self.name:
|
||||
if member.hostname == self.name:
|
||||
continue
|
||||
try:
|
||||
member_conn = psycopg2.connect(member['address'])
|
||||
member_conn = psycopg2.connect(member.address)
|
||||
member_conn.autocommit = True
|
||||
member_cursor = member_conn.cursor()
|
||||
member_cursor.execute(
|
||||
@@ -178,11 +178,11 @@ class Postgresql:
|
||||
r = parseurl(leader_url)
|
||||
return 'user={username} password={password} host={hostname} port={port} sslmode=prefer sslcompression=1'.format(**r)
|
||||
|
||||
def check_recovery_conf(self, leader_hash):
|
||||
def check_recovery_conf(self, leader):
|
||||
if not os.path.isfile(self.recovery_conf):
|
||||
return False
|
||||
|
||||
pattern = leader_hash and 'address' in leader_hash and self.primary_conninfo(leader_hash['address'])
|
||||
pattern = leader and leader.address and self.primary_conninfo(leader.address)
|
||||
|
||||
with open(self.recovery_conf, 'r') as f:
|
||||
for line in f:
|
||||
@@ -193,23 +193,23 @@ class Postgresql:
|
||||
|
||||
return not pattern
|
||||
|
||||
def write_recovery_conf(self, leader_hash):
|
||||
def write_recovery_conf(self, leader):
|
||||
with open(self.recovery_conf, 'w') as f:
|
||||
f.write("""standby_mode = 'on'
|
||||
recovery_target_timeline = 'latest'
|
||||
""")
|
||||
if leader_hash and 'address' in leader_hash:
|
||||
if leader and leader.address:
|
||||
f.write("""
|
||||
primary_slot_name = '{}'
|
||||
primary_conninfo = '{}'
|
||||
""".format(self.name, self.primary_conninfo(leader_hash['address'])))
|
||||
""".format(self.name, self.primary_conninfo(leader.address)))
|
||||
for name, value in self.config.get('recovery_conf', {}).iteritems():
|
||||
f.write("{} = '{}'\n".format(name, value))
|
||||
|
||||
def follow_the_leader(self, leader_hash):
|
||||
if self.check_recovery_conf(leader_hash):
|
||||
def follow_the_leader(self, leader):
|
||||
if self.check_recovery_conf(leader):
|
||||
return
|
||||
self.write_recovery_conf(leader_hash)
|
||||
self.write_recovery_conf(leader)
|
||||
self.restart()
|
||||
|
||||
def promote(self):
|
||||
|
||||
Reference in New Issue
Block a user