diff --git a/patroni/ha.py b/patroni/ha.py index 9e0f9b84..d5f47bd6 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -293,12 +293,25 @@ class Ha(object): return result + def _handle_crash_recovery(self): + if not self._crash_recovery_executed and (self.cluster.is_unlocked() or self._rewind.can_rewind): + self._crash_recovery_executed = True + self._crash_recovery_started = time.time() + msg = 'doing crash recovery in a single user mode' + return self._async_executor.try_run_async(msg, self._rewind.ensure_clean_shutdown) or msg + def _handle_rewind_or_reinitialize(self): leader = self.get_remote_master() if self.is_standby_cluster() else self.cluster.leader if not self._rewind.rewind_or_reinitialize_needed_and_possible(leader): return None if self._rewind.can_rewind: + # rewind is required, but postgres wasn't shut down cleanly. + if self.state_handler.controldata().get('Database cluster state') == 'in archive recovery': + msg = self._handle_crash_recovery() + if msg: + return msg + msg = 'running pg_rewind from ' + leader.name return self._async_executor.try_run_async(msg, self._rewind.execute, args=(leader,)) or msg @@ -325,13 +338,10 @@ class Ha(object): data = self.state_handler.controldata() logger.info('pg_controldata:\n%s\n', '\n'.join(' {0}: {1}'.format(k, v) for k, v in data.items())) - if data.get('Database cluster state') in ('in production', 'shutting down', 'in crash recovery') \ - and not self._crash_recovery_executed and \ - (self.cluster.is_unlocked() or self._rewind.can_rewind): - self._crash_recovery_executed = True - self._crash_recovery_started = time.time() - msg = 'doing crash recovery in a single user mode' - return self._async_executor.try_run_async(msg, self._rewind.ensure_clean_shutdown) or msg + if data.get('Database cluster state') in ('in production', 'shutting down', 'in crash recovery'): + msg = self._handle_crash_recovery() + if msg: + return msg self.load_cluster_from_dcs() diff --git a/patroni/postgresql/rewind.py b/patroni/postgresql/rewind.py index ba6e4532..fb754f58 100644 --- a/patroni/postgresql/rewind.py +++ b/patroni/postgresql/rewind.py @@ -112,7 +112,7 @@ class Rewind(object): in_recovery = timeline = lsn = None data = self._postgresql.controldata() try: - if data.get('Database cluster state') == 'shut down in recovery': + if data.get('Database cluster state') in ('shut down in recovery', 'in archive recovery'): in_recovery = True lsn = data.get('Minimum recovery ending location') timeline = int(data.get("Min recovery ending loc's timeline")) diff --git a/tests/test_ha.py b/tests/test_ha.py index 90467b20..e2ef7b6f 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -284,6 +284,18 @@ class TestHa(PostgresInit): self.ha.patroni.config.set_dynamic_configuration({'maximum_lag_on_failover': 10}) self.assertEqual(self.ha.run_cycle(), 'terminated crash recovery because of startup timeout') + @patch.object(Rewind, 'ensure_clean_shutdown', Mock()) + @patch.object(Rewind, 'rewind_or_reinitialize_needed_and_possible', Mock(return_value=True)) + @patch.object(Rewind, 'can_rewind', PropertyMock(return_value=True)) + def test_crash_recovery_before_rewind(self): + self.p.is_leader = false + self.p.is_running = false + self.p.controldata = lambda: {'Database cluster state': 'in archive recovery', + 'Database system identifier': SYSID} + self.ha._rewind.trigger_check_diverged_lsn() + self.ha.cluster = get_cluster_initialized_with_leader() + self.assertEqual(self.ha.run_cycle(), 'doing crash recovery in a single user mode') + @patch.object(Rewind, 'rewind_or_reinitialize_needed_and_possible', Mock(return_value=True)) @patch.object(Rewind, 'can_rewind', PropertyMock(return_value=True)) def test_recover_with_rewind(self):