mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Introduce Retry.ensure_deadline() method (#2694)
it helps to get rid of recurring patterns:
```python
retry.deadline = retry.stoptime - time.time()
if retry.deadline < XXX:
raise Exception(...) or return False
```
This commit is contained in:
+6
-15
@@ -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:
|
||||
|
||||
+4
-10
@@ -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
|
||||
|
||||
@@ -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 {}
|
||||
|
||||
@@ -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.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user