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());
}