diff --git a/patroni/ha.py b/patroni/ha.py index 6923a126..f41ce505 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -106,12 +106,6 @@ class Ha(object): return 'waiting for leader to bootstrap' def recover(self): - # try to see if we are the former master that crashed. If so - we likely need to run pg_rewind - # in order to join the former standby being promoted. - if self.state_handler.role == 'master': - pg_controldata = self.state_handler.controldata() - if pg_controldata and pg_controldata.get('Database cluster state', '') == 'in production': # crashed master - self.state_handler.require_rewind() self.recovering = True return self.follow("starting as readonly because i had the session lock", "starting as a secondary", True, True) diff --git a/patroni/postgresql.py b/patroni/postgresql.py index cf2fddfa..a5b1e16c 100644 --- a/patroni/postgresql.py +++ b/patroni/postgresql.py @@ -73,14 +73,13 @@ class Postgresql(object): self._connection = None self._cursor_holder = None - self._need_rewind = False self._sysid = None self.replication_slots = [] # list of already existing replication slots self.retry = Retry(max_tries=-1, deadline=5, max_delay=1, retry_exceptions=PostgresConnectionException) self._state = 'stopped' self._state_lock = Lock() - self._role = 'replica' + self._role = self.get_postgres_role_from_data_directory() self._role_lock = Lock() if self.is_running(): @@ -115,9 +114,6 @@ class Postgresql(object): self._sysid = data.get('Database system identifier', "") return self._sysid - def require_rewind(self): - self._need_rewind = True - def get_local_address(self): listen_addresses = self.listen_addresses.split(',') local_address = listen_addresses[0].strip() # take first address from listen_addresses @@ -128,6 +124,9 @@ class Postgresql(object): break return local_address + ':' + self.port + def get_postgres_role_from_data_directory(self): + return 'replica' if os.path.exists(self.recovery_conf) else 'master' + @property def _connect_kwargs(self): r = parseurl('postgres://{0}/postgres'.format(self.local_address)) @@ -351,7 +350,7 @@ class Postgresql(object): logger.error('Cannot start PostgreSQL because one is already running.') return True - self.set_role('replica' if os.path.exists(self.recovery_conf) else 'master') + self.set_role(self.get_postgres_role_from_data_directory()) if os.path.exists(self.postmaster_pid): os.remove(self.postmaster_pid) logger.info('Removed %s', self.postmaster_pid) @@ -564,13 +563,12 @@ recovery_target_timeline = 'latest' def follow(self, leader, recovery=False): if self.check_recovery_conf(leader) and not recovery: return True - change_role = self.role == 'master' - self._need_rewind = (self._need_rewind or change_role) and self.can_rewind - if self._need_rewind: + need_rewind = change_role and self.can_rewind + if need_rewind: logger.info("set the rewind flag after demote") self.write_recovery_conf(leader) - if leader and self._need_rewind: # we have a leader and need to rewind + if leader and need_rewind: # we have a leader and need to rewind if self.is_running(): self.stop() # at present, pg_rewind only runs when the cluster is shut down cleanly @@ -594,7 +592,6 @@ recovery_target_timeline = 'latest' logger.error("unable to rewind the former master") self.remove_data_directory() ret = True - self._need_rewind = False else: # do not rewind until the leader becomes available ret = self.restart() if change_role and ret: @@ -630,7 +627,6 @@ recovery_target_timeline = 'latest' if ret: self.set_role('master') logger.info("cleared rewind flag after becoming the leader") - self._need_rewind = False self.call_nowait(ACTION_ON_ROLE_CHANGE) return ret diff --git a/tests/test_ha.py b/tests/test_ha.py index 64b4d7c5..3071ecc8 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -91,6 +91,7 @@ class TestHa(unittest.TestCase): 'data_dir': 'data/postgresql0', 'superuser': {}, 'admin': {}, 'replication': {'username': '', 'password': '', 'network': ''}}) self.p.set_state('running') + self.p.set_role('replica') self.p.check_replication_lag = true self.p.can_create_replica_without_replication_connection = MagicMock(return_value=False) self.e = Etcd('foo', {'ttl': 30, 'host': 'ok:2379', 'scope': 'test'}) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 932e860d..8bd563b9 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -257,12 +257,10 @@ class TestPostgresql(unittest.TestCase): self.p.follow(Leader(-1, 28, self.other)) self.p.rewind = mock_pg_rewind self.p.follow(self.leader) - self.p.require_rewind() with mock.patch('os.path.islink', MagicMock(return_value=True)): with mock.patch('patroni.postgresql.Postgresql.can_rewind', new_callable=PropertyMock(return_value=True)): with mock.patch('os.unlink', MagicMock(return_value=True)): self.p.follow(self.leader, recovery=True) - self.p.require_rewind() with mock.patch('patroni.postgresql.Postgresql.can_rewind', new_callable=PropertyMock(return_value=True)): self.p.rewind.return_value = True self.p.follow(self.leader, recovery=True)