diff --git a/patroni/ha.py b/patroni/ha.py index 59c57565..95ca6e05 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -497,6 +497,7 @@ class Ha(object): role = 'replica' if self.has_lock() and not self.is_standby_cluster(): + self._rewind.reset_state() # we want to later trigger CHECKPOINT after promote msg = "starting as readonly because i had the session lock" node_to_follow = None else: @@ -769,18 +770,14 @@ class Ha(object): self.state_handler.sync_handler.set_synchronous_standby_names( CaseInsensitiveSet('*') if self.global_config.is_synchronous_mode_strict else CaseInsensitiveSet()) if self.state_handler.role not in ('master', 'promoted', 'primary'): - def on_success(): - self._rewind.reset_state() - logger.info("cleared rewind state after becoming the leader") - def before_promote(): self.notify_citus_coordinator('before_promote') with self._async_response: self._async_response.reset() + self._async_executor.try_run_async('promote', self.state_handler.promote, - args=(self.dcs.loop_wait, self._async_response, - before_promote, on_success)) + args=(self.dcs.loop_wait, self._async_response, before_promote)) return promote_message def fetch_node_status(self, member: Member) -> _MemberStatus: @@ -1614,6 +1611,9 @@ class Ha(object): else: if self._was_paused: self.state_handler.schedule_sanity_checks_after_pause() + # during pause people could manually do something with Postgres, therefore we want + # to double check rewind conditions on replicas and maybe run CHECKPOINT on the primary + self._rewind.reset_state() self._was_paused = False if not self.cluster.has_member(self.state_handler.name): diff --git a/patroni/postgresql/__init__.py b/patroni/postgresql/__init__.py index 632cbee7..04fab0a2 100644 --- a/patroni/postgresql/__init__.py +++ b/patroni/postgresql/__init__.py @@ -1125,8 +1125,8 @@ class Postgresql(object): except Exception as e: logger.error('Exception when calling `%s`: %r', cmd, e) - def promote(self, wait_seconds: int, task: CriticalTask, before_promote: Optional[Callable[..., Any]] = None, - on_success: Optional[Callable[..., Any]] = None) -> Optional[bool]: + def promote(self, wait_seconds: int, task: CriticalTask, + before_promote: Optional[Callable[..., Any]] = None) -> Optional[bool]: if self.role in ('promoted', 'master', 'primary'): return True @@ -1152,8 +1152,6 @@ class Postgresql(object): ret = self.pg_ctl('promote', '-W') if ret: self.set_role('promoted') - if on_success is not None: - on_success() self.call_nowait(CallbackAction.ON_ROLE_CHANGE) ret = self._wait_promote(wait_seconds) return ret diff --git a/patroni/postgresql/rewind.py b/patroni/postgresql/rewind.py index ff6fc751..270d629c 100644 --- a/patroni/postgresql/rewind.py +++ b/patroni/postgresql/rewind.py @@ -280,7 +280,7 @@ class Rewind(object): """After promote issue a CHECKPOINT from a new thread and asynchronously check the result. In case if CHECKPOINT failed, just check that timeline in pg_control was updated.""" - if self._state == REWIND_STATUS.INITIAL and self._postgresql.is_leader(): + if self._state != REWIND_STATUS.CHECKPOINT and self._postgresql.is_leader(): with self._checkpoint_task_lock: if self._checkpoint_task: with self._checkpoint_task: