diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index a0a08978..abd98a54 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -905,13 +905,23 @@ class Kubernetes(AbstractDCS): except (RetryFailedError, K8sException): return False - deadline = retry.stoptime - time.time() - if deadline < 2: + retry.deadline = retry.stoptime - time.time() + if retry.deadline < 1: return False - retry.sleep_func(1) # Give a chance for ObjectCache to receive the latest version + # Try to get the latest version directly from K8s API instead of relying on async cache + try: + kind = retry(self._api.read_namespaced_kind, self.leader_path, self._namespace) + except Exception as e: + logger.error('Failed to get the leader object "%s": %r', self.leader_path, e) + return False + + self._kinds.set(self.leader_path, kind) + + retry.deadline = retry.stoptime - time.time() + if retry.deadline < 0.5: + return False - kind = self._kinds.get(self.leader_path) kind_annotations = kind and kind.metadata.annotations or {} kind_resource_version = kind and kind.metadata.resource_version @@ -919,7 +929,6 @@ class Kubernetes(AbstractDCS): if kind and (kind_annotations.get(self._LEADER) != self._name or kind_resource_version == resource_version): return False - retry.deadline = deadline - 1 # Update deadline and retry return self.patch_or_create(self.leader_path, annotations, kind_resource_version, ips=ips, retry=_retry) def update_leader(self, last_operation, access_is_restricted=False): diff --git a/tests/test_kubernetes.py b/tests/test_kubernetes.py index cb9462e6..aa4493a9 100644 --- a/tests/test_kubernetes.py +++ b/tests/test_kubernetes.py @@ -26,7 +26,7 @@ def mock_list_namespaced_config_map(*args, **kwargs): return k8s_client.V1ConfigMapList(metadata=metadata, items=items, kind='ConfigMapList') -def mock_list_namespaced_endpoints(*args, **kwargs): +def mock_read_namespaced_endpoints(*args, **kwargs): target_ref = k8s_client.V1ObjectReference(kind='Pod', resource_version='10', name='p-0', namespace='default', uid='964dfeae-e79b-4476-8a5a-1920b5c2a69d') address0 = k8s_client.V1EndpointAddress(ip='10.0.0.0', target_ref=target_ref) @@ -35,9 +35,12 @@ def mock_list_namespaced_endpoints(*args, **kwargs): subset = k8s_client.V1EndpointSubset(addresses=[address1, address0], ports=[port]) metadata = k8s_client.V1ObjectMeta(resource_version='1', labels={'f': 'b'}, name='test', annotations={'optime': '1234', 'leader': 'p-0', 'ttl': '30s'}) - endpoint = k8s_client.V1Endpoints(subsets=[subset], metadata=metadata) - metadata = k8s_client.V1ObjectMeta(resource_version='1') - return k8s_client.V1EndpointsList(metadata=metadata, items=[endpoint], kind='V1EndpointsList') + return k8s_client.V1Endpoints(subsets=[subset], metadata=metadata) + + +def mock_list_namespaced_endpoints(*args, **kwargs): + return k8s_client.V1EndpointsList(metadata=k8s_client.V1ObjectMeta(resource_version='1'), + items=[mock_read_namespaced_endpoints()], kind='V1EndpointsList') def mock_list_namespaced_pod(*args, **kwargs): @@ -261,20 +264,25 @@ class TestKubernetesEndpoints(BaseTestKubernetes): def test_update_leader_with_restricted_access(self): self.assertIsNotNone(self.k.update_leader('123', True)) + @patch.object(k8s_client.CoreV1Api, 'read_namespaced_endpoints', create=True) @patch.object(k8s_client.CoreV1Api, 'patch_namespaced_endpoints', create=True) - def test__update_leader_with_retry(self, mock_patch): + def test__update_leader_with_retry(self, mock_patch, mock_read): + mock_read.return_value = mock_read_namespaced_endpoints() mock_patch.side_effect = k8s_client.rest.ApiException(502, '') self.assertFalse(self.k.update_leader('123')) mock_patch.side_effect = RetryFailedError('') self.assertFalse(self.k.update_leader('123')) mock_patch.side_effect = k8s_client.rest.ApiException(409, '') - with patch('time.time', Mock(side_effect=[0, 100, 200])): + with patch('time.time', Mock(side_effect=[0, 100, 200, 0, 0, 0, 0, 100, 200])): self.assertFalse(self.k.update_leader('123')) - with patch('time.sleep', Mock()): self.assertFalse(self.k.update_leader('123')) - mock_patch.side_effect = [k8s_client.rest.ApiException(409, ''), mock_namespaced_kind()] - self.k._kinds._object_cache['test'].metadata.resource_version = '2' - self.assertIsNotNone(self.k._update_leader_with_retry({}, '1', [])) + self.assertFalse(self.k.update_leader('123')) + mock_patch.side_effect = [k8s_client.rest.ApiException(409, ''), mock_namespaced_kind()] + mock_read.return_value.metadata.resource_version = '2' + self.assertIsNotNone(self.k._update_leader_with_retry({}, '1', [])) + mock_patch.side_effect = k8s_client.rest.ApiException(409, '') + mock_read.side_effect = Exception + self.assertFalse(self.k.update_leader('123')) @patch.object(k8s_client.CoreV1Api, 'create_namespaced_endpoints', Mock(side_effect=[k8s_client.rest.ApiException(500, ''),