From 1c7bf2f59e755caccd1c8ad00849ba37dc992855 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Thu, 25 May 2023 10:01:43 +0200 Subject: [PATCH] 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 --- patroni/dcs/etcd3.py | 47 ++++++++++++++++++++++++++++++-------------- tests/test_etcd3.py | 5 ++++- 2 files changed, 36 insertions(+), 16 deletions(-) diff --git a/patroni/dcs/etcd3.py b/patroni/dcs/etcd3.py index 55923df3..6c4312e5 100644 --- a/patroni/dcs/etcd3.py +++ b/patroni/dcs/etcd3.py @@ -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: diff --git a/tests/test_etcd3.py b/tests/test_etcd3.py index 6f56d0e6..4b3a78ab 100644 --- a/tests/test_etcd3.py +++ b/tests/test_etcd3.py @@ -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'))