From 0d75152d540a7d5dc9d1434e4d1a2941f5cf7355 Mon Sep 17 00:00:00 2001 From: Christopher Winslett Date: Mon, 16 Mar 2015 09:47:52 -0700 Subject: [PATCH] make healthiest node work --- helpers/etcd.py | 20 ++++++++++++++------ helpers/ha.py | 2 +- helpers/postgresql.py | 21 ++++++++++++++++++++- 3 files changed, 35 insertions(+), 8 deletions(-) diff --git a/helpers/etcd.py b/helpers/etcd.py index 909639ad..7abb12af 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -37,12 +37,6 @@ class Etcd: def client_url(self, path): return "http://%s/v2/keys/service/%s%s" % (self.host, self.scope, path) - def xlog_position(member): - try: - return self.get_client_path("/service/postgresql/xlog-position/%s" % member)["node"]["value"] - except urllib2.HTTPError: - return None - def current_leader(self): try: hostname = self.get_client_path("/leader")["node"]["value"] @@ -54,6 +48,20 @@ class Etcd: return None raise helpers.errors.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: + return None + raise helpers.errors.CurrentLeaderError("Etcd is not responding properly") + def touch_member(self, member, connection_string): self.put_client_path("/members/%s" % member, {"value": connection_string}) diff --git a/helpers/ha.py b/helpers/ha.py index f4d703ee..8ea95056 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -33,7 +33,7 @@ class Ha: try: if self.state_handler.is_healthy(): if self.is_unlocked(): - if self.state_handler.is_healthiest_node(): + 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() diff --git a/helpers/postgresql.py b/helpers/postgresql.py index 01a0cc1f..85a301d5 100644 --- a/helpers/postgresql.py +++ b/helpers/postgresql.py @@ -104,7 +104,23 @@ class Postgresql: return True - def is_healthiest_node(self): + def is_healthiest_node(self, members): + for member in members: + if member["hostname"] == self.name: + continue + try: + member_conn = psycopg2.connect(member["address"]) + member_conn.autocommit = True + member_cursor = member_conn.cursor() + member_cursor.execute("SELECT '%s'::pg_lsn - pg_last_xlog_replay_location() AS bytes;" % self.xlog_position()) + xlog_diff = member_cursor.fetchone()[0] + print [self.name, member["hostname"], xlog_diff] + if xlog_diff < 0: + member_cursor.close() + return False + member_cursor.close() + except psycopg2.OperationalError: + continue return True def replication_slot_name(self): @@ -150,3 +166,6 @@ recovery_target_timeline = 'latest' def create_replication_user(self): self.query("CREATE USER \"%s\" WITH REPLICATION ENCRYPTED PASSWORD '%s';" % (self.replication["username"], self.replication["password"])) + + def xlog_position(self): + return self.query("SELECT pg_last_xlog_replay_location();").fetchone()[0]