diff --git a/governor.py b/governor.py index fe2ecd8a..ba48fefd 100755 --- a/governor.py +++ b/governor.py @@ -22,9 +22,15 @@ class Governor: host, port = config['restapi']['listen'].split(':') self.api = RestApiServer(self, config['restapi']) self.next_run = time.time() + self.shutdown_member_ttl = 300 def touch_member(self, ttl=None): connection_string = self.postgresql.connection_string + '?application_name=' + self.api.connection_string + if self.ha.cluster: + for m in self.ha.cluster.members: + # Do not update member TTL when it is far from being expired + if m.hostname == self.postgresql.name and m.real_ttl() > self.shutdown_member_ttl: + return True return self.etcd.touch_member(self.postgresql.name, connection_string, ttl) def initialize(self): @@ -91,7 +97,7 @@ def main(): except KeyboardInterrupt: pass finally: - governor.touch_member(300) # schedule member removal + governor.touch_member(governor.shutdown_member_ttl) # schedule member removal governor.postgresql.stop() governor.etcd.delete_leader(governor.postgresql.name) diff --git a/helpers/etcd.py b/helpers/etcd.py index 3a222a3b..f63cc364 100644 --- a/helpers/etcd.py +++ b/helpers/etcd.py @@ -8,7 +8,7 @@ from collections import namedtuple from dns.exception import DNSException from dns import resolver from helpers.errors import CurrentLeaderError, EtcdError, EtcdConnectionFailed -from helpers.utils import sleep +from helpers.utils import calculate_ttl, sleep from requests.exceptions import RequestException if sys.hexversion >= 0x03000000: @@ -30,6 +30,9 @@ class Member(namedtuple('Member', 'hostname,conn_url,api_url,expiration,ttl')): ttl = node.get('ttl', None) return Member(node['key'].split('/')[-1], conn_url, api_url, expiration, ttl) + def real_ttl(self): + return calculate_ttl(self.expiration) + class Cluster(namedtuple('Cluster', 'initialize,leader,last_leader_operation,members')): diff --git a/helpers/utils.py b/helpers/utils.py index e08dd657..c12f389a 100644 --- a/helpers/utils.py +++ b/helpers/utils.py @@ -1,10 +1,39 @@ +import datetime import os +import re import signal import sys import time received_sigchld = False +_DATE_TIME_RE = re.compile(r'''^ +(?P\d{4})\-(?P\d{2})\-(?P\d{2}) # date +T +(?P\d{2}):(?P\d{2}):(?P\d{2})\.(?P\d{6}) # time +\d*Z$''', re.X) + + +def parse_datetime(time_str): + """ + >>> parse_datetime('2015-06-10T12:56:30.552539016Z') + datetime.datetime(2015, 6, 10, 12, 56, 30, 552539) + >>> parse_datetime('2015-06-10 12:56:30.552539016Z') + """ + m = _DATE_TIME_RE.match(time_str) + if not m: + return None + p = dict((n, int(m.group(n))) for n in 'year month day hour minute second microsecond'.split(' ')) + return datetime.datetime(**p) + + +def calculate_ttl(expiration): + expiration = parse_datetime(expiration) + if not expiration: + return None + now = datetime.datetime.utcnow() + return int((expiration - now).total_seconds()) + def lsn_to_bytes(value): """ diff --git a/tests/test_etcd.py b/tests/test_etcd.py index c2e07c55..b9e88df7 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -1,3 +1,4 @@ +import datetime import dns.resolver import json import requests @@ -7,7 +8,7 @@ import unittest from dns.exception import DNSException from helpers.errors import EtcdError, CurrentLeaderError, EtcdConnectionFailed -from helpers.etcd import Client, Cluster, Etcd +from helpers.etcd import Client, Cluster, Etcd, Member class MockResponse: @@ -94,6 +95,18 @@ def socket_getaddrinfo(*args): raise socket.error() +class TestMember(unittest.TestCase): + + def __init__(self, method_name='runTest'): + super(TestMember, self).__init__(method_name) + + def test_real_ttl(self): + now = datetime.datetime.utcnow() + member = Member('a', 'b', 'c', (now + datetime.timedelta(seconds=2)).strftime('%Y-%m-%dT%H:%M:%S.%fZ'), None) + self.assertLess(member.real_ttl(), 2) + self.assertIsNone(Member('a', 'b', 'c', '', None).real_ttl()) + + class TestClient(unittest.TestCase): def __init__(self, method_name='runTest'): diff --git a/tests/test_governor.py b/tests/test_governor.py index 9ecb88f1..b84637a7 100644 --- a/tests/test_governor.py +++ b/tests/test_governor.py @@ -1,3 +1,4 @@ +import datetime import psycopg2 import requests import subprocess @@ -7,6 +8,7 @@ import unittest import yaml from governor import Governor, main +from helpers.etcd import Cluster, Member from test_ha import true, false from test_postgresql import Postgresql, subprocess_call, psycopg2_connect from test_etcd import requests_get, requests_put, requests_delete @@ -67,6 +69,13 @@ class TestGovernor(unittest.TestCase): return False return True + def test_touch_member(self): + now = datetime.datetime.utcnow() + member = Member(self.g.postgresql.name, 'b', 'c', (now + datetime.timedelta( + seconds=self.g.shutdown_member_ttl + 10)).strftime('%Y-%m-%dT%H:%M:%S.%fZ'), None) + self.g.ha.cluster = Cluster(True, member, 0, [member]) + self.g.touch_member() + def test_governor_initialize(self): self.g.etcd.client._base_uri = 'http://remote' self.g.postgresql.data_directory_empty = true