From 064af9822bc80d6cad20417e6e83da027aaaed8c Mon Sep 17 00:00:00 2001 From: Ioannis Canellos Date: Fri, 9 Sep 2016 12:52:44 +0300 Subject: [PATCH] #54: Fix handling of endpoints with multiple addresses. --- pom.xml | 11 +++ spring-cloud-kubernetes-discovery/pom.xml | 32 +++++++ .../discovery/KubernetesDiscoveryClient.java | 8 +- .../discovery/KubernetesServiceInstance.java | 31 ++----- .../KubernetesDiscoveryClientTest.groovy | 84 +++++++++++++++++++ .../ZipkinKubernetesAutoConfiguration.java | 2 +- 6 files changed, 142 insertions(+), 26 deletions(-) create mode 100644 spring-cloud-kubernetes-discovery/src/test/groovy/io/fabric8/spring/cloud/discovery/KubernetesDiscoveryClientTest.groovy diff --git a/pom.xml b/pom.xml index 2bb46d8a..80cbb826 100644 --- a/pom.xml +++ b/pom.xml @@ -89,6 +89,7 @@ 3.5 + 2.19.1 1.2 0.3.4 @@ -262,6 +263,16 @@ + + org.apache.maven.plugins + maven-surefire-plugin + ${maven-surefire-plugin.version} + true + + 1 + false + + io.sundr sundr-maven-plugin diff --git a/spring-cloud-kubernetes-discovery/pom.xml b/spring-cloud-kubernetes-discovery/pom.xml index 0f003c3e..9a7ddfea 100644 --- a/spring-cloud-kubernetes-discovery/pom.xml +++ b/spring-cloud-kubernetes-discovery/pom.xml @@ -51,6 +51,38 @@ spring-cloud-context true + + + + org.springframework.boot + spring-boot-starter-test + test + + + + org.springframework.boot + spring-boot-starter-web + test + + + + io.fabric8 + kubernetes-client + test-jar + test + + + + io.fabric8 + mockwebserver + test + + + + org.spockframework + spock-spring + test + \ No newline at end of file diff --git a/spring-cloud-kubernetes-discovery/src/main/java/io/fabric8/spring/cloud/discovery/KubernetesDiscoveryClient.java b/spring-cloud-kubernetes-discovery/src/main/java/io/fabric8/spring/cloud/discovery/KubernetesDiscoveryClient.java index 16b317eb..e8fbfcef 100644 --- a/spring-cloud-kubernetes-discovery/src/main/java/io/fabric8/spring/cloud/discovery/KubernetesDiscoveryClient.java +++ b/spring-cloud-kubernetes-discovery/src/main/java/io/fabric8/spring/cloud/discovery/KubernetesDiscoveryClient.java @@ -67,7 +67,10 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { return endpoints.getSubsets() .stream() .filter(s -> s.getAddresses().get(0).getTargetRef().getName().equals(podName)) - .map(s -> (ServiceInstance) new KubernetesServiceInstance(serviceName, s.getPorts().iterator().next().getName(), s, false)) + .map(s -> (ServiceInstance) new KubernetesServiceInstance(serviceName, + s.getAddresses().stream().findFirst().orElseThrow(IllegalStateException::new), + s.getPorts().stream().findFirst().orElseThrow(IllegalStateException::new), + false)) .findFirst().orElse(defaultInstance); } catch (Throwable t) { return defaultInstance; @@ -79,7 +82,8 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { Assert.notNull(serviceId, "[Assertion failed] - the object argument must be null"); return Optional.ofNullable(client.endpoints().withName(serviceId).get()).orElse(new Endpoints()) .getSubsets() - .stream().map(s -> new KubernetesServiceInstance(serviceId, s.getPorts().iterator().next().getName(), s, false)) + .stream() + .flatMap(s -> s.getAddresses().stream().map(a -> (ServiceInstance) new KubernetesServiceInstance(serviceId, a ,s.getPorts().stream().findFirst().orElseThrow(IllegalStateException::new), false))) .collect(Collectors.toList()); } diff --git a/spring-cloud-kubernetes-discovery/src/main/java/io/fabric8/spring/cloud/discovery/KubernetesServiceInstance.java b/spring-cloud-kubernetes-discovery/src/main/java/io/fabric8/spring/cloud/discovery/KubernetesServiceInstance.java index 95dfe36c..0ee93efd 100644 --- a/spring-cloud-kubernetes-discovery/src/main/java/io/fabric8/spring/cloud/discovery/KubernetesServiceInstance.java +++ b/spring-cloud-kubernetes-discovery/src/main/java/io/fabric8/spring/cloud/discovery/KubernetesServiceInstance.java @@ -16,14 +16,13 @@ package io.fabric8.spring.cloud.discovery; +import io.fabric8.kubernetes.api.model.EndpointAddress; import io.fabric8.kubernetes.api.model.EndpointPort; -import io.fabric8.kubernetes.api.model.EndpointSubset; import org.springframework.cloud.client.ServiceInstance; import java.net.URI; import java.net.URISyntaxException; import java.util.Collections; -import java.util.HashMap; import java.util.Map; import static io.fabric8.kubernetes.client.utils.Utils.isNotNullOrEmpty; @@ -36,14 +35,14 @@ public class KubernetesServiceInstance implements ServiceInstance { private static final String COLN = ":"; private final String serviceId; - private final String portName; - private final EndpointSubset subset; + private final EndpointAddress endpointAddress; + private final EndpointPort endpointPort; private final Boolean secure; - public KubernetesServiceInstance(String serviceId, String portName, EndpointSubset subset, Boolean secure) { + public KubernetesServiceInstance(String serviceId, EndpointAddress endpointAddress, EndpointPort endpointPort, Boolean secure) { this.serviceId = serviceId; - this.portName = portName; - this.subset = subset; + this.endpointAddress = endpointAddress; + this.endpointPort = endpointPort; this.secure = secure; } @@ -54,26 +53,12 @@ public class KubernetesServiceInstance implements ServiceInstance { @Override public String getHost() { - if (subset.getAddresses().isEmpty()) { - throw new IllegalStateException("Endpoint subset has no addresses."); - } - return subset.getAddresses().get(0).getIp(); + return endpointAddress.getIp(); } @Override public int getPort() { - if (subset.getPorts().isEmpty()) { - throw new IllegalStateException("Endpoint subset has no ports."); - } else if (isNullOrEmpty(portName) && subset.getPorts().size() == 1) { - return subset.getPorts().get(0).getPort(); - } else if (isNotNullOrEmpty(portName)) { - for (EndpointPort port : subset.getPorts()) { - if (portName.endsWith(port.getName())) { - return port.getPort(); - } - } - } - throw new IllegalStateException("Endpoint subset has no matching ports."); + return endpointPort.getPort(); } @Override diff --git a/spring-cloud-kubernetes-discovery/src/test/groovy/io/fabric8/spring/cloud/discovery/KubernetesDiscoveryClientTest.groovy b/spring-cloud-kubernetes-discovery/src/test/groovy/io/fabric8/spring/cloud/discovery/KubernetesDiscoveryClientTest.groovy new file mode 100644 index 00000000..d89f108f --- /dev/null +++ b/spring-cloud-kubernetes-discovery/src/test/groovy/io/fabric8/spring/cloud/discovery/KubernetesDiscoveryClientTest.groovy @@ -0,0 +1,84 @@ +package io.fabric8.spring.cloud.discovery + +import io.fabric8.kubernetes.api.model.EndpointsBuilder +import io.fabric8.kubernetes.client.Config +import io.fabric8.kubernetes.client.KubernetesClient +import io.fabric8.kubernetes.server.mock.KubernetesMockServer +import org.springframework.cloud.client.ServiceInstance +import org.springframework.cloud.client.discovery.DiscoveryClient +import spock.lang.Specification + +class KubernetesDiscoveryClientTest extends Specification { + + private static KubernetesMockServer mockServer = new KubernetesMockServer() + private static KubernetesClient mockClient + + + def setupSpec() { + mockServer.init() + mockClient = mockServer.createClient() + + //Configure the kubernetes master url to point to the mock server + System.setProperty(Config.KUBERNETES_MASTER_SYSTEM_PROPERTY, mockClient.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") + } + + def cleanupSpec() { + mockServer.destroy(); + } + + def "Should be able to handle endpoints single address"() { + given: + mockServer.expect().get().withPath("/api/v1/namespaces/test/endpoints/endpoint").andReturn(200, new EndpointsBuilder() + .withNewMetadata() + .withName("endpoint") + .endMetadata() + .addNewSubset() + .addNewAddress() + .withIp("ip1") + .endAddress() + .addNewPort("http",80,"TCP") + .endSubset() + .build()).once() + + DiscoveryClient discoveryClient = new KubernetesDiscoveryClient(mockClient, new KubernetesDiscoveryProperties()) + when: + List instances = discoveryClient.getInstances("endpoint") + then: + instances != null + instances.size() == 1 + instances.find({s -> s.host == "ip1"}) + } + + + + def "Should be able to handle endpoints multiple addresses"() { + given: + mockServer.expect().get().withPath("/api/v1/namespaces/test/endpoints/endpoint").andReturn(200, new EndpointsBuilder() + .withNewMetadata() + .withName("endpoint") + .endMetadata() + .addNewSubset() + .addNewAddress() + .withIp("ip1") + .endAddress() + .addNewAddress() + .withIp("ip2") + .endAddress() + .addNewPort("http",80,"TCP") + .endSubset() + .build()).once() + + DiscoveryClient discoveryClient = new KubernetesDiscoveryClient(mockClient, new KubernetesDiscoveryProperties()) + when: + List instances = discoveryClient.getInstances("endpoint") + then: + instances != null + instances.size() == 2 + instances.find({s -> s.host == "ip1"}) + instances.find({s -> s.host == "ip2"}) + + } +} diff --git a/spring-cloud-kubernetes-zipkin/src/main/java/io/fabric8/spring/cloud/kubernetes/zipkin/ZipkinKubernetesAutoConfiguration.java b/spring-cloud-kubernetes-zipkin/src/main/java/io/fabric8/spring/cloud/kubernetes/zipkin/ZipkinKubernetesAutoConfiguration.java index 8947ff1b..2210d610 100644 --- a/spring-cloud-kubernetes-zipkin/src/main/java/io/fabric8/spring/cloud/kubernetes/zipkin/ZipkinKubernetesAutoConfiguration.java +++ b/spring-cloud-kubernetes-zipkin/src/main/java/io/fabric8/spring/cloud/kubernetes/zipkin/ZipkinKubernetesAutoConfiguration.java @@ -72,7 +72,7 @@ public class ZipkinKubernetesAutoConfiguration { .orElse(new Endpoints()) .getSubsets() .stream() - .map(s -> new KubernetesServiceInstance(name, s.getPorts().iterator().next().getName(), s, false)) + .flatMap(s -> s.getAddresses().stream().map(a -> (ServiceInstance) new KubernetesServiceInstance(name, a ,s.getPorts().stream().findFirst().orElseThrow(IllegalStateException::new), false))) .collect(Collectors.toList()); }