Merge pull request #17 from zalando/features/refactoring

Store api url of each governor together with connect url of postgres
This commit is contained in:
Feike Steenbergen
2015-05-28 09:26:09 +02:00
5 changed files with 14 additions and 9 deletions
+5 -4
View File
@@ -35,9 +35,12 @@ class Governor:
self.etcd = Etcd(config['etcd'])
self.postgresql = Postgresql(config['postgresql'])
self.ha = Ha(self.postgresql, self.etcd)
host, port = config['restapi']['listen'].split(':')
self.api = RestApiServer(self, config['restapi'])
def touch_member(self, ttl=None):
return self.etcd.touch_member(self.postgresql.name, self.postgresql.connection_string, ttl)
connection_string = self.postgresql.connection_string + '?application_name=' + self.api.connection_string
return self.etcd.touch_member(self.postgresql.name, connection_string, ttl)
def initialize(self):
# wait for etcd to be available
@@ -66,6 +69,7 @@ class Governor:
self.postgresql.load_replication_slots()
def run(self):
self.api.start()
while True:
self.touch_member()
logging.info(self.ha.run_cycle())
@@ -87,9 +91,6 @@ def main():
governor = Governor(config)
try:
governor.initialize()
# Start the http_server to serve a simple healthcheck
host, port = config['restapi']['listen'].split(':')
RestApiServer(governor, host, int(port)).start()
governor.run()
finally:
governor.touch_member(300) # schedule member removal
+6 -4
View File
@@ -34,7 +34,7 @@ class RestApiHandler(BaseHTTPRequestHandler):
def get_postgresql_status(self):
if not self.server.governor.postgresql.is_running():
return {'running': False}
cursor = self.server.cursor()
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()
@@ -59,14 +59,16 @@ class RestApiHandler(BaseHTTPRequestHandler):
class RestApiServer(HTTPServer, Thread):
def __init__(self, governor, listen_address='0.0.0.0', listen_port=8080):
HTTPServer.__init__(self, (listen_address, listen_port), RestApiHandler)
def __init__(self, governor, config):
self.connection_string = 'http://{}/governor'.format(config.get('connect_address', None) or config['listen'])
host, port = config['listen'].split(':')
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):
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
+1
View File
@@ -1,6 +1,7 @@
loop_wait: 10
restapi:
listen: 127.0.0.1:8008
connect_address: 127.0.0.1:8008
etcd:
scope: batman
ttl: 30
+1
View File
@@ -1,6 +1,7 @@
loop_wait: 10
restapi:
listen: 127.0.0.1:8009
connect_address: 127.0.0.1:8009
etcd:
scope: batman
ttl: 30
+1 -1
View File
@@ -28,7 +28,7 @@ def requests_get(url, **kwargs):
if url.startswith('http://local'):
raise requests.exceptions.RequestException()
response = MockResponse()
if url.startswith('http://remote'):
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}}'
elif url.startswith('http://other'):
response.status_code = 404