mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Rename ha.etcd into ha.dcs
This commit is contained in:
+12
-7
@@ -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__':
|
||||
|
||||
@@ -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
|
||||
+10
-10
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
+3
-3
@@ -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')
|
||||
|
||||
Reference in New Issue
Block a user