mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-25 14:53:37 +00:00
Fix possible race conditions in update_leader (#1596)
1. Between get_cluster() and update_leader() calls the K8s leader object might be updated from outside and therefore the resource version will not match (error code=409). Since we are watching for all changes, the ObjectCache likely will have the most up-to-date version and we will take advantage of that. There is still a chance to hit a race-condition, but it would be smaller than before. Actually, other DCS are free of this issue. Etcd - update is based on the value comparison, Zookeeper and Consul are relying on session mechanism. 2. If the update still failed - recheck the resource version of the leader object and that the current node is still the leader there and repeat the call. P.S. The leader race is still relying on the version of the leader object as it was during the get_cluster() call. In addition to that fixed handling of K8s API errors, we should retry on 500, not on 502. Close https://github.com/zalando/patroni/issues/1589
This commit is contained in:
+80
-34
@@ -75,7 +75,7 @@ class CoreV1ApiProxy(object):
|
||||
try:
|
||||
return getattr(self._api, func)(*args, **kwargs)
|
||||
except k8s_client.rest.ApiException as e:
|
||||
if e.status in (502, 503, 504) or e.headers and 'retry-after' in e.headers: # XXX
|
||||
if e.status in (500, 503, 504) or e.headers and 'retry-after' in e.headers: # XXX
|
||||
raise KubernetesRetriableException(e)
|
||||
raise
|
||||
return wrapper
|
||||
@@ -148,6 +148,10 @@ class ObjectCache(Thread):
|
||||
with self._object_cache_lock:
|
||||
return self._object_cache.copy()
|
||||
|
||||
def get(self, name):
|
||||
with self._object_cache_lock:
|
||||
return self._object_cache.get(name)
|
||||
|
||||
def _build_cache(self):
|
||||
objects = self._list()
|
||||
return_type = 'V1' + objects.kind[:-4]
|
||||
@@ -244,11 +248,10 @@ class Kubernetes(AbstractDCS):
|
||||
self._api = CoreV1ApiProxy(config.get('use_endpoints'))
|
||||
self._should_create_config_service = self._api.use_endpoints
|
||||
self.reload_config(config)
|
||||
# leader_observed_record, leader_resource_version, and leader_observed_time are used only for leader race!
|
||||
self._leader_observed_record = {}
|
||||
self._leader_observed_time = None
|
||||
self._leader_resource_version = None
|
||||
self._leader_observed_subsets = []
|
||||
self._config_resource_version = None
|
||||
self.__do_not_watch = False
|
||||
|
||||
self._condition = Condition()
|
||||
@@ -307,14 +310,11 @@ class Kubernetes(AbstractDCS):
|
||||
with self._condition:
|
||||
self._wait_caches()
|
||||
|
||||
pods = self._pods.copy()
|
||||
self.__my_pod = pods.get(self._name)
|
||||
members = [self.member(pod) for pod in pods.values()]
|
||||
members = [self.member(pod) for pod in self._pods.copy().values()]
|
||||
nodes = self._kinds.copy()
|
||||
|
||||
config = nodes.get(self.config_path)
|
||||
metadata = config and config.metadata
|
||||
self._config_resource_version = metadata.resource_version if metadata else None
|
||||
annotations = metadata and metadata.annotations or {}
|
||||
|
||||
# get initialize flag
|
||||
@@ -332,8 +332,6 @@ class Kubernetes(AbstractDCS):
|
||||
leader = nodes.get(self.leader_path)
|
||||
metadata = leader and leader.metadata
|
||||
self._leader_resource_version = metadata.resource_version if metadata else None
|
||||
self._leader_observed_subsets = leader.subsets \
|
||||
if self._api.use_endpoints and leader and leader.subsets else []
|
||||
annotations = metadata and metadata.annotations or {}
|
||||
|
||||
# get last leader operation
|
||||
@@ -419,9 +417,9 @@ class Kubernetes(AbstractDCS):
|
||||
return True
|
||||
return False
|
||||
|
||||
def __target_ref(self, leader_ip, pod):
|
||||
def __target_ref(self, leader_ip, latest_subsets, pod):
|
||||
# we want to re-use existing target_ref if possible
|
||||
for subset in self._leader_observed_subsets:
|
||||
for subset in latest_subsets:
|
||||
for address in subset.addresses or []:
|
||||
if address.ip == leader_ip and address.target_ref and address.target_ref.name == self._name:
|
||||
return address.target_ref
|
||||
@@ -429,23 +427,24 @@ class Kubernetes(AbstractDCS):
|
||||
name=self._name, resource_version=pod.metadata.resource_version)
|
||||
|
||||
def _map_subsets(self, endpoints, ips):
|
||||
leader = self._kinds.get(self.leader_path)
|
||||
latest_subsets = leader and leader.subsets or []
|
||||
if not ips:
|
||||
# We want to have subsets empty
|
||||
if self._leader_observed_subsets:
|
||||
if latest_subsets:
|
||||
endpoints['subsets'] = []
|
||||
return
|
||||
|
||||
pod = self.__my_pod
|
||||
pod = self._pods.get(self._name)
|
||||
leader_ip = ips[0] or pod and pod.status.pod_ip
|
||||
# don't touch subsets if our (leader) ip is unknown or subsets is valid
|
||||
if leader_ip and self.subsets_changed(self._leader_observed_subsets, leader_ip, self.__ports):
|
||||
if leader_ip and self.subsets_changed(latest_subsets, leader_ip, self.__ports):
|
||||
kwargs = {'hostname': pod.spec.hostname, 'node_name': pod.spec.node_name,
|
||||
'target_ref': self.__target_ref(leader_ip, pod)} if pod else {}
|
||||
'target_ref': self.__target_ref(leader_ip, latest_subsets, pod)} if pod else {}
|
||||
address = k8s_client.V1EndpointAddress(ip=leader_ip, **kwargs)
|
||||
endpoints['subsets'] = [k8s_client.V1EndpointSubset(addresses=[address], ports=self.__ports)]
|
||||
|
||||
@catch_kubernetes_errors
|
||||
def patch_or_create(self, name, annotations, resource_version=None, patch=False, retry=True, ips=None):
|
||||
def _patch_or_create(self, name, annotations, resource_version=None, patch=False, retry=None, ips=None):
|
||||
metadata = {'namespace': self._namespace, 'name': name, 'labels': self._labels, 'annotations': annotations}
|
||||
if patch or resource_version:
|
||||
if resource_version is not None:
|
||||
@@ -463,20 +462,23 @@ class Kubernetes(AbstractDCS):
|
||||
body = k8s_client.V1Endpoints(**endpoints)
|
||||
else:
|
||||
body = k8s_client.V1ConfigMap(metadata=metadata)
|
||||
ret = self.retry(func, self._namespace, body) if retry else func(self._namespace, body)
|
||||
ret = retry(func, self._namespace, body) if retry else func(self._namespace, body)
|
||||
if ret:
|
||||
self._kinds.set(name, ret)
|
||||
return ret
|
||||
|
||||
@catch_kubernetes_errors
|
||||
def patch_or_create(self, name, annotations, resource_version=None, patch=False, retry=True, ips=None):
|
||||
if retry is True:
|
||||
retry = self.retry
|
||||
return self._patch_or_create(name, annotations, resource_version, patch, retry, ips)
|
||||
|
||||
def patch_or_create_config(self, annotations, resource_version=None, patch=False, retry=True):
|
||||
# SCOPE-config endpoint requires corresponding service otherwise it might be "cleaned" by k8s master
|
||||
if self._api.use_endpoints and not patch and not resource_version:
|
||||
self._should_create_config_service = True
|
||||
self._create_config_service()
|
||||
ret = self.patch_or_create(self.config_path, annotations, resource_version, patch, retry)
|
||||
if ret:
|
||||
self._config_resource_version = ret.metadata.resource_version
|
||||
return ret
|
||||
return self.patch_or_create(self.config_path, annotations, resource_version, patch, retry)
|
||||
|
||||
def _create_config_service(self):
|
||||
metadata = k8s_client.V1ObjectMeta(namespace=self._namespace, name=self.config_path, labels=self._labels)
|
||||
@@ -495,20 +497,58 @@ class Kubernetes(AbstractDCS):
|
||||
def _update_leader(self):
|
||||
"""Unused"""
|
||||
|
||||
def _update_leader_with_retry(self, annotations, resource_version, ips):
|
||||
retry = self._retry.copy()
|
||||
|
||||
def _retry(*args, **kwargs):
|
||||
return retry(*args, **kwargs)
|
||||
|
||||
try:
|
||||
return self._patch_or_create(self.leader_path, annotations, resource_version, ips=ips, retry=_retry)
|
||||
except k8s_client.rest.ApiException as e:
|
||||
if e.status == 409:
|
||||
logger.warning('Concurrent update of %s', self.leader_path)
|
||||
else:
|
||||
logger.exception('Permission denied' if e.status == 403 else 'Unexpected error from Kubernetes API')
|
||||
return False
|
||||
except RetryFailedError:
|
||||
return False
|
||||
|
||||
deadline = retry.stoptime - time.time()
|
||||
if deadline < 2:
|
||||
return False
|
||||
|
||||
retry.sleep_func(1) # Give a chance for ObjectCache to receive the latest version
|
||||
|
||||
kind = self._kinds.get(self.leader_path)
|
||||
kind_annotations = kind and kind.metadata.annotations or {}
|
||||
kind_resource_version = kind and kind.metadata.resource_version
|
||||
|
||||
# There is different leader or resource_version in cache didn't change
|
||||
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):
|
||||
kind = self._kinds.get(self.leader_path)
|
||||
kind_annotations = kind and kind.metadata.annotations or {}
|
||||
|
||||
if kind and kind_annotations.get(self._LEADER) != self._name:
|
||||
return False
|
||||
|
||||
now = datetime.datetime.now(tzutc).isoformat()
|
||||
leader_observed_record = kind_annotations or self._leader_observed_record
|
||||
annotations = {self._LEADER: self._name, 'ttl': str(self._ttl), 'renewTime': now,
|
||||
'acquireTime': self._leader_observed_record.get('acquireTime') or now,
|
||||
'transitions': self._leader_observed_record.get('transitions') or '0'}
|
||||
'acquireTime': leader_observed_record.get('acquireTime') or now,
|
||||
'transitions': leader_observed_record.get('transitions') or '0'}
|
||||
if last_operation:
|
||||
annotations[self._OPTIME] = last_operation
|
||||
|
||||
resource_version = kind and kind.metadata.resource_version
|
||||
ips = [] if access_is_restricted else self.__ips
|
||||
|
||||
ret = self.patch_or_create(self.leader_path, annotations, self._leader_resource_version, ips=ips)
|
||||
if ret:
|
||||
self._leader_resource_version = ret.metadata.resource_version
|
||||
return ret
|
||||
return self._update_leader_with_retry(annotations, resource_version, ips)
|
||||
|
||||
def attempt_to_acquire_leader(self, permanent=False):
|
||||
now = datetime.datetime.now(tzutc).isoformat()
|
||||
@@ -527,9 +567,7 @@ class Kubernetes(AbstractDCS):
|
||||
annotations['transitions'] = str(transitions)
|
||||
ips = [] if self._api.use_endpoints else None
|
||||
ret = self.patch_or_create(self.leader_path, annotations, self._leader_resource_version, ips=ips)
|
||||
if ret:
|
||||
self._leader_resource_version = ret.metadata.resource_version
|
||||
else:
|
||||
if not ret:
|
||||
logger.info('Could not take out TTL lock')
|
||||
return ret
|
||||
|
||||
@@ -545,6 +583,11 @@ class Kubernetes(AbstractDCS):
|
||||
patch = bool(self.cluster and isinstance(self.cluster.failover, Failover) and self.cluster.failover.index)
|
||||
return self.patch_or_create(self.failover_path, annotations, index, bool(index or patch), False)
|
||||
|
||||
@property
|
||||
def _config_resource_version(self):
|
||||
config = self._kinds.get(self.config_path)
|
||||
return config and config.metadata.resource_version
|
||||
|
||||
def set_config_value(self, value, index=None):
|
||||
return self.patch_or_create_config({self._CONFIG: value}, index, bool(self._config_resource_version), False)
|
||||
|
||||
@@ -567,6 +610,8 @@ class Kubernetes(AbstractDCS):
|
||||
'annotations': {'status': json.dumps(data, separators=(',', ':'))}}
|
||||
body = k8s_client.V1Pod(metadata=k8s_client.V1ObjectMeta(**metadata))
|
||||
ret = self._api.patch_namespaced_pod(self._name, self._namespace, body)
|
||||
if ret:
|
||||
self._pods.set(self._name, ret)
|
||||
if self._should_create_config_service:
|
||||
self._create_config_service()
|
||||
return ret
|
||||
@@ -580,11 +625,12 @@ class Kubernetes(AbstractDCS):
|
||||
"""Unused"""
|
||||
|
||||
def delete_leader(self, last_operation=None):
|
||||
if self.cluster and isinstance(self.cluster.leader, Leader) and self.cluster.leader.name == self._name:
|
||||
kind = self._kinds.get(self.leader_path)
|
||||
if kind and (kind.metadata.annotations or {}).get(self._LEADER) == self._name:
|
||||
annotations = {self._LEADER: None}
|
||||
if last_operation:
|
||||
annotations[self._OPTIME] = last_operation
|
||||
self.patch_or_create(self.leader_path, annotations, self._leader_resource_version, True, False, [])
|
||||
self.patch_or_create(self.leader_path, annotations, kind.metadata.resource_version, True, False, [])
|
||||
self.reset_cluster()
|
||||
|
||||
def cancel_initialization(self):
|
||||
|
||||
@@ -101,8 +101,9 @@ class TestKubernetesConfigMaps(BaseTestKubernetes):
|
||||
def test_set_config_value(self):
|
||||
self.k.set_config_value('{}')
|
||||
|
||||
@patch.object(k8s_client.CoreV1Api, 'patch_namespaced_pod', Mock(return_value=True))
|
||||
def test_touch_member(self):
|
||||
@patch.object(k8s_client.CoreV1Api, 'patch_namespaced_pod')
|
||||
def test_touch_member(self, mock_patch_namespaced_pod):
|
||||
mock_patch_namespaced_pod.return_value.metadata.resource_version = '10'
|
||||
self.k.touch_member({'role': 'replica'})
|
||||
self.k._name = 'p-1'
|
||||
self.k.touch_member({'state': 'running', 'role': 'replica'})
|
||||
@@ -142,15 +143,32 @@ class TestKubernetesEndpoints(BaseTestKubernetes):
|
||||
self.assertIsNotNone(self.k.update_leader('123'))
|
||||
args = mock_patch_namespaced_endpoints.call_args[0]
|
||||
self.assertEqual(args[2].subsets[0].addresses[0].target_ref.resource_version, '10')
|
||||
self.k._leader_observed_subsets = []
|
||||
self.k._kinds._object_cache['test'].subsets[:] = []
|
||||
self.assertIsNotNone(self.k.update_leader('123'))
|
||||
self.k._kinds._object_cache['test'].metadata.annotations['leader'] = 'p-1'
|
||||
self.assertFalse(self.k.update_leader('123'))
|
||||
|
||||
@patch.object(k8s_client.CoreV1Api, 'patch_namespaced_endpoints', mock_namespaced_kind)
|
||||
def test_update_leader_with_restricted_access(self):
|
||||
self.assertIsNotNone(self.k.update_leader('123', True))
|
||||
|
||||
@patch.object(k8s_client.CoreV1Api, 'patch_namespaced_endpoints')
|
||||
def test__update_leader_with_retry(self, mock_patch):
|
||||
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])):
|
||||
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', []))
|
||||
|
||||
@patch.object(k8s_client.CoreV1Api, 'create_namespaced_endpoints',
|
||||
Mock(side_effect=[k8s_client.rest.ApiException(502, ''), k8s_client.rest.ApiException(500, '')]))
|
||||
Mock(side_effect=[k8s_client.rest.ApiException(500, ''), k8s_client.rest.ApiException(502, '')]))
|
||||
def test_delete_sync_state(self):
|
||||
self.assertFalse(self.k.delete_sync_state())
|
||||
|
||||
|
||||
Reference in New Issue
Block a user