diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index c6ba3af9..6fcf4984 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -1215,9 +1215,6 @@ class Kubernetes(AbstractDCS): self._kinds.set(self.leader_path, kind) - if not retry.ensure_deadline(0.5): - return False - kind_annotations = kind and kind.metadata.annotations or EMPTY_DICT kind_resource_version = kind and kind.metadata.resource_version @@ -1225,6 +1222,14 @@ class Kubernetes(AbstractDCS): if kind and (kind_annotations.get(self._LEADER) != self._name or kind_resource_version == resource_version): return False + # We can get 409 because we do at least one retry, and the first update might have succeeded, + # therefore we will check if annotations on the read object match expectations. + if all(kind_annotations.get(k) == v for k, v in annotations.items()): + return True + + if not retry.ensure_deadline(0.5): + return False + return bool(_run_and_handle_exceptions(self._patch_or_create, self.leader_path, annotations, kind_resource_version, ips=ips, retry=_retry)) diff --git a/tests/test_kubernetes.py b/tests/test_kubernetes.py index bf684ef4..dc4a1365 100644 --- a/tests/test_kubernetes.py +++ b/tests/test_kubernetes.py @@ -422,14 +422,18 @@ class TestKubernetesEndpoints(BaseTestKubernetes): self.assertFalse(self.k.update_leader(cluster, '123')) mock_patch.side_effect = RetryFailedError('') self.assertRaises(KubernetesError, self.k.update_leader, cluster, '123') - mock_patch.side_effect = k8s_client.rest.ApiException(409, '') + mock_patch.side_effect = [k8s_client.rest.ApiException(409, ''), + k8s_client.rest.ApiException(409, ''), mock_namespaced_kind()] + mock_read.return_value.metadata.resource_version = '2' with patch('time.time', Mock(side_effect=[0, 0, 100, 200, 0, 0, 0, 0, 0, 100, 200])): self.assertFalse(self.k.update_leader(cluster, '123')) self.assertFalse(self.k.update_leader(cluster, '123')) + mock_patch.side_effect = k8s_client.rest.ApiException(409, '') self.assertFalse(self.k.update_leader(cluster, '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', [])) + self.assertTrue(self.k._update_leader_with_retry({}, '1', [])) + mock_patch.side_effect = [k8s_client.rest.ApiException(409, ''), mock_namespaced_kind()] + self.assertIsNotNone(self.k._update_leader_with_retry({'foo': 'bar'}, '1', [])) mock_patch.side_effect = k8s_client.rest.ApiException(409, '') mock_read.side_effect = RetryFailedError('') self.assertRaises(KubernetesError, self.k.update_leader, cluster, '123')