diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java index 19f81e1e..2bf08a77 100644 --- a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java +++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatch.java @@ -17,7 +17,10 @@ package org.springframework.cloud.kubernetes.discovery; +import io.fabric8.kubernetes.api.model.EndpointAddress; +import io.fabric8.kubernetes.api.model.EndpointSubset; import io.fabric8.kubernetes.api.model.Endpoints; +import io.fabric8.kubernetes.api.model.ObjectReference; import io.fabric8.kubernetes.client.KubernetesClient; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -26,7 +29,9 @@ import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisherAware; import org.springframework.scheduling.annotation.Scheduled; +import java.util.Collection; import java.util.List; +import java.util.Objects; import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; @@ -62,8 +67,12 @@ public class KubernetesCatalogWatch implements ApplicationEventPublisherAware { List endpointsPodNames = endpoints.stream() .flatMap(endpoint -> endpoint.getSubsets().stream()) - .flatMap(subset -> subset.getAddresses().stream()) - .map(endpointAddress -> endpointAddress.getTargetRef().getName()) // pod name unique in namespace + .map(EndpointSubset::getAddresses) + .filter(Objects::nonNull) + .flatMap(Collection::stream) + .map(EndpointAddress::getTargetRef) + .filter(Objects::nonNull) + .map(ObjectReference::getName) // pod name unique in namespace .sorted(String::compareTo).collect(Collectors.toList()); catalogEndpointsState.set(endpointsPodNames); diff --git a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatchTest.java b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatchTest.java index e7e1b671..1f6b126c 100644 --- a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatchTest.java +++ b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesCatalogWatchTest.java @@ -114,6 +114,38 @@ public class KubernetesCatalogWatchTest { assertEquals(expectedPodsList, event.getValue()); } + @Test + public void testEndpointsWithoutAddresses() { + + EndpointsList endpoints = createSingleEndpointEndpointListByPodName("api-pod"); + endpoints.getItems().get(0).getSubsets().get(0).setAddresses(null); + + when(endpointsOperation.list()).thenReturn(endpoints); + when(kubernetesClient.endpoints()).thenReturn(endpointsOperation); + + underTest.catalogServicesWatch(); + // second execution on shuffleServices + underTest.catalogServicesWatch(); + + verify(applicationEventPublisher).publishEvent(any(HeartbeatEvent.class)); + } + + @Test + public void testEndpointsWithoutTargetRefs() { + + EndpointsList endpoints = createSingleEndpointEndpointListByPodName("api-pod"); + endpoints.getItems().get(0).getSubsets().get(0).getAddresses().get(0).setTargetRef(null); + + when(endpointsOperation.list()).thenReturn(endpoints); + when(kubernetesClient.endpoints()).thenReturn(endpointsOperation); + + underTest.catalogServicesWatch(); + // second execution on shuffleServices + underTest.catalogServicesWatch(); + + verify(applicationEventPublisher).publishEvent(any(HeartbeatEvent.class)); + } + private EndpointsList createEndpointsListByServiceName(String... serviceNames) { List endpoints = stream(serviceNames)