diff --git a/governor.py b/governor.py index eb029b22..cfe1c546 100755 --- a/governor.py +++ b/governor.py @@ -16,14 +16,19 @@ class Governor: def __init__(self, config): self.nap_time = config['loop_wait'] - self.etcd = Etcd(config['etcd']) self.postgresql = Postgresql(config['postgresql']) - self.ha = Ha(self.postgresql, self.etcd) + self.ha = Ha(self.postgresql, self.get_dcs(config)) host, port = config['restapi']['listen'].split(':') self.api = RestApiServer(self, config['restapi']) self.next_run = time.time() self.shutdown_member_ttl = 300 + @staticmethod + def get_dcs(config): + if 'etcd' in config: + return Etcd(config['etcd']) + raise Exception('Can not find sutable configuration of distributed configuration store') + def touch_member(self, ttl=None): connection_string = self.postgresql.connection_string + '?application_name=' + self.api.connection_string if self.ha.cluster: @@ -31,7 +36,7 @@ class Governor: # Do not update member TTL when it is far from being expired if m.name == self.postgresql.name and m.real_ttl() > self.shutdown_member_ttl: return True - return self.etcd.touch_member(self.postgresql.name, connection_string, ttl) + return self.ha.dcs.touch_member(self.postgresql.name, connection_string, ttl) def initialize(self): # wait for etcd to be available @@ -42,14 +47,14 @@ class Governor: # is data directory empty? if self.postgresql.data_directory_empty(): # racing to initialize - if self.etcd.race('/initialize', self.postgresql.name): + if self.ha.dcs.race('/initialize', self.postgresql.name): self.postgresql.initialize() - self.etcd.take_leader(self.postgresql.name) + self.ha.dcs.take_leader(self.postgresql.name) self.postgresql.start() self.postgresql.create_replication_user() else: while True: - leader = self.etcd.current_leader() + leader = self.ha.dcs.current_leader() if leader and self.postgresql.sync_from_leader(leader): self.postgresql.write_recovery_conf(leader) self.postgresql.start() @@ -105,7 +110,7 @@ def main(): finally: governor.touch_member(governor.shutdown_member_ttl) # schedule member removal governor.postgresql.stop() - governor.etcd.delete_leader(governor.postgresql.name) + governor.ha.dcs.delete_leader(governor.postgresql.name) if __name__ == '__main__': diff --git a/helpers/errors.py b/helpers/errors.py deleted file mode 100644 index 3a56a1e0..00000000 --- a/helpers/errors.py +++ /dev/null @@ -1,15 +0,0 @@ -class EtcdError(Exception): - - def __init__(self, value): - self.value = value - - def __str__(self): - return repr(self.value) - - -class CurrentLeaderError(EtcdError): - pass - - -class EtcdConnectionFailed(EtcdError): - pass diff --git a/helpers/ha.py b/helpers/ha.py index 52462b92..1ffa06ae 100644 --- a/helpers/ha.py +++ b/helpers/ha.py @@ -10,17 +10,17 @@ class Ha: def __init__(self, state_handler, etcd): self.state_handler = state_handler - self.etcd = etcd + self.dcs = etcd self.cluster = None - def load_cluster_from_etcd(self): - self.cluster = self.etcd.get_cluster() + def load_cluster_from_dcs(self): + self.cluster = self.dcs.get_cluster() def acquire_lock(self): - return self.etcd.attempt_to_acquire_leader(self.state_handler.name) + return self.dcs.attempt_to_acquire_leader(self.state_handler.name) def update_lock(self): - return self.etcd.update_leader(self.state_handler) + return self.dcs.update_leader(self.state_handler) def has_lock(self): lock_owner = self.cluster.leader and self.cluster.leader.name @@ -35,7 +35,7 @@ class Ha: def run_cycle(self): try: - self.load_cluster_from_etcd() + self.load_cluster_from_dcs() if not self.state_handler.is_healthy(): has_lock = self.has_lock() self.state_handler.write_recovery_conf(None if has_lock else self.cluster.leader) @@ -43,7 +43,7 @@ class Ha: if not has_lock: return 'started as a secondary' logger.info('started as readonly because i had the session lock') - self.load_cluster_from_etcd() + self.load_cluster_from_dcs() if self.cluster.is_unlocked(): if self.state_handler.is_healthiest_node(self.cluster): @@ -54,7 +54,7 @@ class Ha: self.state_handler.promote() return 'promoted self to leader by acquiring session lock' else: - self.load_cluster_from_etcd() + self.load_cluster_from_dcs() if self.state_handler.is_leader(): self.demote() return 'demoted self due after trying and failing to obtain lock' @@ -62,7 +62,7 @@ class Ha: self.follow_the_leader() return 'following new leader after trying and failing to obtain lock' else: - self.load_cluster_from_etcd() + self.load_cluster_from_dcs() if self.state_handler.is_leader(): self.demote() return 'demoting self because i am not the healthiest node' @@ -84,7 +84,7 @@ class Ha: else: self.follow_the_leader() return 'no action. i am a secondary and i am following a leader' - except DCSError as e: + except DCSError: logger.error('Error communicating with DCS') if self.state_handler.is_leader(): self.state_handler.demote(None) diff --git a/tests/test_governor.py b/tests/test_governor.py index 0a8d4ab8..daced98b 100644 --- a/tests/test_governor.py +++ b/tests/test_governor.py @@ -83,11 +83,11 @@ class TestGovernor(unittest.TestCase): self.g.touch_member() def test_governor_initialize(self): - self.g.etcd.client._base_uri = 'http://remote' + self.g.ha.dcs.client._base_uri = 'http://remote' self.g.postgresql.data_directory_empty = true - self.g.etcd.race = true + self.g.ha.dcs.race = true self.g.initialize() - self.g.etcd.race = false + self.g.ha.dcs.race = false self.g.initialize() self.g.postgresql.data_directory_empty = false self.g.touch_member = self.touch_member diff --git a/tests/test_ha.py b/tests/test_ha.py index e45885e8..299bc4b6 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -73,9 +73,9 @@ class TestHa(unittest.TestCase): self.p = MockPostgresql() self.e = Etcd({'ttl': 30, 'host': 'remotehost:2379', 'scope': 'test'}) self.ha = Ha(self.p, self.e) - self.ha.load_cluster_from_etcd() + self.ha.load_cluster_from_dcs() self.ha.cluster = Cluster(False, None, None, []) - self.ha.load_cluster_from_etcd = nop + self.ha.load_cluster_from_dcs = nop def test_start_as_slave(self): self.p.is_healthy = false @@ -133,5 +133,5 @@ class TestHa(unittest.TestCase): self.assertEquals(self.ha.run_cycle(), 'no action. i am a secondary and i am following a leader') def test_no_etcd_connection_master_demote(self): - self.ha.load_cluster_from_etcd = dead_etcd + self.ha.load_cluster_from_dcs = dead_etcd self.assertEquals(self.ha.run_cycle(), 'demoted self because DCS is not accessible and i was a leader')