diff --git a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/reactive/KubernetesReactiveDiscoveryClientTests.java b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/reactive/KubernetesReactiveDiscoveryClientTests.java index 80f45d12..b9b5fc62 100644 --- a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/reactive/KubernetesReactiveDiscoveryClientTests.java +++ b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/reactive/KubernetesReactiveDiscoveryClientTests.java @@ -16,18 +16,20 @@ package org.springframework.cloud.kubernetes.discovery.reactive; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import io.fabric8.kubernetes.api.model.Endpoints; import io.fabric8.kubernetes.api.model.EndpointsBuilder; import io.fabric8.kubernetes.api.model.EndpointsList; import io.fabric8.kubernetes.api.model.ServiceBuilder; +import io.fabric8.kubernetes.api.model.ServiceList; import io.fabric8.kubernetes.api.model.ServiceListBuilder; import io.fabric8.kubernetes.client.Config; import io.fabric8.kubernetes.client.KubernetesClient; import io.fabric8.kubernetes.client.server.mock.KubernetesServer; import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import reactor.core.publisher.Flux; @@ -53,7 +55,7 @@ class KubernetesReactiveDiscoveryClientTests { public void setup(@Client KubernetesClient kubernetesClient) { // Configure the kubernetes master url to point to the mock server System.setProperty(Config.KUBERNETES_MASTER_SYSTEM_PROPERTY, - kubernetesClient.getConfiguration().getMasterUrl()); + kubernetesClient.getConfiguration().getMasterUrl()); System.setProperty(Config.KUBERNETES_TRUST_CERT_SYSTEM_PROPERTY, "true"); System.setProperty(Config.KUBERNETES_AUTH_TRYKUBECONFIG_SYSTEM_PROPERTY, "false"); System.setProperty(Config.KUBERNETES_AUTH_TRYSERVICEACCOUNT_SYSTEM_PROPERTY, "false"); @@ -64,228 +66,257 @@ class KubernetesReactiveDiscoveryClientTests { public void verifyDefaults(@Client KubernetesClient kubernetesClient) { KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); assertThat(client.description()).isEqualTo("Kubernetes Reactive Discovery Client"); assertThat(client.getOrder()).isEqualTo(ReactiveDiscoveryClient.DEFAULT_ORDER); } @Test public void shouldReturnFluxOfServices(@Client KubernetesClient kubernetesClient, - @Server KubernetesServer kubernetesServer) { + @Server KubernetesServer kubernetesServer) { kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services") - .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("s1") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().endItem().addNewItem().withNewMetadata().withName("s2") - .withLabels(new HashMap() { - { - put("label", "value"); - put("label2", "value2"); - } - }).endMetadata().endItem().addNewItem().withNewMetadata().withName("s3").endMetadata().endItem() - .build()) - .once(); + .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("s1") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata().endItem().addNewItem().withNewMetadata().withName("s2") + .withLabels(new HashMap() { + { + put("label", "value"); + put("label2", "value2"); + } + }).endMetadata().endItem().addNewItem().withNewMetadata().withName("s3").endMetadata().endItem() + .build()) + .once(); KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); Flux services = client.getServices(); StepVerifier.create(services).expectNext("s1", "s2", "s3").expectComplete().verify(); } @Test public void shouldReturnEmptyFluxOfServicesWhenNoInstancesFound(@Client KubernetesClient kubernetesClient, - @Server KubernetesServer kubernetesServer) { + @Server KubernetesServer kubernetesServer) { kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services") - .andReturn(200, new ServiceListBuilder().build()).once(); + .andReturn(200, new ServiceListBuilder().build()).once(); KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); Flux services = client.getServices(); StepVerifier.create(services).expectNextCount(0).expectComplete().verify(); } @Test - @Disabled // see gh-603 - public void shouldReturnEmptyFluxForNonExistingService(@Client KubernetesClient kubernetesClient) { + public void shouldReturnEmptyFluxForNonExistingService(@Client KubernetesClient kubernetesClient, + @Server KubernetesServer kubernetesServer) { + kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/endpoints?fieldSelector=metadata.name%3Dnonexistent-service") + .andReturn(200, new EndpointsBuilder().build()).once(); + KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); Flux instances = client.getInstances("nonexistent-service"); StepVerifier.create(instances).expectNextCount(0).expectComplete().verify(); } @Test - @Disabled // see gh-603 public void shouldReturnEmptyFluxWhenServiceHasNoSubsets(@Client KubernetesClient kubernetesClient, - @Server KubernetesServer kubernetesServer) { + @Server KubernetesServer kubernetesServer) { kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services") - .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().endItem().build()) - .once(); + .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata().endItem().build()) + .once(); + + kubernetesServer.expect().get() + .withPath("/api/v1/namespaces/test/endpoints?fieldSelector=metadata.name%3Dexisting-service") + .andReturn(200, new EndpointsBuilder().build()).once(); KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); Flux instances = client.getInstances("existing-service"); StepVerifier.create(instances).expectNextCount(0).expectComplete().verify(); } @Test - @Disabled // see gh-603 public void shouldReturnFlux(@Client KubernetesClient kubernetesClient, @Server KubernetesServer kubernetesServer) { - kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services") - .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().endItem().build()) - .once(); + ServiceList services = new ServiceListBuilder().addNewItem() + .withNewMetadata().withName("existing-service").withNamespace("test") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }) + .endMetadata() + .endItem().build(); - Endpoints endPoints = new EndpointsBuilder().withNewMetadata().withName("endpoint").withNamespace("test") - .endMetadata().addNewSubset().addNewAddress().withIp("ip1").withNewTargetRef().withUid("uid1") - .endTargetRef().endAddress().addNewPort("http", "http_tcp", 80, "TCP").endSubset().build(); + Endpoints endPoint = new EndpointsBuilder().withNewMetadata() + .withName("existing-service") + .withNamespace("test") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata() + .addNewSubset().addNewAddress().withIp("ip1") + .withNewTargetRef().withUid("uid1") + .endTargetRef().endAddress().addNewPort("http", "http_tcp", 80, "TCP") + .endSubset() + .build(); - kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/endpoints/existing-service") - .andReturn(200, endPoints).once(); + List endpointsList = new ArrayList<>(); + endpointsList.add(endPoint); + + EndpointsList endpoints = new EndpointsList(); + endpoints.setItems(endpointsList); + + kubernetesServer.expect().get() + .withPath("/api/v1/namespaces/test/endpoints?fieldSelector=metadata.name%3Dexisting-service") + .andReturn(200, endpoints).once(); kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services/existing-service") - .andReturn(200, new ServiceBuilder().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().build()) - .once(); + .andReturn(200, services.getItems().get(0)).once(); KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); + properties.getMetadata().setAddAnnotations(false); + properties.getMetadata().setAddLabels(false); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); Flux instances = client.getInstances("existing-service"); StepVerifier.create(instances).expectNextCount(1).expectComplete().verify(); } @Test - @Disabled // see gh-603 public void shouldReturnFluxWithPrefixedMetadata(@Client KubernetesClient kubernetesClient, - @Server KubernetesServer kubernetesServer) { + @Server KubernetesServer kubernetesServer) { kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services") - .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().endItem().build()) - .once(); + .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata().endItem().build()) + .once(); - Endpoints endPoints = new EndpointsBuilder().withNewMetadata().withName("endpoint").withNamespace("test") - .endMetadata().addNewSubset().addNewAddress().withIp("ip1").withNewTargetRef().withUid("uid1") - .endTargetRef().endAddress().addNewPort("http", "http_tcp", 80, "TCP").endSubset().build(); + Endpoints endPoint = new EndpointsBuilder().withNewMetadata().withName("endpoint").withNamespace("test") + .endMetadata().addNewSubset().addNewAddress().withIp("ip1").withNewTargetRef().withUid("uid1") + .endTargetRef().endAddress().addNewPort("http", "http_tcp", 80, "TCP").endSubset().build(); - kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/endpoints/existing-service") - .andReturn(200, endPoints).once(); + List endpointsList = new ArrayList<>(); + endpointsList.add(endPoint); + + EndpointsList endpoints = new EndpointsList(); + endpoints.setItems(endpointsList); + + kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/endpoints?fieldSelector=metadata.name%3Dexisting-service") + .andReturn(200, endpoints).once(); kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services/existing-service") - .andReturn(200, new ServiceBuilder().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().build()) - .once(); + .andReturn(200, new ServiceBuilder().withNewMetadata().withName("existing-service") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata().build()) + .once(); KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); properties.getMetadata().setAnnotationsPrefix("annotation."); properties.getMetadata().setLabelsPrefix("label."); properties.getMetadata().setPortsPrefix("port."); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); Flux instances = client.getInstances("existing-service"); StepVerifier.create(instances).expectNextCount(1).expectComplete().verify(); } @Test - @Disabled // see gh-603 public void shouldReturnFluxWhenServiceHasMultiplePortsAndPrimaryPortNameIsSet( - @Client KubernetesClient kubernetesClient, @Server KubernetesServer kubernetesServer) { + @Client KubernetesClient kubernetesClient, @Server KubernetesServer kubernetesServer) { kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services") - .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().endItem().build()) - .once(); + .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata().endItem().build()) + .once(); - Endpoints endPoints = new EndpointsBuilder().withNewMetadata().withName("endpoint").withNamespace("test") - .endMetadata().addNewSubset().addNewAddress().withIp("ip1").withNewTargetRef().withUid("uid1") - .endTargetRef().endAddress().addNewPort("http", "http_tcp", 80, "TCP") - .addNewPort("https", "https_tcp", 443, "TCP").endSubset().build(); + Endpoints endPoint = new EndpointsBuilder().withNewMetadata().withName("endpoint").withNamespace("test") + .endMetadata().addNewSubset().addNewAddress().withIp("ip1").withNewTargetRef().withUid("uid1") + .endTargetRef().endAddress().addNewPort("http", "http_tcp", 80, "TCP") + .addNewPort("https", "https_tcp", 443, "TCP").endSubset().build(); - kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/endpoints/existing-service") - .andReturn(200, endPoints).once(); + List endpointsList = new ArrayList<>(); + endpointsList.add(endPoint); + + EndpointsList endpoints = new EndpointsList(); + endpoints.setItems(endpointsList); + + kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/endpoints?fieldSelector=metadata.name%3Dexisting-service") + .andReturn(200, endpoints).once(); kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services/existing-service") - .andReturn(200, new ServiceBuilder().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().build()) - .once(); + .andReturn(200, new ServiceBuilder().withNewMetadata().withName("existing-service") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata().build()) + .once(); KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); - properties.setPrimaryPortName("https"); + properties.setPrimaryPortName("https_tcp"); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); Flux instances = client.getInstances("existing-service"); StepVerifier.create(instances).expectNextCount(1).expectComplete().verify(); } @Test public void shouldReturnFluxOfServicesAcrossAllNamespaces(@Client KubernetesClient kubernetesClient, - @Server KubernetesServer kubernetesServer) { + @Server KubernetesServer kubernetesServer) { kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services") - .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().endItem().build()) - .once(); + .andReturn(200, new ServiceListBuilder().addNewItem().withNewMetadata().withName("existing-service") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata().endItem().build()) + .once(); Endpoints endpoints = new EndpointsBuilder().withNewMetadata().withName("endpoint").withNamespace("test") - .endMetadata().addNewSubset().addNewAddress().withIp("ip1").withNewTargetRef().withUid("uid1") - .endTargetRef().endAddress().addNewPort("http", "http_tcp", 80, "TCP") - .addNewPort("https", "https_tcp", 443, "TCP").endSubset().build(); + .endMetadata().addNewSubset().addNewAddress().withIp("ip1").withNewTargetRef().withUid("uid1") + .endTargetRef().endAddress().addNewPort("http", "http_tcp", 80, "TCP") + .addNewPort("https", "https_tcp", 443, "TCP").endSubset().build(); EndpointsList endpointsList = new EndpointsList(); endpointsList.setItems(singletonList(endpoints)); kubernetesServer.expect().get().withPath("/api/v1/endpoints?fieldSelector=metadata.name%3Dexisting-service") - .andReturn(200, endpointsList).once(); + .andReturn(200, endpointsList).once(); kubernetesServer.expect().get().withPath("/api/v1/namespaces/test/services/existing-service") - .andReturn(200, new ServiceBuilder().withNewMetadata().withName("existing-service") - .withLabels(new HashMap() { - { - put("label", "value"); - } - }).endMetadata().build()) - .once(); + .andReturn(200, new ServiceBuilder().withNewMetadata().withName("existing-service") + .withLabels(new HashMap() { + { + put("label", "value"); + } + }).endMetadata().build()) + .once(); KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); properties.setAllNamespaces(true); ReactiveDiscoveryClient client = new KubernetesReactiveDiscoveryClient(kubernetesClient, properties, - KubernetesClient::services); + KubernetesClient::services); Flux instances = client.getInstances("existing-service"); StepVerifier.create(instances).expectNextCount(1).expectComplete().verify(); }