Merge branch 'master' into feature/awstags

This commit is contained in:
Oleksii Kliukin
2015-06-05 09:50:17 +02:00
9 changed files with 105 additions and 50 deletions
+17 -16
View File
@@ -7,8 +7,10 @@ from threading import Thread
if sys.hexversion >= 0x03000000:
from http.server import BaseHTTPRequestHandler, HTTPServer
from socketserver import ThreadingMixIn
else:
from BaseHTTPServer import BaseHTTPRequestHandler, HTTPServer
from SocketServer import ThreadingMixIn
logger = logging.getLogger(__name__)
@@ -28,16 +30,14 @@ class RestApiHandler(BaseHTTPRequestHandler):
def get_postgresql_status(self):
try:
cursor = self.server._cursor()
cursor.execute("""SELECT to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'),
pg_is_in_recovery(),
CASE WHEN pg_is_in_recovery()
THEN null
ELSE pg_current_xlog_location() END,
pg_last_xlog_receive_location(),
pg_last_xlog_replay_location(),
pg_is_in_recovery() AND pg_is_xlog_replay_paused()""")
row = cursor.fetchone()
row = self.server.query("""SELECT to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'),
pg_is_in_recovery(),
CASE WHEN pg_is_in_recovery()
THEN null
ELSE pg_current_xlog_location() END,
pg_last_xlog_receive_location(),
pg_last_xlog_replay_location(),
pg_is_in_recovery() AND pg_is_xlog_replay_paused()""")[0]
return {
'running': True,
'postmaster_start_time': row[0],
@@ -54,7 +54,7 @@ class RestApiHandler(BaseHTTPRequestHandler):
return {'running': self.server.governor.postgresql.is_running()}
class RestApiServer(HTTPServer, Thread):
class RestApiServer(ThreadingMixIn, HTTPServer, Thread):
def __init__(self, governor, config):
self.connection_string = 'http://{}/governor'.format(config.get('connect_address', None) or config['listen'])
@@ -62,10 +62,11 @@ class RestApiServer(HTTPServer, Thread):
HTTPServer.__init__(self, (host, int(port)), RestApiHandler)
Thread.__init__(self, target=self.serve_forever)
self.governor = governor
self._cursor_holder = None
self.daemon = True
def _cursor(self):
if not self._cursor_holder or self._cursor_holder.closed:
self._cursor_holder = self.governor.postgresql.connection().cursor()
return self._cursor_holder
def query(self, sql, *params):
cursor = self.governor.postgresql.connection().cursor()
cursor.execute(sql, params)
ret = [r for r in cursor]
cursor.close()
return ret
+21 -3
View File
@@ -1,14 +1,32 @@
import logging
import requests
import sys
from requests.exceptions import RequestException
from collections import namedtuple
from helpers.errors import CurrentLeaderError, EtcdError
from helpers.utils import sleep
if sys.hexversion >= 0x03000000:
from urllib.parse import urlparse, urlunparse, parse_qsl
else:
from urlparse import urlparse, urlunparse, parse_qsl
logger = logging.getLogger(__name__)
Member = namedtuple('Member', 'hostname,address,ttl')
class Member(namedtuple('Member', 'hostname,conn_url,api_url,ttl')):
@staticmethod
def fromNode(node):
scheme, netloc, path, params, query, fragment = urlparse(node['value'])
conn_url = urlunparse((scheme, netloc, path, params, '', fragment))
api_url = None
for name, value in parse_qsl(query):
if name == 'application_name' and value:
api_url = value
break
return Member(node['key'].split('/')[-1], conn_url, api_url, node.get('ttl', None))
class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members')):
@@ -91,7 +109,7 @@ class Etcd:
initialize = True if node else False
# 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']]
members = [Member.fromNode(n) for n in node['nodes']]
# get last leader operation
last_leader_operation = 0
@@ -110,7 +128,7 @@ class Etcd:
leader = m
break
if not leader:
leader = Member(node['value'], None, None)
leader = Member(node['value'], None, None, None)
return Cluster(initialize, leader, last_leader_operation, members)
elif status_code == 404:
+21 -17
View File
@@ -45,7 +45,7 @@ class Postgresql:
self.recovery_conf = os.path.join(self.data_dir, 'recovery.conf')
self.configuration_to_save = (os.path.join(self.data_dir, 'pg_hba.conf'),
os.path.join(self.data_dir, 'postgresql.conf'))
self.pid_path = os.path.join(self.data_dir, 'postmaster.pid')
self.postmaster_pid = os.path.join(self.data_dir, 'postmaster.pid')
self.trigger_file = config.get('recovery_conf', {}).get('trigger_file', None) or 'promote'
self.trigger_file = os.path.abspath(os.path.join(self.data_dir, self.trigger_file))
self.is_promoted = False
@@ -67,8 +67,14 @@ class Postgresql:
self.on_change_callback = on_change_callback
def get_local_address(self):
# TODO: try to get unix_socket_directory from postmaster.pid
return self.listen_addresses.split(',')[0].strip() + ':' + self.port
listen_addresses = self.listen_addresses.split(',')
local_address = listen_addresses[0].strip() # take first address from listen_addresses
for la in listen_addresses:
if la.strip() in ['*', '0.0.0.0']: # we are listening on *
local_address = 'localhost' # connection via localhost is preferred
break
return local_address + ':' + self.port
def connection(self):
if not self._connection or self._connection.closed != 0:
@@ -119,7 +125,7 @@ class Postgresql:
os.path.exists(self.trigger_file) and os.unlink(self.trigger_file)
def sync_from_leader(self, leader):
r = parseurl(leader.address)
r = parseurl(leader.conn_url)
pgpass = 'pgpass'
with open(pgpass, 'w') as f:
@@ -238,9 +244,9 @@ class Postgresql:
logger.error('Cannot start PostgreSQL because one is already running.')
return False
if os.path.exists(self.pid_path):
os.remove(self.pid_path)
logger.info('Removed %s', self.pid_path)
if os.path.exists(self.postmaster_pid):
os.remove(self.postmaster_pid)
logger.info('Removed %s', self.postmaster_pid)
ret = subprocess.call(self._pg_ctl + ['start', '-o', self.server_options()]) == 0
ret and self.load_replication_slots()
@@ -284,7 +290,7 @@ class Postgresql:
if member.hostname == self.name:
continue
try:
r = parseurl(member.address)
r = parseurl(member.conn_url)
member_conn = psycopg2.connect(**r)
member_conn.autocommit = True
member_cursor = member_conn.cursor()
@@ -304,12 +310,10 @@ class Postgresql:
def write_pg_hba(self):
with open(os.path.join(self.data_dir, 'pg_hba.conf'), 'a') as f:
f.write('\nhost replication {username} {network} md5\n'.format(**self.replication))
# allow TCP connections from the host's own address
f.write("\nhost postgres postgres samehost trust\n")
# allow TCP connections from the rest of the world with a password, prefer ssl
if self.config['parameters'].get('ssl', 'off').lower() == 'on':
f.write("\nhostssl all all 0.0.0.0/0 md5\n")
f.write("\nhost all all 0.0.0.0/0 md5\n")
for line in self.config.get('pg_hba', []):
if line.split()[0].strip() == 'hostssl' and self.config['parameters'].get('ssl', 'off').lower() != 'on':
continue
f.write(line + '\n')
@staticmethod
def primary_conninfo(leader_url):
@@ -320,7 +324,7 @@ class Postgresql:
if not os.path.isfile(self.recovery_conf):
return False
pattern = leader and leader.address and self.primary_conninfo(leader.address)
pattern = leader and leader.conn_url and self.primary_conninfo(leader.conn_url)
with open(self.recovery_conf, 'r') as f:
for line in f:
@@ -336,11 +340,11 @@ class Postgresql:
f.write("""standby_mode = 'on'
recovery_target_timeline = 'latest'
""")
if leader and leader.address:
if leader and leader.conn_url:
f.write("""
primary_slot_name = '{}'
primary_conninfo = '{}'
""".format(self.name, self.primary_conninfo(leader.address)))
""".format(self.name, self.primary_conninfo(leader.conn_url)))
for name, value in self.config.get('recovery_conf', {}).items():
f.write("{} = '{}'\n".format(name, value))
+26
View File
@@ -6,6 +6,32 @@ import time
received_sigchld = False
def lsn_to_bytes(value):
"""
>>> lsn_to_bytes('1/66000060')
6006243424
>>> lsn_to_bytes('j/66000060')
0
"""
try:
e = value.split('/')
if len(e) == 2 and len(e[0]) > 0 and len(e[1]) > 0:
return (int(e[0], 16) << 32) | int(e[1], 16)
except ValueError:
pass
return 0
def bytes_to_lsn(value):
"""
>>> bytes_to_lsn(6006243424)
'1/66000060'
"""
id = value >> 32
off = value & 0xffffffff
return '%x/%x' % (id, off)
def sigterm_handler(signo, stack_frame):
sys.exit()
+3
View File
@@ -12,6 +12,9 @@ postgresql:
connect_address: 127.0.0.1:5432
data_dir: data/postgresql0
maximum_lag_on_failover: 1048576 # 1 megabyte in bytes
pg_hba:
- host all all 0.0.0.0/0 md5
- hostssl all all 0.0.0.0/0 md5
replication:
username: replicator
password: rep-pass
+3
View File
@@ -12,6 +12,9 @@ postgresql:
connect_address: 127.0.0.1:5433
data_dir: data/postgresql1
maximum_lag_on_failover: 1048576 # 1 megabyte in bytes
pg_hba:
- host all all 0.0.0.0/0 md5
- hostssl all all 0.0.0.0/0 md5
replication:
username: replicator
password: rep-pass
+1 -2
View File
@@ -44,8 +44,7 @@ class MockRestApiServer(RestApiServer):
def __init__(self, Handler, path, *args):
self.governor = MockGovernor()
if len(args) > 0:
self._cursor = args[0]
self._cursor_holder = None
self.query = args[0]
Handler(MockRequest(path), ('0.0.0.0', 8080), self)
+2 -2
View File
@@ -30,11 +30,11 @@ def requests_get(url, **kwargs):
raise requests.exceptions.RequestException()
response = MockResponse()
if url.startswith('http://remote') or url.startswith('http://127.0.0.1'):
response.content = '{"action":"get","node":{"key":"/service/batman5","dir":true,"nodes":[{"key":"/service/batman5/initialize","value":"postgresql0","modifiedIndex":1582,"createdIndex":1582},{"key":"/service/batman5/leader","value":"postgresql1","expiration":"2015-05-15T09:11:00.037397538Z","ttl":21,"modifiedIndex":20728,"createdIndex":20434},{"key":"/service/batman5/optime","dir":true,"nodes":[{"key":"/service/batman5/optime/leader","value":"2164261704","modifiedIndex":20729,"createdIndex":20729}],"modifiedIndex":20437,"createdIndex":20437},{"key":"/service/batman5/members","dir":true,"nodes":[{"key":"/service/batman5/members/postgresql1","value":"postgres://replicator:[email protected]:5434/postgres","expiration":"2015-05-15T09:10:59.949384522Z","ttl":21,"modifiedIndex":20727,"createdIndex":20727},{"key":"/service/batman5/members/postgresql0","value":"postgres://replicator:[email protected]:5433/postgres","expiration":"2015-05-15T09:11:09.611860899Z","ttl":30,"modifiedIndex":20730,"createdIndex":20730}],"modifiedIndex":1581,"createdIndex":1581}],"modifiedIndex":1581,"createdIndex":1581}}'
response.content = '{"action":"get","node":{"key":"/service/batman5","dir":true,"nodes":[{"key":"/service/batman5/initialize","value":"postgresql0","modifiedIndex":1582,"createdIndex":1582},{"key":"/service/batman5/leader","value":"postgresql1","expiration":"2015-05-15T09:11:00.037397538Z","ttl":21,"modifiedIndex":20728,"createdIndex":20434},{"key":"/service/batman5/optime","dir":true,"nodes":[{"key":"/service/batman5/optime/leader","value":"2164261704","modifiedIndex":20729,"createdIndex":20729}],"modifiedIndex":20437,"createdIndex":20437},{"key":"/service/batman5/members","dir":true,"nodes":[{"key":"/service/batman5/members/postgresql1","value":"postgres://replicator:[email protected]:5434/postgres?application_name=http://127.0.0.1:8009/governor","expiration":"2015-05-15T09:10:59.949384522Z","ttl":21,"modifiedIndex":20727,"createdIndex":20727},{"key":"/service/batman5/members/postgresql0","value":"postgres://replicator:[email protected]:5433/postgres?application_name=http://127.0.0.1:8008/governor","expiration":"2015-05-15T09:11:09.611860899Z","ttl":30,"modifiedIndex":20730,"createdIndex":20730}],"modifiedIndex":1581,"createdIndex":1581}],"modifiedIndex":1581,"createdIndex":1581}}'
elif url.startswith('http://other'):
response.status_code = 404
elif url.startswith('http://noleader'):
response.content = '{"action":"get","node":{"key":"/service/batman5","dir":true,"nodes":[{"key":"/service/batman5/initialize","value":"postgresql0","modifiedIndex":1582,"createdIndex":1582},{"key":"/service/batman5/leader","value":"postgresql1","expiration":"2015-05-15T09:11:00.037397538Z","ttl":21,"modifiedIndex":20728,"createdIndex":20434},{"key":"/service/batman5/optime","dir":true,"nodes":[{"key":"/service/batman5/optime/leader","value":"2164261704","modifiedIndex":20729,"createdIndex":20729}],"modifiedIndex":20437,"createdIndex":20437},{"key":"/service/batman5/members","dir":true,"nodes":[{"key":"/service/batman5/members/postgresql0","value":"postgres://replicator:[email protected]:5433/postgres","expiration":"2015-05-15T09:11:09.611860899Z","ttl":30,"modifiedIndex":20730,"createdIndex":20730}],"modifiedIndex":1581,"createdIndex":1581}],"modifiedIndex":1581,"createdIndex":1581}}'
response.content = '{"action":"get","node":{"key":"/service/batman5","dir":true,"nodes":[{"key":"/service/batman5/initialize","value":"postgresql0","modifiedIndex":1582,"createdIndex":1582},{"key":"/service/batman5/leader","value":"postgresql1","expiration":"2015-05-15T09:11:00.037397538Z","ttl":21,"modifiedIndex":20728,"createdIndex":20434},{"key":"/service/batman5/optime","dir":true,"nodes":[{"key":"/service/batman5/optime/leader","value":"2164261704","modifiedIndex":20729,"createdIndex":20729}],"modifiedIndex":20437,"createdIndex":20437},{"key":"/service/batman5/members","dir":true,"nodes":[{"key":"/service/batman5/members/postgresql0","value":"postgres://replicator:[email protected]:5433/postgres?application_name=http://127.0.0.1:8008/governor","expiration":"2015-05-15T09:11:09.611860899Z","ttl":30,"modifiedIndex":20730,"createdIndex":20730}],"modifiedIndex":1581,"createdIndex":1581}],"modifiedIndex":1581,"createdIndex":1581}}'
else:
response.status_code = 404
response.ok = False
+11 -10
View File
@@ -110,9 +110,10 @@ class TestPostgresql(unittest.TestCase):
def set_up(self):
subprocess.call = subprocess_call
shutil.copy = nop
self.p = Postgresql({'name': 'test0', 'data_dir': 'data/test0', 'listen': '127.0.0.1, 127.0.0.2:5432',
'connect_address': '127.0.0.2:5432', 'superuser': {'password': ''},
'admin': {'username': 'admin', 'password': 'admin'},
self.p = Postgresql({'name': 'test0', 'data_dir': 'data/test0', 'listen': '127.0.0.1, *:5432',
'connect_address': '127.0.0.2:5432',
'pg_hba': ['hostssl all all 0.0.0.0/0 md5', 'host all all 0.0.0.0/0 md5'],
'superuser': {'password': ''}, 'admin': {'username': 'admin', 'password': 'admin'},
'replication': {'username': 'replicator',
'password': 'rep-pass',
'network': '127.0.0.1/32'},
@@ -120,7 +121,7 @@ class TestPostgresql(unittest.TestCase):
psycopg2.connect = psycopg2_connect
if not os.path.exists(self.p.data_dir):
os.makedirs(self.p.data_dir)
self.leader = Member('leader', 'postgres://replicator:[email protected]:5434/postgres', 28)
self.leader = Member('leader', 'postgres://replicator:[email protected]:5434/postgres', None, 28)
def tear_down(self):
shutil.rmtree('data')
@@ -147,12 +148,12 @@ class TestPostgresql(unittest.TestCase):
self.p.follow_the_leader(None)
self.p.demote(self.leader)
self.p.follow_the_leader(self.leader)
self.p.follow_the_leader(Member('leader', 'postgres://replicator:[email protected]:5435/postgres', 28))
self.p.follow_the_leader(Member('leader', 'postgres://replicator:[email protected]:5435/postgres', None, 28))
def test_create_replication_slots(self):
self.p.start()
me = Member('test0', 'postgres://replicator:[email protected]:5434/postgres', 28)
other = Member('test1', 'postgres://replicator:[email protected]:5433/postgres', 28)
me = Member('test0', 'postgres://replicator:[email protected]:5434/postgres', None, 28)
other = Member('test1', 'postgres://replicator:[email protected]:5433/postgres', None, 28)
cluster = Cluster(True, self.leader, 0, [me, other, self.leader])
self.p.create_replication_slots(cluster)
@@ -167,9 +168,9 @@ class TestPostgresql(unittest.TestCase):
self.assertRaises(psycopg2.OperationalError, self.p.query, 'blabla')
def test_is_healthiest_node(self):
leader = Member('leader', 'postgres://replicator:[email protected]:5435/postgres', 28)
me = Member('test0', 'postgres://replicator:[email protected]:5434/postgres', 28)
other = Member('test1', 'postgres://replicator:[email protected]:5433/postgres', 28)
leader = Member('leader', 'postgres://replicator:[email protected]:5435/postgres', None, 28)
me = Member('test0', 'postgres://replicator:[email protected]:5434/postgres', None, 28)
other = Member('test1', 'postgres://replicator:[email protected]:5433/postgres', None, 28)
cluster = Cluster(True, leader, 0, [me, other, leader])
self.assertTrue(self.p.is_healthiest_node(cluster))
self.p.is_leader = false