Fix a problem with etcd3.update_leader() (#2693)

It didn't took into account the fact that we can get a new lease after changing a TTL. In this case we have to update the exiting leader key with the new lease.
To solve it we introduce the on transaction 'failure' callback. The whole workflow looks like (schematically):
```python
txn(
    compare=(old_value == self._name)),
    success=put(self.leader_path, self._name, self._lease),
    failure=txn(
        compare=(create_revision == '0'),
        success=put(self.leader_path, self._name, self._lease)
    )
)
```

The problem was introduced in d98d6d0b02c9cc67464a7b1b31b1c5570c26e12d
This commit is contained in:
Alexander Kukushkin
2023-05-25 10:01:43 +02:00
committed by GitHub
parent 822b6ec711
commit 1c7bf2f59e
2 changed files with 36 additions and 16 deletions
+32 -15
View File
@@ -16,7 +16,7 @@ from threading import Condition, Lock, Thread
from typing import Any, Callable, Collection, Dict, Iterator, List, Optional, Tuple, Type, TYPE_CHECKING, Union
from . import ClusterConfig, Cluster, Failover, Leader, Member, SyncState,\
TimelineHistory, ReturnFalseException, catch_return_false_exception, citus_group_re
TimelineHistory, catch_return_false_exception, citus_group_re
from .etcd import AbstractEtcdClientWithFailover, AbstractEtcd, catch_etcd_errors, DnsCachingResolver, Retry
from ..exceptions import DCSError, PatroniException
from ..utils import deep_compare, enable_keepalive, iter_response_objects, RetryFailedError, USER_AGENT
@@ -343,9 +343,13 @@ class Etcd3Client(AbstractEtcdClientWithFailover):
def lease_keepalive(self, ID: str, retry: Optional[Retry] = None) -> Optional[str]:
return self.call_rpc('/lease/keepalive', {'ID': ID}, retry).get('result', {}).get('TTL')
def txn(self, compare: Dict[str, Any], success: Dict[str, Any], retry: Optional[Retry] = None) -> Dict[str, Any]:
ret = self.call_rpc('/kv/txn', {'compare': [compare], 'success': [success]}, retry)
return ret if ret.get('succeeded') else {}
def txn(self, compare: Dict[str, Any], success: Dict[str, Any],
failure: Optional[Dict[str, Any]] = None, retry: Optional[Retry] = None) -> Dict[str, Any]:
fields = {'compare': [compare], 'success': [success]}
if failure:
fields['failure'] = [failure]
ret = self.call_rpc('/kv/txn', fields, retry)
return ret if failure or ret.get('succeeded') else {}
@_handle_auth_errors
def put(self, key: str, value: str, lease: Optional[str] = None, create_revision: Optional[str] = None,
@@ -360,7 +364,7 @@ class Etcd3Client(AbstractEtcdClientWithFailover):
else:
return self.call_rpc('/kv/put', fields, retry)
compare['key'] = fields['key']
return self.txn(compare, {'request_put': fields}, retry)
return self.txn(compare, {'request_put': fields}, retry=retry)
@_handle_auth_errors
def deleterange(self, key: str, range_end: Union[bytes, str, None] = None,
@@ -369,7 +373,7 @@ class Etcd3Client(AbstractEtcdClientWithFailover):
if mod_revision is None:
return self.call_rpc('/kv/deleterange', fields, retry)
compare = {'target': 'MOD', 'mod_revision': mod_revision, 'key': fields['key']}
return self.txn(compare, {'request_delete_range': fields}, retry)
return self.txn(compare, {'request_delete_range': fields}, retry=retry)
def deleteprefix(self, key: str, retry: Optional[Retry] = None) -> Dict[str, Any]:
return self.deleterange(key, prefix_range_end(key), retry=retry)
@@ -603,7 +607,13 @@ class PatroniEtcd3Client(Etcd3Client):
if self._kv_cache:
value = delete = None
if method == '/kv/txn' and ret.get('succeeded'):
# For the 'failure' case we only support a second (nested) transaction that attempts to
# update/delete the same keys. Anything more complex than that we don't need and therefore it doesn't
# make sense to write a universal response analyzer and we can just check expected JSON path.
if method == '/kv/txn'\
and (ret.get('succeeded') or 'failure' in fields and 'request_txn' in fields['failure'][0]
and ret.get('responses', [{'response_txn': {'succeeded': False}}])[0]
.get('response_txn', {}).get('succeeded')):
on_success = fields['success'][0]
value = on_success.get('request_put')
delete = on_success.get('request_delete_range')
@@ -813,7 +823,7 @@ class Etcd3(AbstractEtcd):
return retry(*args, **kwargs)
try:
return _retry(self._client.put, self.leader_path, self._name, self._lease, '0')
return _retry(self._client.put, self.leader_path, self._name, self._lease, create_revision='0')
except LeaseNotFound:
logger.error('Our lease disappeared from Etcd. Will try to get a new one and retry attempt')
self._lease = None
@@ -825,7 +835,7 @@ class Etcd3(AbstractEtcd):
if retry.deadline < 1:
raise Etcd3Error('_do_attempt_to_acquire_leader timeout')
return _retry(self._client.put, self.leader_path, self._name, self._lease, '0')
return _retry(self._client.put, self.leader_path, self._name, self._lease, create_revision='0')
@catch_return_false_exception
def attempt_to_acquire_leader(self) -> bool:
@@ -884,16 +894,23 @@ class Etcd3(AbstractEtcd):
if retry.deadline < 1:
raise Etcd3Error('update_leader timeout')
try:
self._run_and_handle_exceptions(self._client.put, self.leader_path,
self._name, self._lease, '0', retry=_retry)
except ReturnFalseException:
pass
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
compare1 = {'key': fields['key'], 'target': 'VALUE', 'value': fields['value']}
request_put = {'request_put': fields}
# If the first comparison failed we will try to create the new leader key in a transaction
compare2 = {'key': fields['key'], 'target': 'CREATE', 'create_revision': '0'}
request_txn = {'request_txn': {'compare': [compare2], 'success': [request_put]}}
ret = self._run_and_handle_exceptions(self._client.txn, compare1,
request_put, request_txn, retry=_retry)
return ret.get('succeeded', False)\
or ret.get('responses', [{}])[0].get('response_txn', {}).get('succeeded', False)
return bool(self._lease)
@catch_etcd_errors
def initialize(self, create_new: bool = True, sysid: str = ""):
return self.retry(self._client.put, self.initialize_path, sysid, None, '0' if create_new else None)
return self.retry(self._client.put, self.initialize_path, sysid, create_revision='0' if create_new else None)
@catch_etcd_errors
def _delete_leader(self) -> bool:
+4 -1
View File
@@ -236,12 +236,15 @@ class TestEtcd3(BaseTestEtcd3):
def test__update_leader(self):
self.etcd3._lease = None
self.etcd3.update_leader('123', failsafe={'foo': 'bar'})
with patch.object(Etcd3Client, 'txn', Mock(return_value={'succeeded': True})):
self.etcd3.update_leader('123', failsafe={'foo': 'bar'})
self.etcd3._last_lease_refresh = 0
self.etcd3.update_leader('124')
with patch.object(PatroniEtcd3Client, 'lease_keepalive', Mock(return_value=True)),\
patch('time.time', Mock(side_effect=[0, 100, 200, 300])):
self.assertRaises(Etcd3Error, self.etcd3.update_leader, '126')
self.etcd3._lease = self.etcd3.cluster.leader.session
self.etcd3.update_leader('124')
self.etcd3._last_lease_refresh = 0
with patch.object(PatroniEtcd3Client, 'lease_keepalive', Mock(side_effect=Unknown)):
self.assertFalse(self.etcd3.update_leader('125'))