diff --git a/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesDiscoveryClientUtils.java b/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesDiscoveryClientUtils.java index 38bae8f2..9b50ebb2 100644 --- a/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesDiscoveryClientUtils.java +++ b/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesDiscoveryClientUtils.java @@ -21,6 +21,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.function.Predicate; import java.util.function.Supplier; import io.kubernetes.client.informer.SharedInformerFactory; @@ -33,6 +34,9 @@ import org.apache.commons.logging.LogFactory; import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryProperties; import org.springframework.core.log.LogAccessor; +import org.springframework.expression.Expression; +import org.springframework.expression.spel.standard.SpelExpressionParser; +import org.springframework.expression.spel.support.SimpleEvaluationContext; import static org.springframework.cloud.kubernetes.commons.config.ConfigUtils.keysWithPrefix; import static org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryConstants.NAMESPACE_METADATA_KEY; @@ -45,6 +49,11 @@ final class KubernetesDiscoveryClientUtils { private static final LogAccessor LOG = new LogAccessor(LogFactory.getLog(KubernetesDiscoveryClientUtils.class)); + private static final SpelExpressionParser PARSER = new SpelExpressionParser(); + + private static final SimpleEvaluationContext EVALUATION_CONTEXT = SimpleEvaluationContext.forReadOnlyDataBinding() + .withInstanceMethods().build(); + private KubernetesDiscoveryClientUtils() { } @@ -107,6 +116,24 @@ final class KubernetesDiscoveryClientUtils { return serviceMetadata; } + static Predicate filter(KubernetesDiscoveryProperties properties) { + String spelExpression = properties.filter(); + Predicate predicate; + if (spelExpression == null || spelExpression.isEmpty()) { + LOG.debug(() -> "filter not defined, returning always true predicate"); + predicate = service -> true; + } + else { + Expression filterExpr = PARSER.parseExpression(spelExpression); + predicate = service -> { + Boolean include = filterExpr.getValue(EVALUATION_CONTEXT, service, Boolean.class); + return Optional.ofNullable(include).orElse(false); + }; + LOG.debug(() -> "returning predicate based on filter expression: " + spelExpression); + } + return predicate; + } + static void postConstruct(List sharedInformerFactories, KubernetesDiscoveryProperties properties, Supplier informersReadyFunc, List> serviceListers) { @@ -118,16 +145,16 @@ final class KubernetesDiscoveryClientUtils { })) { if (properties.waitCacheReady()) { throw new IllegalStateException( - "Timeout waiting for informers cache to be ready, is the kubernetes service up?"); + "Timeout waiting for informers cache to be ready, is the kubernetes service up?"); } else { - LOG.warn(() -> "Timeout waiting for informers cache to be ready, " + - "ignoring the failure because waitForInformerCacheReady property is false"); + LOG.warn(() -> "Timeout waiting for informers cache to be ready, " + + "ignoring the failure because waitForInformerCacheReady property is false"); } } else { LOG.info(() -> "Cache fully loaded (total " + serviceListers.stream().mapToLong(x -> x.list().size()).sum() - + " services), discovery client is now available"); + + " services), discovery client is now available"); } } diff --git a/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClient.java b/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClient.java index deb201eb..bcc069d8 100644 --- a/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClient.java +++ b/spring-cloud-kubernetes-client-discovery/src/main/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClient.java @@ -22,6 +22,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Optional; +import java.util.function.Predicate; import java.util.function.Supplier; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -44,6 +45,7 @@ import org.springframework.core.log.LogAccessor; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; +import static org.springframework.cloud.kubernetes.client.discovery.KubernetesDiscoveryClientUtils.filter; import static org.springframework.cloud.kubernetes.client.discovery.KubernetesDiscoveryClientUtils.matchesServiceLabels; import static org.springframework.cloud.kubernetes.client.discovery.KubernetesDiscoveryClientUtils.postConstruct; import static org.springframework.cloud.kubernetes.client.discovery.KubernetesDiscoveryClientUtils.serviceMetadata; @@ -72,6 +74,8 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient { private final KubernetesDiscoveryProperties properties; + private final Predicate filter; + @Deprecated(forRemoval = true) public KubernetesInformerDiscoveryClient(String namespace, SharedInformerFactory sharedInformerFactory, Lister serviceLister, Lister endpointsLister, @@ -82,6 +86,7 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient { this.endpointsListers = List.of(endpointsLister); this.informersReadyFunc = () -> serviceInformer.hasSynced() && endpointsInformer.hasSynced(); this.properties = properties; + filter = filter(properties); } public KubernetesInformerDiscoveryClient(SharedInformerFactory sharedInformerFactory, @@ -93,6 +98,7 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient { this.endpointsListers = List.of(endpointsLister); this.informersReadyFunc = () -> serviceInformer.hasSynced() && endpointsInformer.hasSynced(); this.properties = properties; + filter = filter(properties); } public KubernetesInformerDiscoveryClient(List sharedInformerFactories, @@ -112,6 +118,7 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient { }; this.properties = properties; + filter = filter(properties); } @Override @@ -125,11 +132,11 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient { List services = serviceListers.stream().flatMap(x -> x.list().stream()) .filter(scv -> scv.getMetadata() != null).filter(svc -> serviceId.equals(svc.getMetadata().getName())) - .toList(); + .filter(filter).toList(); if (services.size() == 0 || services.stream().noneMatch(service -> matchesServiceLabels(service, properties))) { return List.of(); } - return services.stream().flatMap(s -> getServiceInstanceDetails(s, serviceId)).toList(); + return services.stream().flatMap(service -> getServiceInstanceDetails(service, serviceId)).toList(); } private Stream getServiceInstanceDetails(V1Service service, String serviceId) { @@ -229,8 +236,8 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient { @Override public List getServices() { List services = serviceListers.stream().flatMap(serviceLister -> serviceLister.list().stream()) - .filter(service -> matchesServiceLabels(service, properties)).map(s -> s.getMetadata().getName()) - .distinct().toList(); + .filter(service -> matchesServiceLabels(service, properties)).filter(filter) + .map(s -> s.getMetadata().getName()).distinct().toList(); LOG.debug(() -> "will return services : " + services); return services; } diff --git a/spring-cloud-kubernetes-client-discovery/src/test/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesDiscoveryClientFilterTests.java b/spring-cloud-kubernetes-client-discovery/src/test/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesDiscoveryClientFilterTests.java new file mode 100644 index 00000000..f7308940 --- /dev/null +++ b/spring-cloud-kubernetes-client-discovery/src/test/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesDiscoveryClientFilterTests.java @@ -0,0 +1,123 @@ +/* + * Copyright 2013-2023 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 + * + * https://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.client.discovery; + +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.function.Predicate; + +import io.kubernetes.client.openapi.models.V1Service; +import io.kubernetes.client.openapi.models.V1ServiceBuilder; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; + +import org.springframework.boot.test.system.CapturedOutput; +import org.springframework.boot.test.system.OutputCaptureExtension; +import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryProperties; + +/** + * @author wind57 + */ +@ExtendWith(OutputCaptureExtension.class) +class KubernetesDiscoveryClientFilterTests { + + @Test + void testEmptyExpression(CapturedOutput output) { + + String spelFilter = null; + KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(false, false, Set.of(), true, 60L, + false, spelFilter, Set.of(), Map.of(), null, null, 0, false); + + Predicate predicate = KubernetesDiscoveryClientUtils.filter(properties); + Assertions.assertNotNull(predicate); + Assertions.assertTrue(output.getOut().contains("filter not defined, returning always true predicate")); + } + + @Test + void testExpressionPresent(CapturedOutput output) { + + String spelFilter = "some"; + KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(false, false, Set.of(), true, 60L, + false, spelFilter, Set.of(), Map.of(), null, null, 0, false); + + Predicate predicate = KubernetesDiscoveryClientUtils.filter(properties); + Assertions.assertNotNull(predicate); + Assertions.assertTrue(output.getOut().contains("returning predicate based on filter expression: some")); + } + + @Test + void testTwoServicesBothMatch() { + String spelFilter = """ + #root.metadata.namespace matches "^.+A$" + """; + KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(false, false, Set.of(), true, 60L, + false, spelFilter, Set.of(), Map.of(), null, null, 0, false); + + V1Service a = new V1ServiceBuilder().withNewMetadata().withNamespace("namespace-A").withName("a").and().build(); + + V1Service b = new V1ServiceBuilder().withNewMetadata().withNamespace("namespace-A").withName("a").and().build(); + + List unfiltered = List.of(a, b); + Predicate predicate = KubernetesDiscoveryClientUtils.filter(properties); + List filtered = unfiltered.stream().filter(predicate) + .sorted(Comparator.comparing(service -> service.getMetadata().getName())).toList(); + Assertions.assertEquals(filtered.get(0).getMetadata().getName(), "a"); + Assertions.assertEquals(filtered.get(1).getMetadata().getName(), "a"); + } + + @Test + void testTwoServicesNoneMatch() { + String spelFilter = """ + #root.metadata.namespace matches "^.+A$" + """; + KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(false, false, Set.of(), true, 60L, + false, spelFilter, Set.of(), Map.of(), null, null, 0, false); + + V1Service a = new V1ServiceBuilder().withNewMetadata().withNamespace("namespace-B").withName("a").and().build(); + + V1Service b = new V1ServiceBuilder().withNewMetadata().withNamespace("namespace-B").withName("a").and().build(); + + List unfiltered = List.of(a, b); + Predicate predicate = KubernetesDiscoveryClientUtils.filter(properties); + List filtered = unfiltered.stream().filter(predicate) + .sorted(Comparator.comparing(service -> service.getMetadata().getName())).toList(); + Assertions.assertEquals(filtered.size(), 0); + } + + @Test + void testTwoServicesOneMatch() { + String spelFilter = """ + #root.metadata.namespace matches "^.+A$" + """; + KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(false, false, Set.of(), true, 60L, + false, spelFilter, Set.of(), Map.of(), null, null, 0, false); + + V1Service a = new V1ServiceBuilder().withNewMetadata().withNamespace("namespace-B").withName("a").and().build(); + + V1Service b = new V1ServiceBuilder().withNewMetadata().withNamespace("namespace-B").withName("a").and().build(); + + List unfiltered = List.of(a, b); + Predicate predicate = KubernetesDiscoveryClientUtils.filter(properties); + List filtered = unfiltered.stream().filter(predicate) + .sorted(Comparator.comparing(service -> service.getMetadata().getName())).toList(); + Assertions.assertEquals(filtered.size(), 0); + } + +} diff --git a/spring-cloud-kubernetes-client-discovery/src/test/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClientTests.java b/spring-cloud-kubernetes-client-discovery/src/test/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClientTests.java index 52484266..6042dc89 100644 --- a/spring-cloud-kubernetes-client-discovery/src/test/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClientTests.java +++ b/spring-cloud-kubernetes-client-discovery/src/test/java/org/springframework/cloud/kubernetes/client/discovery/KubernetesInformerDiscoveryClientTests.java @@ -16,6 +16,8 @@ package org.springframework.cloud.kubernetes.client.discovery; +import java.util.Comparator; +import java.util.List; import java.util.Map; import java.util.Set; @@ -381,6 +383,71 @@ class KubernetesInformerDiscoveryClientTests { "namespace2", null)); } + @Test + void testBothServicesMatchesFilter() { + Lister serviceLister = setupServiceLister(SERVICE_1, SERVICE_3); + Lister endpointsLister = setupEndpointsLister(ENDPOINTS_1, ENDPOINTS_3); + + String spelFilter = """ + #root.metadata.namespace matches "^.+1$" + """; + KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(false, false, Set.of(), true, 60L, + false, spelFilter, Set.of(), Map.of(), null, KubernetesDiscoveryProperties.Metadata.DEFAULT, 0, false); + + KubernetesInformerDiscoveryClient discoveryClient = new KubernetesInformerDiscoveryClient( + SHARED_INFORMER_FACTORY, serviceLister, endpointsLister, null, null, properties); + + assertThat(discoveryClient.getServices()).contains("test-svc-1", "test-svc-3"); + + List one = discoveryClient.getInstances("test-svc-1"); + assertThat(one.get(0).getMetadata().get("k8s_namespace")).isEqualTo("namespace1"); + + List two = discoveryClient.getInstances("test-svc-3"); + assertThat(two.get(0).getMetadata().get("k8s_namespace")).isEqualTo("namespace1"); + + } + + @Test + void testOneServiceMatchesFilter() { + Lister serviceLister = setupServiceLister(SERVICE_1, SERVICE_2); + Lister endpointsLister = setupEndpointsLister(ENDPOINTS_1, ENDPOINTS_2); + + // without filter, both match + String spelFilter = ""; + KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(false, false, Set.of(), true, 60L, + false, spelFilter, Set.of(), Map.of(), null, KubernetesDiscoveryProperties.Metadata.DEFAULT, 0, false); + + KubernetesInformerDiscoveryClient discoveryClient = new KubernetesInformerDiscoveryClient( + SHARED_INFORMER_FACTORY, serviceLister, endpointsLister, null, null, properties); + + // only one here because of distinct + assertThat(discoveryClient.getServices()).contains("test-svc-1"); + + List result = discoveryClient.getInstances("test-svc-1").stream() + .sorted(Comparator.comparing(res -> res.getMetadata().get("k8s_namespace"))).toList(); + assertThat(result.get(0).getMetadata().get("k8s_namespace")).isEqualTo("namespace1"); + assertThat(result.get(1).getMetadata().get("k8s_namespace")).isEqualTo("namespace2"); + + // with filter, only one matches + + spelFilter = """ + #root.metadata.namespace matches "^.+1$" + """; + properties = new KubernetesDiscoveryProperties(false, false, Set.of(), true, 60L, false, spelFilter, Set.of(), + Map.of(), null, KubernetesDiscoveryProperties.Metadata.DEFAULT, 0, false); + discoveryClient = new KubernetesInformerDiscoveryClient(SHARED_INFORMER_FACTORY, serviceLister, endpointsLister, + null, null, properties); + + // only one here because of distinct + assertThat(discoveryClient.getServices()).contains("test-svc-1"); + + result = discoveryClient.getInstances("test-svc-1").stream() + .sorted(Comparator.comparing(res -> res.getMetadata().get("k8s_namespace"))).toList(); + assertThat(result.size()).isEqualTo(1); + assertThat(result.get(0).getMetadata().get("k8s_namespace")).isEqualTo("namespace1"); + + } + private Lister setupServiceLister(V1Service... services) { Cache serviceCache = new Cache<>(); Lister serviceLister = new Lister<>(serviceCache); diff --git a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-client-configmap-event-reload/src/test/java/org/springframework/cloud/kubernetes/client/configmap/event/reload/ConfigMapEventReloadIT.java b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-client-configmap-event-reload/src/test/java/org/springframework/cloud/kubernetes/client/configmap/event/reload/ConfigMapEventReloadIT.java index b2014d7e..55dd4719 100644 --- a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-client-configmap-event-reload/src/test/java/org/springframework/cloud/kubernetes/client/configmap/event/reload/ConfigMapEventReloadIT.java +++ b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-client-configmap-event-reload/src/test/java/org/springframework/cloud/kubernetes/client/configmap/event/reload/ConfigMapEventReloadIT.java @@ -301,7 +301,8 @@ class ConfigMapEventReloadIT { .orElse(List.of())); if (secretsDisabled) { - V1EnvVar secretsDisabledEnvVar = new V1EnvVar().name("SPRING_CLOUD_KUBERNETES_SECRETS_ENABLED").value("FALSE"); + V1EnvVar secretsDisabledEnvVar = new V1EnvVar().name("SPRING_CLOUD_KUBERNETES_SECRETS_ENABLED") + .value("FALSE"); envVars.add(secretsDisabledEnvVar); deployment.getSpec().getTemplate().getSpec().getContainers().get(0).setEnv(envVars); } diff --git a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-client-discovery-it/src/test/java/org/springframework/cloud/kubernetes/client/discovery/it/KubernetesClientDiscoveryFilterIT.java b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-client-discovery-it/src/test/java/org/springframework/cloud/kubernetes/client/discovery/it/KubernetesClientDiscoveryFilterIT.java new file mode 100644 index 00000000..3895434c --- /dev/null +++ b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-client-discovery-it/src/test/java/org/springframework/cloud/kubernetes/client/discovery/it/KubernetesClientDiscoveryFilterIT.java @@ -0,0 +1,238 @@ +/* + * Copyright 2013-2023 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 + * + * https://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.client.discovery.it; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; + +import io.kubernetes.client.openapi.models.V1Deployment; +import io.kubernetes.client.openapi.models.V1EnvVar; +import io.kubernetes.client.openapi.models.V1Ingress; +import io.kubernetes.client.openapi.models.V1Service; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.testcontainers.k3s.K3sContainer; +import reactor.netty.http.client.HttpClient; +import reactor.util.retry.Retry; +import reactor.util.retry.RetryBackoffSpec; + +import org.springframework.cloud.kubernetes.commons.discovery.DefaultKubernetesServiceInstance; +import org.springframework.cloud.kubernetes.integration.tests.commons.Commons; +import org.springframework.cloud.kubernetes.integration.tests.commons.Phase; +import org.springframework.cloud.kubernetes.integration.tests.commons.native_client.Util; +import org.springframework.core.ParameterizedTypeReference; +import org.springframework.http.HttpMethod; +import org.springframework.http.client.reactive.ReactorClientHttpConnector; +import org.springframework.web.reactive.function.client.WebClient; + +/** + * @author wind57 + */ +class KubernetesClientDiscoveryFilterIT { + + private static final String FILTER_BOTH_NAMESPACES = "#root.metadata.namespace matches '^.*uat$'"; + + private static final String FILTER_SINGLE_NAMESPACE = "#root.metadata.namespace matches 'a-uat$'"; + + private static final String NAMESPACE_A_UAT = "a-uat"; + + private static final String NAMESPACE_B_UAT = "b-uat"; + + private static final String NAMESPACE = "default"; + + private static final String IMAGE_NAME = "spring-cloud-kubernetes-client-discovery-it"; + + private static Util util; + + private static final K3sContainer K3S = Commons.container(); + + @BeforeAll + static void beforeAll() throws Exception { + K3S.start(); + Commons.validateImage(IMAGE_NAME, K3S); + Commons.loadSpringCloudKubernetesImage(IMAGE_NAME, K3S); + + util = new Util(K3S); + } + + @BeforeEach + void beforeEach() { + util.createNamespace(NAMESPACE_A_UAT); + util.createNamespace(NAMESPACE_B_UAT); + util.setUpClusterWide(NAMESPACE, Set.of(NAMESPACE, NAMESPACE_A_UAT, NAMESPACE_B_UAT)); + util.wiremock(NAMESPACE_A_UAT, "/wiremock", Phase.CREATE); + util.wiremock(NAMESPACE_B_UAT, "/wiremock", Phase.CREATE); + } + + @AfterEach + void afterEach() { + util.wiremock(NAMESPACE_A_UAT, "/wiremock", Phase.DELETE); + util.wiremock(NAMESPACE_B_UAT, "/wiremock", Phase.DELETE); + util.deleteNamespace(NAMESPACE_A_UAT); + util.deleteNamespace(NAMESPACE_B_UAT); + } + + @AfterAll + static void after() throws Exception { + Commons.cleanUp(IMAGE_NAME, K3S); + } + + /** + *
+	 *     - service "wiremock" is present in namespace "a-uat"
+	 *     - service "wiremock" is present in namespace "b-uat"
+	 *
+	 *     - we search with a predicate : "#root.metadata.namespace matches '^uat.*$'"
+	 *
+	 *     As such, both services are found via 'getInstances' call.
+	 * 
+ */ + @Test + void filterMatchesBothNamespacesViaThePredicate() { + + manifests(Phase.CREATE, FILTER_BOTH_NAMESPACES); + + WebClient clientServices = builder().baseUrl("http://localhost/services").build(); + + @SuppressWarnings("unchecked") + List services = (List) clientServices.method(HttpMethod.GET).retrieve().bodyToMono(List.class) + .retryWhen(retrySpec()).block(); + + Assertions.assertEquals(services.size(), 1); + Assertions.assertTrue(services.contains("service-wiremock")); + + WebClient client = builder().baseUrl("http://localhost/service-instances/service-wiremock").build(); + List serviceInstances = client.method(HttpMethod.GET).retrieve() + .bodyToMono(new ParameterizedTypeReference>() { + + }).retryWhen(retrySpec()).block(); + + Assertions.assertEquals(serviceInstances.size(), 2); + List sorted = serviceInstances.stream() + .sorted(Comparator.comparing(DefaultKubernetesServiceInstance::getNamespace)).toList(); + + DefaultKubernetesServiceInstance first = sorted.get(0); + Assertions.assertEquals(first.getServiceId(), "service-wiremock"); + Assertions.assertNotNull(first.getInstanceId()); + Assertions.assertEquals(first.getPort(), 8080); + Assertions.assertEquals(first.getNamespace(), "a-uat"); + Assertions.assertEquals(first.getMetadata(), + Map.of("app", "service-wiremock", "http", "8080", "k8s_namespace", "a-uat", "type", "ClusterIP")); + + DefaultKubernetesServiceInstance second = sorted.get(1); + Assertions.assertEquals(second.getServiceId(), "service-wiremock"); + Assertions.assertNotNull(second.getInstanceId()); + Assertions.assertEquals(second.getPort(), 8080); + Assertions.assertEquals(second.getNamespace(), "b-uat"); + Assertions.assertEquals(second.getMetadata(), + Map.of("app", "service-wiremock", "http", "8080", "k8s_namespace", "b-uat", "type", "ClusterIP")); + + manifests(Phase.DELETE, FILTER_BOTH_NAMESPACES); + } + + /** + *
+	 *     - service "wiremock" is present in namespace "a-uat"
+	 *     - service "wiremock" is present in namespace "b-uat"
+	 *
+	 *     - we search with a predicate : "#root.metadata.namespace matches 'a-uat$'"
+	 *
+	 *     As such, only service from 'a-uat' namespace matches.
+	 * 
+ */ + @Test + void filterMatchesOneNamespaceViaThePredicate() { + manifests(Phase.CREATE, FILTER_SINGLE_NAMESPACE); + + WebClient clientServices = builder().baseUrl("http://localhost/services").build(); + + @SuppressWarnings("unchecked") + List services = (List) clientServices.method(HttpMethod.GET).retrieve().bodyToMono(List.class) + .retryWhen(retrySpec()).block(); + + Assertions.assertEquals(services.size(), 1); + Assertions.assertTrue(services.contains("service-wiremock")); + + WebClient client = builder().baseUrl("http://localhost/service-instances/service-wiremock").build(); + List serviceInstances = client.method(HttpMethod.GET).retrieve() + .bodyToMono(new ParameterizedTypeReference>() { + + }).retryWhen(retrySpec()).block(); + + Assertions.assertEquals(serviceInstances.size(), 1); + + DefaultKubernetesServiceInstance first = serviceInstances.get(0); + Assertions.assertEquals(first.getServiceId(), "service-wiremock"); + Assertions.assertNotNull(first.getInstanceId()); + Assertions.assertEquals(first.getPort(), 8080); + Assertions.assertEquals(first.getNamespace(), "a-uat"); + Assertions.assertEquals(first.getMetadata(), + Map.of("app", "service-wiremock", "http", "8080", "k8s_namespace", "a-uat", "type", "ClusterIP")); + + manifests(Phase.DELETE, FILTER_SINGLE_NAMESPACE); + } + + private static void manifests(Phase phase, String serviceFilter) { + + V1Deployment deployment = (V1Deployment) util.yaml("kubernetes-discovery-deployment.yaml"); + V1Service service = (V1Service) util.yaml("kubernetes-discovery-service.yaml"); + V1Ingress ingress = (V1Ingress) util.yaml("kubernetes-discovery-ingress.yaml"); + + List envVars = new ArrayList<>( + Optional.ofNullable(deployment.getSpec().getTemplate().getSpec().getContainers().get(0).getEnv()) + .orElse(List.of())); + V1EnvVar namespaceAUat = new V1EnvVar().name("SPRING_CLOUD_KUBERNETES_DISCOVERY_NAMESPACES_0") + .value(NAMESPACE_A_UAT); + V1EnvVar namespaceBUat = new V1EnvVar().name("SPRING_CLOUD_KUBERNETES_DISCOVERY_NAMESPACES_1") + .value(NAMESPACE_B_UAT); + V1EnvVar filter = new V1EnvVar().name("SPRING_CLOUD_KUBERNETES_DISCOVERY_FILTER").value(serviceFilter); + V1EnvVar debug = new V1EnvVar().name("LOGGING_LEVEL_ORG_SPRINGFRAMEWORK_CLOUD_KUBERNETES_CLIENT_DISCOVERY") + .value("DEBUG"); + envVars.add(namespaceAUat); + envVars.add(namespaceBUat); + envVars.add(filter); + envVars.add(debug); + deployment.getSpec().getTemplate().getSpec().getContainers().get(0).setEnv(envVars); + + if (phase.equals(Phase.CREATE)) { + util.createAndWait(NAMESPACE, null, deployment, service, ingress, true); + } + else if (phase.equals(Phase.DELETE)) { + util.deleteAndWait(NAMESPACE, deployment, service, ingress); + } + + } + + private WebClient.Builder builder() { + return WebClient.builder().clientConnector(new ReactorClientHttpConnector(HttpClient.create())); + } + + private RetryBackoffSpec retrySpec() { + return Retry.fixedDelay(15, Duration.ofSeconds(2)).filter(Objects::nonNull); + } + +}