From fda256fc7eff8b2fb22eaeda5b8b961092a73579 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jakub=20Kubry=C5=84ski?= Date: Thu, 28 Feb 2019 16:53:53 +0100 Subject: [PATCH] Fix #341 - allow defining a primary port name when multiple ports are defined for a service (#342) --- .../discovery/KubernetesDiscoveryClient.java | 25 ++++++++++- .../KubernetesDiscoveryProperties.java | 14 ++++++ .../KubernetesDiscoveryClientTest.java | 43 ++++++++++++++++--- 3 files changed, 74 insertions(+), 8 deletions(-) diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClient.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClient.java index 60e766e5..44a286a5 100644 --- a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClient.java +++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClient.java @@ -153,8 +153,8 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { if (endpointAddress.getTargetRef() != null) { instanceId = endpointAddress.getTargetRef().getUid(); } - final EndpointPort endpointPort = s.getPorts().stream().findFirst() - .orElseThrow(IllegalStateException::new); + + EndpointPort endpointPort = findEndpointPort(s); instances.add(new KubernetesServiceInstance(instanceId, serviceId, endpointAddress, endpointPort, endpointMetadata, this.isServicePortSecureResolver @@ -170,6 +170,27 @@ public class KubernetesDiscoveryClient implements DiscoveryClient { return instances; } + private EndpointPort findEndpointPort(EndpointSubset s) { + List ports = s.getPorts(); + EndpointPort endpointPort; + if (ports.size() == 1) { + endpointPort = ports.get(0); + } + else { + Predicate portPredicate; + if (!StringUtils.isEmpty(properties.getPrimaryPortName())) { + portPredicate = port -> properties.getPrimaryPortName() + .equalsIgnoreCase(port.getName()); + } + else { + portPredicate = port -> true; + } + endpointPort = ports.stream().filter(portPredicate).findAny() + .orElseThrow(IllegalStateException::new); + } + return endpointPort; + } + private List getSubsetsFromEndpoints(Endpoints endpoints) { if (endpoints == null) { return new ArrayList<>(); diff --git a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java index 06ff845d..00e65ecd 100644 --- a/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java +++ b/spring-cloud-kubernetes-discovery/src/main/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryProperties.java @@ -60,6 +60,12 @@ public class KubernetesDiscoveryProperties { */ private Map serviceLabels = new HashMap<>(); + /** + * If set then the port with a given name is used as primary when multiple ports are + * defined for a service. + */ + private String primaryPortName; + private Metadata metadata = new Metadata(); public boolean isEnabled() { @@ -102,6 +108,14 @@ public class KubernetesDiscoveryProperties { this.serviceLabels = serviceLabels; } + public String getPrimaryPortName() { + return primaryPortName; + } + + public void setPrimaryPortName(String primaryPortName) { + this.primaryPortName = primaryPortName; + } + public Metadata getMetadata() { return this.metadata; } diff --git a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientTest.java b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientTest.java index e2c17cbc..1e9b2604 100644 --- a/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientTest.java +++ b/spring-cloud-kubernetes-discovery/src/test/java/org/springframework/cloud/kubernetes/discovery/KubernetesDiscoveryClientTest.java @@ -57,12 +57,10 @@ public class KubernetesDiscoveryClientTest { @Test public void getInstancesShouldBeAbleToHandleEndpointsSingleAddress() { mockServer.expect().get().withPath("/api/v1/namespaces/test/endpoints/endpoint") - .andReturn(200, - new EndpointsBuilder().withNewMetadata().withName("endpoint") - .endMetadata().addNewSubset().addNewAddress() - .withIp("ip1").withNewTargetRef().withUid("uid").endTargetRef().endAddress() - .addNewPort("http", 80, "TCP") - .endSubset().build()) + .andReturn(200, new EndpointsBuilder().withNewMetadata() + .withName("endpoint").endMetadata().addNewSubset().addNewAddress() + .withIp("ip1").withNewTargetRef().withUid("uid").endTargetRef() + .endAddress().addNewPort("http", 80, "TCP").endSubset().build()) .once(); mockServer.expect().get().withPath("/api/v1/namespaces/test/services/endpoint") @@ -86,6 +84,39 @@ public class KubernetesDiscoveryClientTest { .filteredOn(s -> s.getInstanceId().equals("uid")).hasSize(1); } + @Test + public void getInstancesShouldBeAbleToHandleEndpointsSingleAddressAndMultiplePorts() { + mockServer.expect().get().withPath("/api/v1/namespaces/test/endpoints/endpoint") + .andReturn(200, new EndpointsBuilder().withNewMetadata() + .withName("endpoint").endMetadata().addNewSubset().addNewAddress() + .withIp("ip1").withNewTargetRef().withUid("uid").endTargetRef() + .endAddress().addNewPort("mgmt", 9000, "TCP") + .addNewPort("http", 80, "TCP").endSubset().build()) + .once(); + + mockServer.expect().get().withPath("/api/v1/namespaces/test/services/endpoint") + .andReturn(200, new ServiceBuilder().withNewMetadata() + .withName("endpoint").withLabels(new HashMap() { + { + put("l", "v"); + } + }).endMetadata().build()) + .once(); + + final KubernetesDiscoveryProperties properties = new KubernetesDiscoveryProperties(); + properties.setPrimaryPortName("http"); + final DiscoveryClient discoveryClient = new KubernetesDiscoveryClient(mockClient, + properties, KubernetesClient::services, + new DefaultIsServicePortSecureResolver(properties)); + + final List instances = discoveryClient.getInstances("endpoint"); + + assertThat(instances).hasSize(1) + .filteredOn(s -> s.getHost().equals("ip1") && !s.isSecure()).hasSize(1) + .filteredOn(s -> s.getInstanceId().equals("uid")).hasSize(1) + .filteredOn(s -> 80 == s.getPort()).hasSize(1); + } + @Test public void getInstancesShouldBeAbleToHandleEndpointsMultipleAddresses() { mockServer.expect().get().withPath("/api/v1/namespaces/test/endpoints/endpoint")