diff --git a/pom.xml b/pom.xml index e4480642..62ab5879 100644 --- a/pom.xml +++ b/pom.xml @@ -91,12 +91,7 @@ - - org.codehaus.groovy - groovy-all - ${groovy.version} - - + org.springframework.cloud spring-cloud-kubernetes-dependencies diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClient.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClient.java index 9bbe92ed..b9eac58d 100644 --- a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClient.java +++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClient.java @@ -14,13 +14,13 @@ * limitations under the License. * */ - package org.springframework.cloud.kubernetes.discovery; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.function.Predicate; import java.util.stream.Collectors; import io.fabric8.kubernetes.api.model.EndpointAddress; @@ -31,9 +31,13 @@ import io.fabric8.kubernetes.client.KubernetesClient; import io.fabric8.kubernetes.client.utils.Utils; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; + import org.springframework.cloud.client.DefaultServiceInstance; import org.springframework.cloud.client.ServiceInstance; import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.expression.Expression; +import org.springframework.expression.spel.standard.SpelExpressionParser; +import org.springframework.expression.spel.support.SimpleEvaluationContext; import org.springframework.util.Assert; public class KubernetesDiscoveryClient implements DiscoveryClient { @@ -42,10 +46,16 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { private static final String HOSTNAME = "HOSTNAME"; private KubernetesClient client; - private KubernetesDiscoveryProperties properties; + private final KubernetesDiscoveryProperties properties; + private final SpelExpressionParser parser = new SpelExpressionParser(); + private final SimpleEvaluationContext evalCtxt = SimpleEvaluationContext + .forReadOnlyDataBinding() + .withInstanceMethods() + .build(); public KubernetesDiscoveryClient(KubernetesClient client, - KubernetesDiscoveryProperties kubernetesDiscoveryProperties) { + KubernetesDiscoveryProperties kubernetesDiscoveryProperties) { + this.client = client; this.properties = kubernetesDiscoveryProperties; } @@ -67,9 +77,9 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { String serviceName = properties.getServiceName(); String podName = System.getenv(HOSTNAME); ServiceInstance defaultInstance = new DefaultServiceInstance(serviceName, - "localhost", - 8080, - false); + "localhost", + 8080, + false); Endpoints endpoints = client.endpoints().withName(serviceName).get(); Optional service = Optional.ofNullable(client.services().withName(serviceName).get()); @@ -86,14 +96,14 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { List subsets = endpoints.getSubsets(); if (subsets != null) { - for (EndpointSubset subset : subsets) { - List addresses = subset.getAddresses(); - for (EndpointAddress address : addresses) { + for (EndpointSubset s : subsets) { + List addresses = s.getAddresses(); + for (EndpointAddress a : addresses) { return new KubernetesServiceInstance(serviceName, - address, - subset.getPorts().stream().findFirst().orElseThrow(IllegalStateException::new), - labels, - false); + a, + s.getPorts().stream().findFirst().orElseThrow(IllegalStateException::new), + labels, + false); } } } @@ -107,7 +117,7 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { @Override public List getInstances(String serviceId) { Assert.notNull(serviceId, - "[Assertion failed] - the object argument must be null"); + "[Assertion failed] - the object argument must be null"); Optional service = Optional.ofNullable(client.services().withName(serviceId).get()); final Map labels; if (service.isPresent()) { @@ -120,14 +130,14 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { List subsets = endpoints.get().getSubsets(); List instances = new ArrayList<>(); if (subsets != null) { - for (EndpointSubset subset : subsets) { - List addresses = subset.getAddresses(); - for (EndpointAddress address : addresses) { + for (EndpointSubset s : subsets) { + List addresses = s.getAddresses(); + for (EndpointAddress a : addresses) { instances.add(new KubernetesServiceInstance(serviceId, - address, - subset.getPorts().stream().findFirst().orElseThrow(IllegalStateException::new), - labels, - false)); + a, + s.getPorts().stream().findFirst().orElseThrow(IllegalStateException::new), + labels, + false)); } } } @@ -137,9 +147,30 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { @Override public List getServices() { - return client.services().list() - .getItems() - .stream().map(s -> s.getMetadata().getName()) - .collect(Collectors.toList()); + String spelExpression = properties.getFilter(); + Predicate filteredServices; + if (spelExpression == null || spelExpression.isEmpty()) { + filteredServices = (Service instance) -> true; + } else { + Expression filterExpr = parser.parseExpression(spelExpression); + filteredServices = (Service instance) -> { + Boolean include = filterExpr.getValue(evalCtxt, instance, Boolean.class); + if (include == null) { + return false; + } + return include; + }; + } + return getServices(filteredServices); } + + public List getServices(Predicate filter) { + return client.services().list() + .getItems() + .stream() + .filter(filter) + .map(s -> s.getMetadata().getName()) + .collect(Collectors.toList()); + } + } diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java index 15b4a0c6..5c80ea1f 100644 --- a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java +++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java @@ -29,6 +29,11 @@ public class KubernetesDiscoveryProperties extends AutoServiceRegistrationProper @Value("${spring.application.name:unknown}") private String serviceName = "unknown"; + /** + * SpEL expression to filter services + **/ + private String filter; + public boolean isEnabled() { return enabled; } @@ -41,11 +46,20 @@ public class KubernetesDiscoveryProperties extends AutoServiceRegistrationProper return serviceName; } + public String getFilter() { + return filter; + } + + public void setFilter(String filter){ + this.filter = filter; + } + @Override public String toString() { return "KubernetesDiscoveryProperties{" + "enabled=" + enabled + ", serviceName='" + serviceName + '\'' + + ", filter='" + filter + '\'' + '}'; } } diff --git a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientTest.java b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientTest.java new file mode 100644 index 00000000..c5fce284 --- /dev/null +++ b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientTest.java @@ -0,0 +1,141 @@ +/* + * Copyright 2013-2018 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ +package org.springframework.cloud.kubernetes.discovery; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; + +import io.fabric8.kubernetes.api.model.DoneableService; +import io.fabric8.kubernetes.api.model.ObjectMeta; +import io.fabric8.kubernetes.api.model.Service; +import io.fabric8.kubernetes.api.model.ServiceList; +import io.fabric8.kubernetes.client.KubernetesClient; +import io.fabric8.kubernetes.client.dsl.MixedOperation; +import io.fabric8.kubernetes.client.dsl.Resource; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.MockitoJUnitRunner; + +import static org.junit.Assert.assertEquals; +import static org.mockito.Mockito.when; + +@RunWith(MockitoJUnitRunner.class) +public class KubernetesDiscoveryClientTest { + + @Mock + private KubernetesClient kubernetesClient; + + @Mock + private KubernetesDiscoveryProperties properties; + + @Mock + private MixedOperation> serviceOperation; + + @InjectMocks + private KubernetesDiscoveryClient underTest; + + @Test + public void testFilteredServices() { + List springBootServiceNames = Arrays.asList("serviceA", "serviceB"); + List services = createSpringBootServiceByName(springBootServiceNames); + + // Add non spring boot service + Service service = new Service(); + ObjectMeta objectMeta = new ObjectMeta(); + objectMeta.setName("ServiceNonSpringBoot"); + service.setMetadata(objectMeta); + services.add(service); + + ServiceList serviceList = new ServiceList(); + serviceList.setItems(services); + when(serviceOperation.list()) + .thenReturn(serviceList); + when(kubernetesClient.services()).thenReturn(serviceOperation); + + when(properties.getFilter()).thenReturn("metadata.additionalProperties['spring-boot']"); + + List filteredServices = underTest.getServices(); + + System.out.println("Filtered Services: " + filteredServices); + assertEquals(springBootServiceNames, filteredServices); + + } + + @Test + public void testFilteredServicesByPrefix() { + List springBootServiceNames = Arrays.asList("serviceA", "serviceB", "serviceC"); + List services = createSpringBootServiceByName(springBootServiceNames); + + // Add non spring boot service + Service service = new Service(); + ObjectMeta objectMeta = new ObjectMeta(); + objectMeta.setName("anotherService"); + service.setMetadata(objectMeta); + services.add(service); + + ServiceList serviceList = new ServiceList(); + serviceList.setItems(services); + when(serviceOperation.list()) + .thenReturn(serviceList); + when(kubernetesClient.services()).thenReturn(serviceOperation); + + when(properties.getFilter()).thenReturn("metadata.name.startsWith('service')"); + + List filteredServices = underTest.getServices(); + + System.out.println("Filtered Services: " + filteredServices); + assertEquals(springBootServiceNames, filteredServices); + + } + + @Test + public void testNoExpression() { + List springBootServiceNames = Arrays.asList("serviceA", "serviceB", "serviceC"); + List services = createSpringBootServiceByName(springBootServiceNames); + + ServiceList serviceList = new ServiceList(); + serviceList.setItems(services); + when(serviceOperation.list()) + .thenReturn(serviceList); + when(kubernetesClient.services()).thenReturn(serviceOperation); + + when(properties.getFilter()).thenReturn(""); + + List filteredServices = underTest.getServices(); + + System.out.println("Filtered Services: " + filteredServices); + assertEquals(springBootServiceNames, filteredServices); + + } + + private List createSpringBootServiceByName(List serviceNames) { + List serviceCollection = new ArrayList<>(serviceNames.size()); + for (String serviceName : serviceNames) { + Service service = new Service(); + ObjectMeta objectMeta = new ObjectMeta(); + objectMeta.setName(serviceName); + objectMeta.setAdditionalProperty("spring-boot", "true"); + service.setMetadata(objectMeta); + serviceCollection.add(service); + } + return serviceCollection; + } + +}