diff --git a/patroni/dcs/etcd3.py b/patroni/dcs/etcd3.py index 0bfc4fae..e28da844 100644 --- a/patroni/dcs/etcd3.py +++ b/patroni/dcs/etcd3.py @@ -630,6 +630,16 @@ class PatroniEtcd3Client(Etcd3Client): return ret + def txn(self, compare: Dict[str, Any], success: Dict[str, Any], + failure: Optional[Dict[str, Any]] = None, retry: Optional[Retry] = None) -> Dict[str, Any]: + ret = super(PatroniEtcd3Client, self).txn(compare, success, failure, retry) + # Here we abuse the fact that the `failure` is only set in the call from update_leader(). + # In all other cases the txn() call failure may be an indicator of a stale cache, + # and therefore we want to restart watcher. + if not failure and not ret: + self._restart_watcher() + return ret + class Etcd3(AbstractEtcd): diff --git a/tests/test_etcd3.py b/tests/test_etcd3.py index ace59a62..2e3ed59d 100644 --- a/tests/test_etcd3.py +++ b/tests/test_etcd3.py @@ -127,10 +127,16 @@ class TestPatroniEtcd3Client(BaseTestEtcd3): request = {'key': base64_encode('/patroni/test/leader')} mock_urlopen.return_value = MockResponse() mock_urlopen.return_value.content = '{"succeeded":true,"header":{"revision":"1"}}' - self.client.call_rpc('/kv/txn', {'success': [{'request_delete_range': request}]}) self.client.call_rpc('/kv/put', request) self.client.call_rpc('/kv/deleterange', request) + @patch.object(urllib3.PoolManager, 'urlopen') + def test_txn(self, mock_urlopen): + mock_urlopen.return_value = MockResponse() + mock_urlopen.return_value.content = '{"header":{"revision":"1"}}' + self.client.txn({'target': 'MOD', 'mod_revision': '1'}, + {'request_delete_range': {'key': base64_encode('/patroni/test/leader')}}) + @patch('time.time', Mock(side_effect=[1, 10.9, 100])) def test__wait_cache(self): with self.kv_cache.condition: