diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index 87ae4f6a..5327be0f 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -549,13 +549,11 @@ class Consul(AbstractDCS): except InvalidSession: logger.error('Our session disappeared from Consul. Will try to get a new one and retry attempt') self._session = None - retry.deadline = retry.stoptime - time.time() + retry.ensure_deadline(0) retry(self._do_refresh_session) - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: - raise ConsulError('_do_attempt_to_acquire_leader timeout') + retry.ensure_deadline(1, ConsulError('_do_attempt_to_acquire_leader timeout')) return retry(self._client.kv.put, self.leader_path, self._name, acquire=self._session) @@ -564,9 +562,7 @@ class Consul(AbstractDCS): retry = self._retry.copy() self._run_and_handle_exceptions(self._do_refresh_session, retry=retry) - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: - raise ConsulError('attempt_to_acquire_leader timeout') + retry.ensure_deadline(1, ConsulError('attempt_to_acquire_leader timeout')) ret = self._run_and_handle_exceptions(self._do_attempt_to_acquire_leader, retry, retry=None) if not ret: @@ -614,16 +610,12 @@ class Consul(AbstractDCS): self._run_and_handle_exceptions(self._do_refresh_session, True, retry=retry) if self._session and leader.session != self._session: - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: - raise ConsulError('update_leader timeout') + retry.ensure_deadline(1, ConsulError('update_leader timeout')) logger.warning('Recreating the leader key due to session mismatch') self._run_and_handle_exceptions(self._client.kv.delete, self.leader_path, cas=leader.version) - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 0.5: - raise ConsulError('update_leader timeout') + retry.ensure_deadline(0.5, ConsulError('update_leader timeout')) self._run_and_handle_exceptions(self._client.kv.put, self.leader_path, self._name, acquire=self._session) @@ -659,8 +651,7 @@ class Consul(AbstractDCS): retry = self._retry.copy() ret = retry(self._client.kv.put, self.sync_path, value, cas=version) if ret: # We have no other choise, only read after write :( - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 0.5: + if not retry.ensure_deadline(0.5): return False _, ret = self.retry(self._client.kv.get, self.sync_path) if ret and (ret.get('Value') or b'').decode('utf-8') == value: diff --git a/patroni/dcs/etcd3.py b/patroni/dcs/etcd3.py index 4d0ff75d..b89e36a5 100644 --- a/patroni/dcs/etcd3.py +++ b/patroni/dcs/etcd3.py @@ -827,13 +827,11 @@ class Etcd3(AbstractEtcd): except LeaseNotFound: logger.error('Our lease disappeared from Etcd. Will try to get a new one and retry attempt') self._lease = None - retry.deadline = retry.stoptime - time.time() + retry.ensure_deadline(0) _retry(self._do_refresh_lease) - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: - raise Etcd3Error('_do_attempt_to_acquire_leader timeout') + retry.ensure_deadline(1, Etcd3Error('_do_attempt_to_acquire_leader timeout')) return _retry(self._client.put, self.leader_path, self._name, self._lease, create_revision='0') @@ -847,9 +845,7 @@ class Etcd3(AbstractEtcd): self._run_and_handle_exceptions(self._do_refresh_lease, retry=_retry) - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: - raise Etcd3Error('attempt_to_acquire_leader timeout') + retry.ensure_deadline(1, Etcd3Error('attempt_to_acquire_leader timeout')) ret = self._run_and_handle_exceptions(self._do_attempt_to_acquire_leader, retry, retry=None) if not ret: @@ -887,9 +883,7 @@ class Etcd3(AbstractEtcd): self._run_and_handle_exceptions(self._do_refresh_lease, True, retry=_retry) if self._lease and leader.session != self._lease: - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: - raise Etcd3Error('update_leader timeout') + retry.ensure_deadline(1, Etcd3Error('update_leader timeout')) fields = {'key': base64_encode(self.leader_path), 'value': base64_encode(self._name), 'lease': self._lease} # First we try to update lease on existing leader key "hoping" that we still owning it diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index ac65652d..0d457ba4 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -1153,8 +1153,7 @@ class Kubernetes(AbstractDCS): raise KubernetesError(e) # if we are here, that means update failed with 409 - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 1: + if not retry.ensure_deadline(1): return False # No time for retry. Tell ha.py that we have to demote due to failed update. # Try to get the latest version directly from K8s API instead of relying on async cache @@ -1168,8 +1167,7 @@ class Kubernetes(AbstractDCS): self._kinds.set(self.leader_path, kind) - retry.deadline = retry.stoptime - time.time() - if retry.deadline < 0.5: + if not retry.ensure_deadline(0.5): return False kind_annotations = kind and kind.metadata.annotations or {} diff --git a/patroni/utils.py b/patroni/utils.py index f6cfa919..eb02c561 100644 --- a/patroni/utils.py +++ b/patroni/utils.py @@ -531,6 +531,21 @@ class Retry(object): """Get the current stop time.""" return self._cur_stoptime or 0 + def ensure_deadline(self, timeout: float, raise_ex: Optional[Exception] = None) -> bool: + """Calculates, sets, and checks the remaining deadline time. + + :param timeout: if the *deadline* is smaller than the provided *timeout* value raise *raise_ex* exception + :param raise_ex: the exception object that will be raised if the *deadline* is smaller than provided *timeout* + :returns: `False` if *deadline* is smaller than a provided *timeout* and *raise_ex* isn't set. Otherwise `True` + :raises Exception: if calculated deadline is smaller than provided *timeout* + """ + self.deadline = self.stoptime - time.time() + if self.deadline < timeout: + if raise_ex: + raise raise_ex + return False + return True + def __call__(self, func: Callable[..., Any], *args: Any, **kwargs: Any) -> Any: """Call a function *func* with arguments ``*args`` and ``*kwargs`` in a loop.