diff --git a/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/KubernetesDiscoveryClientAutoConfiguration.java b/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/KubernetesDiscoveryClientAutoConfiguration.java index d834a1ea..1bc108b6 100644 --- a/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/KubernetesDiscoveryClientAutoConfiguration.java +++ b/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/KubernetesDiscoveryClientAutoConfiguration.java @@ -17,6 +17,7 @@ package org.springframework.cloud.kubernetes.fabric8.discovery; import io.fabric8.kubernetes.client.KubernetesClient; +import org.apache.commons.logging.LogFactory; import org.springframework.boot.actuate.health.HealthIndicator; import org.springframework.boot.autoconfigure.AutoConfigureAfter; @@ -40,6 +41,7 @@ import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.env.Environment; +import org.springframework.core.log.LogAccessor; /** * Auto configuration for discovery clients. @@ -56,6 +58,9 @@ import org.springframework.core.env.Environment; @AutoConfigureAfter({ Fabric8AutoConfiguration.class, KubernetesDiscoveryPropertiesAutoConfiguration.class }) public class KubernetesDiscoveryClientAutoConfiguration { + private static final LogAccessor LOG = new LogAccessor( + LogFactory.getLog(KubernetesDiscoveryClientAutoConfiguration.class)); + @Bean @ConditionalOnMissingBean public KubernetesClientServicesFunction servicesFunction(KubernetesDiscoveryProperties properties, @@ -77,6 +82,7 @@ public class KubernetesDiscoveryClientAutoConfiguration { @ConditionalOnDiscoveryHealthIndicatorEnabled public KubernetesDiscoveryClientHealthIndicatorInitializer indicatorInitializer( ApplicationEventPublisher applicationEventPublisher, PodUtils podUtils) { + LOG.debug(() -> "Will publish InstanceRegisteredEvent from blocking implementation"); return new KubernetesDiscoveryClientHealthIndicatorInitializer(podUtils, applicationEventPublisher); } diff --git a/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClient.java b/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClient.java index f6f7c57d..351d6244 100644 --- a/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClient.java +++ b/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClient.java @@ -45,7 +45,7 @@ public class KubernetesReactiveDiscoveryClient implements ReactiveDiscoveryClien @Override public String description() { - return "Kubernetes Reactive Discovery Client"; + return "Fabric8 Kubernetes Reactive Discovery Client"; } @Override diff --git a/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClientAutoConfiguration.java b/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClientAutoConfiguration.java index 73b8c292..948d37bf 100644 --- a/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClientAutoConfiguration.java +++ b/spring-cloud-kubernetes-fabric8-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClientAutoConfiguration.java @@ -17,6 +17,7 @@ package org.springframework.cloud.kubernetes.fabric8.discovery.reactive; import io.fabric8.kubernetes.client.KubernetesClient; +import org.apache.commons.logging.LogFactory; import org.springframework.boot.autoconfigure.AutoConfigureAfter; import org.springframework.boot.autoconfigure.AutoConfigureBefore; @@ -32,15 +33,19 @@ import org.springframework.cloud.client.discovery.composite.reactive.ReactiveCom import org.springframework.cloud.client.discovery.health.DiscoveryClientHealthIndicatorProperties; import org.springframework.cloud.client.discovery.health.reactive.ReactiveDiscoveryClientHealthIndicator; import org.springframework.cloud.client.discovery.simple.reactive.SimpleReactiveDiscoveryClientAutoConfiguration; +import org.springframework.cloud.kubernetes.commons.PodUtils; import org.springframework.cloud.kubernetes.commons.discovery.ConditionalOnKubernetesDiscoveryEnabled; +import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryClientHealthIndicatorInitializer; import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryProperties; import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryPropertiesAutoConfiguration; import org.springframework.cloud.kubernetes.fabric8.discovery.KubernetesClientServicesFunction; import org.springframework.cloud.kubernetes.fabric8.discovery.KubernetesClientServicesFunctionProvider; import org.springframework.cloud.kubernetes.fabric8.discovery.KubernetesDiscoveryClientAutoConfiguration; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.env.Environment; +import org.springframework.core.log.LogAccessor; /** * Auto configuration for reactive discovery client. @@ -58,6 +63,9 @@ import org.springframework.core.env.Environment; KubernetesDiscoveryClientAutoConfiguration.class, KubernetesDiscoveryPropertiesAutoConfiguration.class }) public class KubernetesReactiveDiscoveryClientAutoConfiguration { + private static final LogAccessor LOG = new LogAccessor( + LogFactory.getLog(KubernetesReactiveDiscoveryClientAutoConfiguration.class)); + @Bean @ConditionalOnMissingBean public KubernetesClientServicesFunction servicesFunction(KubernetesDiscoveryProperties properties, @@ -73,6 +81,18 @@ public class KubernetesReactiveDiscoveryClientAutoConfiguration { return new KubernetesReactiveDiscoveryClient(client, properties, kubernetesClientServicesFunction); } + /** + * Post an event so that health indicator is initialized. + */ + @Bean + @ConditionalOnClass(name = "org.springframework.boot.actuate.health.ReactiveHealthIndicator") + @ConditionalOnDiscoveryHealthIndicatorEnabled + KubernetesDiscoveryClientHealthIndicatorInitializer reactiveIndicatorInitializer( + ApplicationEventPublisher applicationEventPublisher, PodUtils podUtils) { + LOG.debug(() -> "Will publish InstanceRegisteredEvent from reactive implementation"); + return new KubernetesDiscoveryClientHealthIndicatorInitializer(podUtils, applicationEventPublisher); + } + @Bean @ConditionalOnClass(name = "org.springframework.boot.actuate.health.ReactiveHealthIndicator") @ConditionalOnDiscoveryHealthIndicatorEnabled diff --git a/spring-cloud-kubernetes-fabric8-discovery/src/test/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClientTests.java b/spring-cloud-kubernetes-fabric8-discovery/src/test/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClientTests.java index 33625931..75fee4c5 100644 --- a/spring-cloud-kubernetes-fabric8-discovery/src/test/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClientTests.java +++ b/spring-cloud-kubernetes-fabric8-discovery/src/test/java/org/springframework/cloud/kubernetes/fabric8/discovery/reactive/KubernetesReactiveDiscoveryClientTests.java @@ -76,7 +76,7 @@ class KubernetesReactiveDiscoveryClientTests { void verifyDefaults() { ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, KubernetesDiscoveryProperties.DEFAULT, KubernetesClient::services); - assertThat(client.description()).isEqualTo("Kubernetes Reactive Discovery Client"); + assertThat(client.description()).isEqualTo("Fabric8 Kubernetes Reactive Discovery Client"); assertThat(client.getOrder()).isEqualTo(ReactiveDiscoveryClient.DEFAULT_ORDER); } diff --git a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8ApplicationDiscoveryListener.java b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8ApplicationDiscoveryListener.java new file mode 100644 index 00000000..78534a95 --- /dev/null +++ b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8ApplicationDiscoveryListener.java @@ -0,0 +1,45 @@ +/* + * 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.fabric8.discovery; + +import io.fabric8.kubernetes.api.model.Pod; +import org.apache.commons.logging.LogFactory; + +import org.springframework.cloud.client.discovery.event.InstanceRegisteredEvent; +import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryClientHealthIndicatorInitializer; +import org.springframework.context.ApplicationListener; +import org.springframework.core.log.LogAccessor; +import org.springframework.stereotype.Component; + +/** + * @author wind57 + */ +@Component +public class Fabric8ApplicationDiscoveryListener implements ApplicationListener> { + + private static final LogAccessor LOG = new LogAccessor( + LogFactory.getLog(Fabric8ApplicationDiscoveryListener.class)); + + @Override + public void onApplicationEvent(InstanceRegisteredEvent event) { + Pod pod = (Pod) ((KubernetesDiscoveryClientHealthIndicatorInitializer.RegisteredEventSource) event.getSource()) + .pod(); + LOG.info(() -> "received InstanceRegisteredEvent from pod with 'app' label value : " + + pod.getMetadata().getLabels().get("app")); + } + +} diff --git a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8DiscoveryController.java b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8DiscoveryController.java index f0e36e25..e56a8622 100644 --- a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8DiscoveryController.java +++ b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8DiscoveryController.java @@ -20,6 +20,7 @@ import java.util.List; import io.fabric8.kubernetes.api.model.Endpoints; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.cloud.client.ServiceInstance; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; @@ -33,8 +34,10 @@ public class Fabric8DiscoveryController { private final KubernetesDiscoveryClient discoveryClient; - public Fabric8DiscoveryController(KubernetesDiscoveryClient discoveryClient) { - this.discoveryClient = discoveryClient; + public Fabric8DiscoveryController(ObjectProvider discoveryClient) { + KubernetesDiscoveryClient[] local = new KubernetesDiscoveryClient[1]; + discoveryClient.ifAvailable(x -> local[0] = x); + this.discoveryClient = local[0]; } @GetMapping("/services") diff --git a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8ReactiveDiscoveryController.java b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8ReactiveDiscoveryController.java new file mode 100644 index 00000000..d497f6f1 --- /dev/null +++ b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8ReactiveDiscoveryController.java @@ -0,0 +1,55 @@ +/* + * 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.fabric8.discovery; + +import java.util.List; + +import reactor.core.publisher.Mono; + +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.kubernetes.fabric8.discovery.reactive.KubernetesReactiveDiscoveryClient; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.RestController; + +/** + * @author wind57 + */ +@RestController +public class Fabric8ReactiveDiscoveryController { + + private final KubernetesReactiveDiscoveryClient reactiveDiscoveryClient; + + public Fabric8ReactiveDiscoveryController( + ObjectProvider reactiveDiscoveryClient) { + KubernetesReactiveDiscoveryClient[] local = new KubernetesReactiveDiscoveryClient[1]; + reactiveDiscoveryClient.ifAvailable(x -> local[0] = x); + this.reactiveDiscoveryClient = local[0]; + } + + @GetMapping("/reactive/services") + public Mono> allServices() { + return reactiveDiscoveryClient.getServices().collectList(); + } + + @GetMapping("/reactive/service-instances/{serviceId}") + public Mono> serviceInstances(@PathVariable("serviceId") String serviceId) { + return reactiveDiscoveryClient.getInstances(serviceId).collectList(); + } + +} diff --git a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/resources/application.yaml b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/resources/application.yaml new file mode 100644 index 00000000..cd402822 --- /dev/null +++ b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/main/resources/application.yaml @@ -0,0 +1,8 @@ +management: + endpoint: + health: + show-details: always + endpoints: + web: + exposure: + include: "*" diff --git a/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/test/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8DiscoveryClientHealthIT.java b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/test/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8DiscoveryClientHealthIT.java new file mode 100644 index 00000000..ee31842e --- /dev/null +++ b/spring-cloud-kubernetes-integration-tests/spring-cloud-kubernetes-fabric8-client-discovery/src/test/java/org/springframework/cloud/kubernetes/fabric8/discovery/Fabric8DiscoveryClientHealthIT.java @@ -0,0 +1,317 @@ +/* + * 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.fabric8.discovery; + +import java.io.InputStream; +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; + +import io.fabric8.kubernetes.api.model.EnvVar; +import io.fabric8.kubernetes.api.model.EnvVarBuilder; +import io.fabric8.kubernetes.api.model.Service; +import io.fabric8.kubernetes.api.model.apps.Deployment; +import io.fabric8.kubernetes.api.model.networking.v1.Ingress; +import io.fabric8.kubernetes.client.KubernetesClient; +import org.assertj.core.api.Assertions; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.testcontainers.containers.Container; +import org.testcontainers.k3s.K3sContainer; +import reactor.netty.http.client.HttpClient; +import reactor.util.retry.Retry; +import reactor.util.retry.RetryBackoffSpec; + +import org.springframework.boot.test.json.BasicJsonTester; +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.fabric8_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 Fabric8DiscoveryClientHealthIT { + + private static final String REACTIVE_STATUS = "$.components.reactiveDiscoveryClients.components.['Fabric8 Kubernetes Reactive Discovery Client'].status"; + + private static final String BLOCKING_STATUS = "$.components.discoveryComposite.components.discoveryClient.status"; + + private static final String NAMESPACE = "default"; + + private static final String IMAGE_NAME = "spring-cloud-kubernetes-fabric8-client-discovery"; + + private static KubernetesClient client; + + private static final BasicJsonTester BASIC_JSON_TESTER = new BasicJsonTester(Fabric8DiscoveryClientHealthIT.class); + + 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); + client = util.client(); + util.setUp(NAMESPACE); + } + + @AfterAll + static void after() throws Exception { + Commons.cleanUp(IMAGE_NAME, K3S); + } + + /** + * Reactive is disabled, only blocking is active. As such, + * KubernetesInformerDiscoveryClientAutoConfiguration::indicatorInitializer will post + * an InstanceRegisteredEvent. + * + * We assert for logs and call '/health' endpoint to see that blocking discovery + * client was initialized. + */ + @Test + void testBlockingConfiguration() { + + manifests(true, false, Phase.CREATE); + + assertLogStatement("Will publish InstanceRegisteredEvent from blocking implementation"); + assertLogStatement("publishing InstanceRegisteredEvent"); + assertLogStatement("Discovery Client has been initialized"); + assertLogStatement( + "received InstanceRegisteredEvent from pod with 'app' label value : spring-cloud-kubernetes-fabric8-client-discovery"); + + WebClient healthClient = builder().baseUrl("http://localhost/actuator/health").build(); + + String healthResult = healthClient.method(HttpMethod.GET).retrieve().bodyToMono(String.class) + .retryWhen(retrySpec()).block(); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)) + .extractingJsonPathStringValue("$.components.discoveryComposite.status").isEqualTo("UP"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)).extractingJsonPathStringValue(BLOCKING_STATUS) + .isEqualTo("UP"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)) + .extractingJsonPathArrayValue( + "$.components.discoveryComposite.components.discoveryClient.details.services") + .containsExactlyInAnyOrder("spring-cloud-kubernetes-fabric8-client-discovery", "kubernetes"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)).doesNotHaveJsonPath(REACTIVE_STATUS); + + manifests(true, false, Phase.DELETE); + } + + /** + * Both blocking and reactive are enabled. + */ + @Test + void testDefaultConfiguration() { + + manifests(false, false, Phase.CREATE); + + assertLogStatement("Will publish InstanceRegisteredEvent from blocking implementation"); + assertLogStatement("publishing InstanceRegisteredEvent"); + assertLogStatement("Discovery Client has been initialized"); + assertLogStatement("received InstanceRegisteredEvent from pod with 'app' label value : " + + "spring-cloud-kubernetes-fabric8-client-discovery"); + + WebClient healthClient = builder().baseUrl("http://localhost/actuator/health").build(); + + String healthResult = healthClient.method(HttpMethod.GET).retrieve().bodyToMono(String.class) + .retryWhen(retrySpec()).block(); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)) + .extractingJsonPathStringValue("$.components.discoveryComposite.status").isEqualTo("UP"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)) + .extractingJsonPathStringValue("$.components.discoveryComposite.components.discoveryClient.status") + .isEqualTo("UP"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)) + .extractingJsonPathArrayValue( + "$.components.discoveryComposite.components.discoveryClient.details.services") + .containsExactlyInAnyOrder("spring-cloud-kubernetes-fabric8-client-discovery", "kubernetes"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)) + .extractingJsonPathStringValue("$.components.reactiveDiscoveryClients.status").isEqualTo("UP"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)).extractingJsonPathStringValue( + "$.components.reactiveDiscoveryClients.components.['Fabric8 Kubernetes Reactive Discovery Client'].status") + .isEqualTo("UP"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)).extractingJsonPathArrayValue( + "$.components.reactiveDiscoveryClients.components.['Fabric8 Kubernetes Reactive Discovery Client'].details.services") + .containsExactlyInAnyOrder("spring-cloud-kubernetes-fabric8-client-discovery", "kubernetes"); + + manifests(false, false, Phase.DELETE); + } + + /** + * Reactive is enabled, blocking is disabled. As such, + * KubernetesInformerDiscoveryClientAutoConfiguration::indicatorInitializer will post + * an InstanceRegisteredEvent. + * + * We assert for logs and call '/health' endpoint to see that blocking discovery + * client was initialized. + */ + @Test + void testReactiveConfiguration() { + + manifests(false, true, Phase.CREATE); + + assertLogStatement("Will publish InstanceRegisteredEvent from reactive implementation"); + assertLogStatement("publishing InstanceRegisteredEvent"); + assertLogStatement("Discovery Client has been initialized"); + assertLogStatement( + "received InstanceRegisteredEvent from pod with 'app' label value : spring-cloud-kubernetes-fabric8-client-discovery"); + + WebClient healthClient = builder().baseUrl("http://localhost/actuator/health").build(); + + String healthResult = healthClient.method(HttpMethod.GET).retrieve().bodyToMono(String.class) + .retryWhen(retrySpec()).block(); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)) + .extractingJsonPathStringValue("$.components.reactiveDiscoveryClients.status").isEqualTo("UP"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)).extractingJsonPathStringValue(REACTIVE_STATUS) + .isEqualTo("UP"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)).extractingJsonPathArrayValue( + "$.components.reactiveDiscoveryClients.components.['Fabric8 Kubernetes Reactive Discovery Client'].details.services") + .containsExactlyInAnyOrder("spring-cloud-kubernetes-fabric8-client-discovery", "kubernetes"); + + Assertions.assertThat(BASIC_JSON_TESTER.from(healthResult)).doesNotHaveJsonPath(BLOCKING_STATUS); + + // test for services also: + + WebClient servicesClient = builder().baseUrl("http://localhost/reactive/services").build(); + + List servicesResult = servicesClient.method(HttpMethod.GET).retrieve() + .bodyToMono(new ParameterizedTypeReference>() { + }).retryWhen(retrySpec()).block(); + + Assertions.assertThat(servicesResult).contains("spring-cloud-kubernetes-fabric8-client-discovery"); + Assertions.assertThat(servicesResult).contains("kubernetes"); + + manifests(false, true, Phase.DELETE); + } + + private static void manifests(boolean disableReactive, boolean disableBlocking, Phase phase) { + + InputStream deploymentStream = util.inputStream("fabric8-discovery-deployment.yaml"); + InputStream serviceStream = util.inputStream("fabric8-discovery-service.yaml"); + InputStream ingressStream = util.inputStream("fabric8-discovery-ingress.yaml"); + + Deployment deployment = client.apps().deployments().load(deploymentStream).get(); + List envVars = new ArrayList<>( + deployment.getSpec().getTemplate().getSpec().getContainers().get(0).getEnv()); + + EnvVar debugLevelForCommons = new EnvVarBuilder() + .withName("LOGGING_LEVEL_ORG_SPRINGFRAMEWORK_CLOUD_KUBERNETES_COMMONS_DISCOVERY").withValue("DEBUG") + .build(); + + if (!disableBlocking) { + EnvVar debugBlockingEnvVar = new EnvVarBuilder() + .withName("LOGGING_LEVEL_ORG_SPRINGFRAMEWORK_CLOUD_CLIENT_DISCOVERY_HEALTH").withValue("DEBUG") + .build(); + + EnvVar debugLevelForBlocking = new EnvVarBuilder() + .withName("LOGGING_LEVEL_ORG_SPRINGFRAMEWORK_CLOUD_KUBERNETES_FABRIC8_DISCOVERY").withValue("DEBUG") + .build(); + + envVars.add(debugBlockingEnvVar); + envVars.add(debugLevelForBlocking); + } + + if (!disableReactive) { + EnvVar debugReactiveEnvVar = new EnvVarBuilder() + .withName("LOGGING_LEVEL_ORG_SPRINGFRAMEWORK_CLOUD_CLIENT_DISCOVERY_HEALTH_REACTIVE") + .withValue("DEBUG").build(); + + EnvVar debugLevelForReactive = new EnvVarBuilder() + .withName("LOGGING_LEVEL_ORG_SPRINGFRAMEWORK_CLOUD_KUBERNETES_FABRIC8_DISCOVERY_REACTIVE") + .withValue("DEBUG").build(); + + EnvVar debugLevelForBlocking = new EnvVarBuilder() + .withName("LOGGING_LEVEL_ORG_SPRINGFRAMEWORK_CLOUD_KUBERNETES_FABRIC8_DISCOVERY").withValue("DEBUG") + .build(); + + envVars.add(debugReactiveEnvVar); + envVars.add(debugLevelForBlocking); + envVars.add(debugLevelForReactive); + } + + Service service = client.services().load(serviceStream).get(); + Ingress ingress = client.network().v1().ingresses().load(ingressStream).get(); + + if (disableBlocking) { + EnvVar disableBlockingEnvVar = new EnvVarBuilder().withName("SPRING_CLOUD_DISCOVERY_BLOCKING_ENABLED") + .withValue("FALSE").build(); + envVars.add(disableBlockingEnvVar); + } + + if (disableReactive) { + EnvVar disableReactiveEnvVar = new EnvVarBuilder().withName("SPRING_CLOUD_DISCOVERY_REACTIVE_ENABLED") + .withValue("FALSE").build(); + envVars.add(disableReactiveEnvVar); + } + + envVars.add(debugLevelForCommons); + deployment.getSpec().getTemplate().getSpec().getContainers().get(0).setEnv(envVars); + + if (phase.equals(Phase.CREATE)) { + util.createAndWait(NAMESPACE, null, deployment, service, ingress, true); + } + else { + 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(1)).filter(Objects::nonNull); + } + + private void assertLogStatement(String message) { + try { + String appPodName = K3S.execInContainer("sh", "-c", + "kubectl get pods -l app=" + IMAGE_NAME + " -o=name --no-headers | tr -d '\n'").getStdout(); + + Container.ExecResult execResult = K3S.execInContainer("sh", "-c", "kubectl logs " + appPodName.trim()); + String ok = execResult.getStdout(); + Assertions.assertThat(ok).contains(message); + } + catch (Exception e) { + e.printStackTrace(); + throw new RuntimeException(e); + } + + } + +}