From a56b346295131c4ddd5f26bb8e2f930875de4a8a Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Tue, 12 May 2015 11:58:47 +0200 Subject: [PATCH 1/7] Hardened the health checks and implemented a real status page. For the postgresql helper, hardened the code to get "the" cursor of the postgresql instance. For the statuspage, a small status json is returned. To find out what status a PostgreSQL cluster is in we use the cursor (instead of the provided query() function), as the query function does some retrying etc. For the healthcheck we want to simple provide an answer to a simple query, if we have to reconnect, we are not *that* healthy anyway. Dropped catching exceptions in the do_GET block, as the HTTPServer will do that nicely for us anyway. --- helpers/postgresql.py | 17 +++++++---- helpers/statuspage.py | 66 +++++++++++++++++++++++++++++++------------ 2 files changed, 60 insertions(+), 23 deletions(-) diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 4bde365c..7f4ccaad 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -12,6 +12,13 @@ class Postgresql: def __init__(self, config, aws_host_address=None): self.name = config["name"] self.host, self.port = config["listen"].split(":") + self.libpq_parameters = { + 'host' : aws_host_address or self.host, + 'port' : self.port, + 'fallback_application_name' : 'Governor', + 'connect_timeout' : 5, + 'options' : '-c statement_timeout=2000' + } self.data_dir = config["data_dir"] self.replication = config["replication"] self.superuser = config.get('superuser') @@ -20,15 +27,15 @@ class Postgresql: self.config = config self.cursor_holder = None - connection_host = aws_host_address or self.host - self.connection_string = "postgres://%s:%s@%s:%s/postgres" % (self.replication["username"], self.replication["password"], connection_host, self.port) + self.connection_string = "postgres://%s:%s@%s:%s/postgres" % (self.replication["username"], self.replication["password"], self.libpq_parameters['host'], self.port) self.conn = None def cursor(self): - if not self.cursor_holder: - self.conn = psycopg2.connect("postgres://%s:%s/postgres" % (self.host, self.port)) - self.conn.autocommit = True + if (self.cursor_holder is None) or self.cursor_holder.closed: + if (self.conn is None) or self.conn.closed: + self.conn = psycopg2.connect(**self.libpq_parameters) + self.conn.autocommit = True self.cursor_holder = self.conn.cursor() return self.cursor_holder diff --git a/helpers/statuspage.py b/helpers/statuspage.py index 67fdb4ab..9aae769b 100644 --- a/helpers/statuspage.py +++ b/helpers/statuspage.py @@ -2,28 +2,59 @@ # -*- coding: utf-8 -*- from BaseHTTPServer import BaseHTTPRequestHandler, HTTPServer +import json class StatusPage(BaseHTTPRequestHandler): def do_GET(self): - try: - if self.path == '/pg_master': - response = (200 if self.server.postgresql.is_leader else 503) - self.send_response(response) - elif self.path == '/pg_slave': - response = (503 if self.server.postgresql.is_leader else 200) - self.send_response(response) - elif self.path == '/pg_status': - self.send_response(200) - self.end_headers() - self.wfile.write(self.server.postgresql.status()) - else: - self.send_response(404) - except Exception, e: - self.send_response(500) - self.end_headers() - self.wfile.write(repr(e)) + if self.path == '/pg_master': + self.pg_master() + elif self.path == '/pg_slave': + self.pg_slave() + elif self.path == '/pg_status': + self.pg_status() + else: + self.send_response(404) + + def pg_master(self): + if not self.pg_is_in_recovery(): + self.send_response(200) + return + + self.send_response(503) + + def pg_slave(self): + if self.pg_is_in_recovery(): + self.send_response(200) + return + + self.send_response(503) + + def pg_is_in_recovery(self): + cursor = self.server.postgresql.cursor() + cursor.execute('SELECT pg_is_in_recovery()') + res = cursor.fetchone() + return res[0] + + def pg_status(self): + cursor = self.server.postgresql.cursor() + cursor.execute(""" + SELECT pg_is_in_recovery(), + to_char(pg_last_xact_replay_timestamp(), 'YYYY-MM-DD HH24:MI:SS.MS TZ'), + extract(epoch from now() - pg_last_xact_replay_timestamp()), + inet_server_addr(), + inet_server_port(), + to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ') + """) + res = cursor.fetchone() + status = {'role': ('master' if not res[0] else 'slave'), 'recovery': {'last_transaction_replayed': res[1], + 'delay': res[2]}, 'server': {'hostaddr': res[3], 'port': res[4], 'start_time': res[5]}} + + self.send_response(200) + self.send_header('Content-Type', 'application/json') + self.end_headers() + self.wfile.write(json.dumps(status)) def getHTTPServer(postgresql, http_port=8081, listen_address='0.0.0.0'): @@ -36,7 +67,6 @@ def getHTTPServer(postgresql, http_port=8081, listen_address='0.0.0.0'): if __name__ == '__main__': import sys import logging - from BaseHTTPServer import HTTPServer logging.basicConfig(format='%(levelname)-6s %(asctime)s - %(message)s', level=logging.DEBUG) logging.debug('Starting as a standalone application') From d25fdd41a6a4f376a3c54877ec18dcd3b5790b8a Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Tue, 12 May 2015 14:14:20 +0200 Subject: [PATCH 2/7] Refactoring of the status page for the healthcheck. Less functions, as the code is readable enough without them. Always return some content to the client, instead of a response only. --- helpers/statuspage.py | 44 +++++++++++++++++++++---------------------- 1 file changed, 21 insertions(+), 23 deletions(-) diff --git a/helpers/statuspage.py b/helpers/statuspage.py index 9aae769b..ee535b54 100644 --- a/helpers/statuspage.py +++ b/helpers/statuspage.py @@ -9,27 +9,23 @@ class StatusPage(BaseHTTPRequestHandler): def do_GET(self): if self.path == '/pg_master': - self.pg_master() + if not self.pg_is_in_recovery(): + response, content = 200, 'I am currently a master' + else: + response, content = 503, 'I am not a master' elif self.path == '/pg_slave': - self.pg_slave() + if self.pg_is_in_recovery(): + response, content = 200, 'I am currently a slave' + else: + response, content = 503, 'I am not a slave' elif self.path == '/pg_status': - self.pg_status() + response, content = 200, self.pg_status() else: - self.send_response(404) + response, content = 404, 'Page not found' - def pg_master(self): - if not self.pg_is_in_recovery(): - self.send_response(200) - return - - self.send_response(503) - - def pg_slave(self): - if self.pg_is_in_recovery(): - self.send_response(200) - return - - self.send_response(503) + self.send_response(response) + self.end_headers() + self.wfile.write(content) def pg_is_in_recovery(self): cursor = self.server.postgresql.cursor() @@ -48,13 +44,12 @@ class StatusPage(BaseHTTPRequestHandler): to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ') """) res = cursor.fetchone() - status = {'role': ('master' if not res[0] else 'slave'), 'recovery': {'last_transaction_replayed': res[1], + status = {'role': ('master' if not res[0] else 'slave'), 'recovery': {'last_transaction_timestamp': res[1], 'delay': res[2]}, 'server': {'hostaddr': res[3], 'port': res[4], 'start_time': res[5]}} - self.send_response(200) self.send_header('Content-Type', 'application/json') - self.end_headers() - self.wfile.write(json.dumps(status)) + + return json.dumps(status) def getHTTPServer(postgresql, http_port=8081, listen_address='0.0.0.0'): @@ -84,5 +79,8 @@ if __name__ == '__main__': postgres_config['listen'] = sys.argv[1] postgresql = Postgresql(postgres_config, aws_host_address) - getHTTPServer(postgresql, 8081, '0.0.0.0').serve_forever() - logging.debug('Abc') + http_port = 8081 + if len(sys.argv) > 2: + http_port = int(sys.argv[2]) + + getHTTPServer(postgresql, http_port, '0.0.0.0').serve_forever() From 362b6b4fa2f23329cff0b0cb8c4f8aab09f8c2c2 Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Tue, 12 May 2015 14:48:40 +0200 Subject: [PATCH 3/7] Made healtcheck port configurable, fixed standalone status page. --- governor.py | 2 +- helpers/statuspage.py | 8 +++++--- postgres0.yml | 1 + postgres1.yml | 1 + 4 files changed, 8 insertions(+), 4 deletions(-) diff --git a/governor.py b/governor.py index 1431cc66..638305df 100755 --- a/governor.py +++ b/governor.py @@ -80,7 +80,7 @@ def main(): governor = Governor(config) # Start the http_server to serve a simple healthcheck - http_server = getHTTPServer(governor.postgresql, http_port=8008, listen_address='0.0.0.0') + http_server = getHTTPServer(governor.postgresql, http_port=config.get('healtcheck_port', 8080), listen_address='0.0.0.0') http_thread = threading.Thread(target=http_server.serve_forever, args=()) http_thread.daemon = True diff --git a/helpers/statuspage.py b/helpers/statuspage.py index ee535b54..cec68432 100644 --- a/helpers/statuspage.py +++ b/helpers/statuspage.py @@ -44,8 +44,8 @@ class StatusPage(BaseHTTPRequestHandler): to_char(pg_postmaster_start_time(), 'YYYY-MM-DD HH24:MI:SS.MS TZ') """) res = cursor.fetchone() - status = {'role': ('master' if not res[0] else 'slave'), 'recovery': {'last_transaction_timestamp': res[1], - 'delay': res[2]}, 'server': {'hostaddr': res[3], 'port': res[4], 'start_time': res[5]}} + status = {'role': ('master' if not res[0] else 'slave'), 'recovery': {'last_transaction_timestamp': res[1]}, + 'server': {'hostaddr': res[3], 'port': res[4], 'start_time': res[5]}} self.send_header('Content-Type', 'application/json') @@ -71,8 +71,10 @@ if __name__ == '__main__': postgres_config = { 'name': 'dummy', 'listen': 'localhost:5432', - 'data_dir': None, + 'data_dir': 'nonsense', 'replication': {'username': None, 'password': None}, + 'superuser': None, + 'admin': None, } aws_host_address = None if len(sys.argv) > 1: diff --git a/postgres0.yml b/postgres0.yml index e7d05c5a..14c22892 100644 --- a/postgres0.yml +++ b/postgres0.yml @@ -1,5 +1,6 @@ loop_wait: 10 aws_use_host_address: "on" +healthcheck_port: 8080 etcd: scope: batman ttl: 30 diff --git a/postgres1.yml b/postgres1.yml index f18ebd17..bb203c6b 100644 --- a/postgres1.yml +++ b/postgres1.yml @@ -1,5 +1,6 @@ loop_wait: 10 aws_use_host_address: "on" +healthcheck_port: 8081 etcd: scope: batman ttl: 30 From c4168bfeb5bef26382f4f4f56807c48f0cd8240a Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Tue, 12 May 2015 14:56:02 +0200 Subject: [PATCH 4/7] Bugfix: Headers sent before http status for the pg_status page. --- helpers/statuspage.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/helpers/statuspage.py b/helpers/statuspage.py index cec68432..2819beeb 100644 --- a/helpers/statuspage.py +++ b/helpers/statuspage.py @@ -8,6 +8,7 @@ import json class StatusPage(BaseHTTPRequestHandler): def do_GET(self): + content_type='text/plain' if self.path == '/pg_master': if not self.pg_is_in_recovery(): response, content = 200, 'I am currently a master' @@ -20,10 +21,12 @@ class StatusPage(BaseHTTPRequestHandler): response, content = 503, 'I am not a slave' elif self.path == '/pg_status': response, content = 200, self.pg_status() + content_type = 'application/json' else: response, content = 404, 'Page not found' self.send_response(response) + self.send_header('Content-Type', content_type) self.end_headers() self.wfile.write(content) @@ -47,8 +50,6 @@ class StatusPage(BaseHTTPRequestHandler): status = {'role': ('master' if not res[0] else 'slave'), 'recovery': {'last_transaction_timestamp': res[1]}, 'server': {'hostaddr': res[3], 'port': res[4], 'start_time': res[5]}} - self.send_header('Content-Type', 'application/json') - return json.dumps(status) From 2b56379b85ffea736c18087060818b69ff8086ee Mon Sep 17 00:00:00 2001 From: Feike Steenbergen Date: Tue, 12 May 2015 15:48:30 +0200 Subject: [PATCH 5/7] Change default port for the health check to 8008. --- governor.py | 2 +- postgres0.yml | 2 +- postgres1.yml | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/governor.py b/governor.py index 638305df..2f9132cf 100755 --- a/governor.py +++ b/governor.py @@ -80,7 +80,7 @@ def main(): governor = Governor(config) # Start the http_server to serve a simple healthcheck - http_server = getHTTPServer(governor.postgresql, http_port=config.get('healtcheck_port', 8080), listen_address='0.0.0.0') + http_server = getHTTPServer(governor.postgresql, http_port=config.get('healtcheck_port', 8008), listen_address='0.0.0.0') http_thread = threading.Thread(target=http_server.serve_forever, args=()) http_thread.daemon = True diff --git a/postgres0.yml b/postgres0.yml index 14c22892..ca577a0b 100644 --- a/postgres0.yml +++ b/postgres0.yml @@ -1,6 +1,6 @@ loop_wait: 10 aws_use_host_address: "on" -healthcheck_port: 8080 +healthcheck_port: 8008 etcd: scope: batman ttl: 30 diff --git a/postgres1.yml b/postgres1.yml index bb203c6b..8f0e493e 100644 --- a/postgres1.yml +++ b/postgres1.yml @@ -1,6 +1,6 @@ loop_wait: 10 aws_use_host_address: "on" -healthcheck_port: 8081 +healthcheck_port: 8009 etcd: scope: batman ttl: 30 From 8ad751b6610aca0f3064a831d3c55a4c366d12e9 Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Wed, 13 May 2015 10:05:12 +0200 Subject: [PATCH 6/7] add a simplest SIGCHLD handler in order to avoid zombies in the docker container governed by the script --- governor.py | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/governor.py b/governor.py index 2f9132cf..924535f6 100755 --- a/governor.py +++ b/governor.py @@ -19,6 +19,14 @@ def sigterm_handler(signo, stack_frame): sys.exit() +# handle SIGCHILD, since we are the equivalent of the INIT process +def sigchld_handler(signo, stack_frame): + try: + os.waitpid(-1, os.WNOHANG) + except OSError: + pass + + class Governor: INSTANCE_METADATA_URL = "http://169.254.169.254/latest/meta-data/" @@ -96,4 +104,5 @@ def main(): if __name__ == '__main__': logging.basicConfig(format='%(asctime)s %(levelname)s: %(message)s', level=logging.INFO) signal.signal(signal.SIGTERM, sigterm_handler) + signal.signal(signal.SIGCHLD, sigchld_handler) main() From 609a4b2fc60509f8891589a96c2688664bd8789e Mon Sep 17 00:00:00 2001 From: Oleksii Kliukin Date: Wed, 13 May 2015 12:22:17 +0200 Subject: [PATCH 7/7] loop until we reap all terminated children. --- governor.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/governor.py b/governor.py index 924535f6..d568da86 100755 --- a/governor.py +++ b/governor.py @@ -22,7 +22,10 @@ def sigterm_handler(signo, stack_frame): # handle SIGCHILD, since we are the equivalent of the INIT process def sigchld_handler(signo, stack_frame): try: - os.waitpid(-1, os.WNOHANG) + while True: + ret = os.waitpid(-1, os.WNOHANG) + if ret == (0, 0): + break except OSError: pass