Add support for filtering by service-labels to Informer based DiscoveryClient. Fixes #810 (#877)

This commit is contained in:
Ryan Baxter
2021-10-04 19:19:38 -04:00
committed by GitHub
parent 493ae0c2ac
commit c48fdfd121
4 changed files with 92 additions and 17 deletions

View File

@@ -48,6 +48,7 @@ import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscover
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.env.Environment;
@Configuration(proxyBeanMethods = false)
@ConditionalOnKubernetesDiscoveryEnabled
@@ -80,11 +81,15 @@ public class KubernetesDiscoveryClientAutoConfiguration {
@Bean
@ConditionalOnMissingBean
public SpringCloudKubernetesInformerFactoryProcessor discoveryInformerConfigurer(
KubernetesNamespaceProvider kubernetesNamespaceProvider,
KubernetesDiscoveryProperties kubernetesDiscoveryProperties, ApiClient apiClient,
CatalogSharedInformerFactory sharedInformerFactory) {
return new SpringCloudKubernetesInformerFactoryProcessor(kubernetesDiscoveryProperties,
kubernetesNamespaceProvider, apiClient, sharedInformerFactory);
KubernetesNamespaceProvider kubernetesNamespaceProvider, ApiClient apiClient,
CatalogSharedInformerFactory sharedInformerFactory, Environment environment) {
// Injecting KubernetesDiscoveryProperties here would cause it to be
// initialize too early
// Instead get the all-namespaces property value from the Environment directly
boolean allNamespaces = environment.getProperty("spring.cloud.kubernetes.discovery.all-namespaces",
Boolean.class, false);
return new SpringCloudKubernetesInformerFactoryProcessor(kubernetesNamespaceProvider, apiClient,
sharedInformerFactory, allNamespaces);
}
@Bean

View File

@@ -102,7 +102,7 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
V1Service service = properties.isAllNamespaces() ? this.serviceLister.list().stream()
.filter(svc -> serviceId.equals(svc.getMetadata().getName())).findFirst().orElse(null)
: this.serviceLister.namespace(this.namespace).get(serviceId);
if (service == null) {
if (service == null || !matchServiceLabels(service)) {
// no such service present in the cluster
return new ArrayList<>();
}
@@ -205,8 +205,8 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
public List<String> getServices() {
List<V1Service> services = this.properties.isAllNamespaces() ? this.serviceLister.list()
: this.serviceLister.namespace(this.namespace).list();
return services.stream().filter(s -> s.getMetadata() != null) // safeguard
.map(s -> s.getMetadata().getName()).collect(Collectors.toList());
return services.stream().filter(this::matchServiceLabels).map(s -> s.getMetadata().getName())
.collect(Collectors.toList());
}
@Override
@@ -230,4 +230,28 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
+ " services) , discovery client is now available");
}
private boolean matchServiceLabels(V1Service service) {
if (log.isDebugEnabled()) {
log.debug("Kubernetes Service Label Properties:");
if (this.properties.getServiceLabels() != null) {
this.properties.getServiceLabels().forEach((key, value) -> log.debug(key + ":" + value));
}
log.debug("Service " + service.getMetadata().getName() + " labels:");
if (service.getMetadata() != null && service.getMetadata().getLabels() != null) {
service.getMetadata().getLabels().forEach((key, value) -> log.debug(key + ":" + value));
}
}
// safeguard
if (service.getMetadata() == null) {
return false;
}
if (properties.getServiceLabels() == null || properties.getServiceLabels().isEmpty()) {
return true;
}
return properties.getServiceLabels().keySet().stream()
.allMatch(k -> service.getMetadata().getLabels() != null
&& service.getMetadata().getLabels().containsKey(k)
&& service.getMetadata().getLabels().get(k).equals(properties.getServiceLabels().get(k)));
}
}

View File

@@ -38,7 +38,6 @@ import org.springframework.beans.factory.support.AbstractBeanDefinition;
import org.springframework.beans.factory.support.BeanDefinitionRegistry;
import org.springframework.beans.factory.support.RootBeanDefinition;
import org.springframework.cloud.kubernetes.commons.KubernetesNamespaceProvider;
import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryProperties;
import org.springframework.core.ResolvableType;
/**
@@ -54,26 +53,24 @@ class SpringCloudKubernetesInformerFactoryProcessor extends KubernetesInformerFa
private final SharedInformerFactory sharedInformerFactory;
private final KubernetesDiscoveryProperties kubernetesDiscoveryProperties;
private final boolean allNamespaces;
private final KubernetesNamespaceProvider kubernetesNamespaceProvider;
@Autowired
SpringCloudKubernetesInformerFactoryProcessor(KubernetesDiscoveryProperties kubernetesDiscoveryProperties,
KubernetesNamespaceProvider kubernetesNamespaceProvider, ApiClient apiClient,
SharedInformerFactory sharedInformerFactory) {
SpringCloudKubernetesInformerFactoryProcessor(KubernetesNamespaceProvider kubernetesNamespaceProvider,
ApiClient apiClient, SharedInformerFactory sharedInformerFactory, boolean allNamespaces) {
super();
this.apiClient = apiClient;
this.sharedInformerFactory = sharedInformerFactory;
this.kubernetesNamespaceProvider = kubernetesNamespaceProvider;
this.kubernetesDiscoveryProperties = kubernetesDiscoveryProperties;
this.allNamespaces = allNamespaces;
}
@Override
public void postProcessBeanFactory(ConfigurableListableBeanFactory beanFactory) throws BeansException {
String namespace = kubernetesDiscoveryProperties.isAllNamespaces() ? Namespaces.NAMESPACE_ALL
: kubernetesNamespaceProvider.getNamespace() == null ? Namespaces.NAMESPACE_DEFAULT
: kubernetesNamespaceProvider.getNamespace();
String namespace = allNamespaces ? Namespaces.NAMESPACE_ALL : kubernetesNamespaceProvider.getNamespace() == null
? Namespaces.NAMESPACE_DEFAULT : kubernetesNamespaceProvider.getNamespace();
this.apiClient.setHttpClient(this.apiClient.getHttpClient().newBuilder().readTimeout(Duration.ZERO).build());

View File

@@ -17,6 +17,7 @@
package org.springframework.cloud.kubernetes.client.discovery;
import java.util.HashMap;
import java.util.Map;
import io.kubernetes.client.informer.SharedInformerFactory;
import io.kubernetes.client.informer.cache.Cache;
@@ -59,6 +60,11 @@ public class KubernetesInformerDiscoveryClientTests {
.metadata(new V1ObjectMeta().name("test-svc-1").namespace("namespace2"))
.spec(new V1ServiceSpec().loadBalancerIP("1.1.1.1")).status(new V1ServiceStatus());
private static final V1Service testService3 = new V1Service()
.metadata(new V1ObjectMeta().name("test-svc-3").namespace("namespace1").putLabelsItem("spring", "true")
.putLabelsItem("k8s", "true"))
.spec(new V1ServiceSpec().loadBalancerIP("1.1.1.1")).status(new V1ServiceStatus());
private static final V1Endpoints testEndpoints1 = new V1Endpoints()
.metadata(new V1ObjectMeta().name("test-svc-1").namespace("namespace1"))
.addSubsetsItem(new V1EndpointSubset().addPortsItem(new V1EndpointPort().port(8080))
@@ -91,6 +97,11 @@ public class KubernetesInformerDiscoveryClientTests {
.addPortsItem(new V1EndpointPort().name("tcp2").port(443))
.addAddressesItem(new V1EndpointAddress().ip("1.1.1.1")));
private static final V1Endpoints testEndpoints3 = new V1Endpoints()
.metadata(new V1ObjectMeta().name("test-svc-3").namespace("namespace1"))
.addSubsetsItem(new V1EndpointSubset().addPortsItem(new V1EndpointPort().port(8080))
.addAddressesItem(new V1EndpointAddress().ip("2.2.2.2")));
@Test
public void testDiscoveryGetServicesAllNamespaceShouldWork() {
Lister<V1Service> serviceLister = setupServiceLister(testService1, testService2);
@@ -106,6 +117,44 @@ public class KubernetesInformerDiscoveryClientTests {
verify(kubernetesDiscoveryProperties, times(1)).isAllNamespaces();
}
@Test
public void testDiscoveryWithServiceLabels() {
Lister<V1Service> serviceLister = setupServiceLister(testService1, testService2, testService3);
Map<String, String> labels = new HashMap<>();
labels.put("k8s", "true");
labels.put("spring", "true");
when(kubernetesDiscoveryProperties.getServiceLabels()).thenReturn(labels);
KubernetesInformerDiscoveryClient discoveryClient = new KubernetesInformerDiscoveryClient("",
sharedInformerFactory, serviceLister, null, null, null, kubernetesDiscoveryProperties);
assertThat(discoveryClient.getServices().toArray()).containsOnly(testService3.getMetadata().getName());
verify(kubernetesDiscoveryProperties, times(1)).isAllNamespaces();
}
@Test
public void testDiscoveryInstancesWithServiceLabels() {
Lister<V1Service> serviceLister = setupServiceLister(testService1, testService2, testService3);
Lister<V1Endpoints> endpointsLister = setupEndpointsLister(testEndpoints1, testEndpoints3);
Map<String, String> labels = new HashMap<>();
labels.put("k8s", "true");
labels.put("spring", "true");
when(kubernetesDiscoveryProperties.isAllNamespaces()).thenReturn(true);
when(kubernetesDiscoveryProperties.getServiceLabels()).thenReturn(labels);
KubernetesInformerDiscoveryClient discoveryClient = new KubernetesInformerDiscoveryClient("",
sharedInformerFactory, serviceLister, endpointsLister, null, null, kubernetesDiscoveryProperties);
assertThat(discoveryClient.getInstances("test-svc-1").toArray()).isEmpty();
assertThat(discoveryClient.getInstances("test-svc-3").toArray())
.containsOnly(new KubernetesServiceInstance("", "test-svc-3", "2.2.2.2", 8080, new HashMap<>(), false));
}
@Test
public void testDiscoveryGetServicesOneNamespaceShouldWork() {
Lister<V1Service> serviceLister = setupServiceLister(testService1, testService2);