Cleanup k8s client loadbalancer part 1 (#1612)

This commit is contained in:
erabii
2024-04-03 18:47:35 +03:00
committed by GitHub
parent 67f89c6a46
commit 17471a1ed6
4 changed files with 224 additions and 82 deletions

View File

@@ -22,7 +22,6 @@ import java.util.List;
import io.kubernetes.client.openapi.ApiException;
import io.kubernetes.client.openapi.apis.CoreV1Api;
import io.kubernetes.client.openapi.models.V1Service;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import reactor.core.publisher.Flux;
@@ -32,47 +31,88 @@ import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscover
import org.springframework.cloud.kubernetes.commons.loadbalancer.KubernetesServiceInstanceMapper;
import org.springframework.cloud.kubernetes.commons.loadbalancer.KubernetesServicesListSupplier;
import org.springframework.core.env.Environment;
import org.springframework.core.log.LogAccessor;
import static org.springframework.cloud.kubernetes.client.KubernetesClientUtils.getApplicationNamespace;
/**
* @author Ryan Baxter
*/
public class KubernetesClientServicesListSupplier extends KubernetesServicesListSupplier<V1Service> {
private static final Log LOG = LogFactory.getLog(KubernetesClientServicesListSupplier.class);
private static final LogAccessor LOG = new LogAccessor(
LogFactory.getLog(KubernetesClientServicesListSupplier.class));
private final CoreV1Api coreV1Api;
private final String namespace;
private final KubernetesNamespaceProvider kubernetesNamespaceProvider;
public KubernetesClientServicesListSupplier(Environment environment,
KubernetesServiceInstanceMapper<V1Service> mapper, KubernetesDiscoveryProperties discoveryProperties,
CoreV1Api coreV1Api, KubernetesNamespaceProvider kubernetesNamespaceProvider) {
super(environment, mapper, discoveryProperties);
this.coreV1Api = coreV1Api;
this.namespace = kubernetesNamespaceProvider.getNamespace();
this.kubernetesNamespaceProvider = kubernetesNamespaceProvider;
}
@Override
public Flux<List<ServiceInstance>> get() {
LOG.info("Getting services with id " + this.getServiceId());
List<ServiceInstance> result = new ArrayList<>();
List<V1Service> services;
try {
if (discoveryProperties.allNamespaces()) {
services = coreV1Api.listServiceForAllNamespaces(null, null, "metadata.name=" + this.getServiceId(),
null, null, null, null, null, null, null, null).getItems();
}
else {
services = coreV1Api.listNamespacedService(namespace, null, null, null,
"metadata.name=" + this.getServiceId(), null, null, null, null, null, null, null).getItems();
}
services.forEach(service -> result.add(mapper.map(service)));
String serviceName = getServiceId();
LOG.debug(() -> "serviceID : " + serviceName);
if (discoveryProperties.allNamespaces()) {
LOG.debug(() -> "discovering services in all namespaces");
List<V1Service> services = services(null, serviceName);
services.forEach(service -> addMappedService(mapper, result, service));
}
catch (ApiException e) {
LOG.warn("Error retrieving service with name " + this.getServiceId(), e);
else if (!discoveryProperties.namespaces().isEmpty()) {
List<String> selectiveNamespaces = discoveryProperties.namespaces().stream().sorted().toList();
LOG.debug(() -> "discovering services in selective namespaces : " + selectiveNamespaces);
selectiveNamespaces.forEach(selectiveNamespace -> {
List<V1Service> services = services(selectiveNamespace, serviceName);
services.forEach(service -> addMappedService(mapper, result, service));
});
}
LOG.info("Returning services: " + result);
else {
String namespace = getApplicationNamespace(null, "loadbalancer-service", kubernetesNamespaceProvider);
LOG.debug(() -> "discovering services in namespace : " + namespace);
List<V1Service> services = services(namespace, serviceName);
services.forEach(service -> addMappedService(mapper, result, service));
}
LOG.debug(() -> "found services : " + result);
return Flux.defer(() -> Flux.just(result));
}
private void addMappedService(KubernetesServiceInstanceMapper<V1Service> mapper, List<ServiceInstance> services,
V1Service service) {
services.add(mapper.map(service));
}
private List<V1Service> services(String namespace, String serviceName) {
if (namespace == null) {
try {
return coreV1Api.listServiceForAllNamespaces(null, null, "metadata.name=" + serviceName, null, null,
null, null, null, null, null, null).getItems();
}
catch (ApiException apiException) {
LOG.warn(apiException, "Error retrieving services (in all namespaces) with name " + serviceName);
return List.of();
}
}
else {
try {
// there is going to be a single service here, if found
return coreV1Api.listNamespacedService(namespace, null, null, null, "metadata.name=" + serviceName,
null, null, null, null, null, null, null).getItems();
}
catch (ApiException apiException) {
LOG.warn(apiException,
"Error retrieving service with name " + serviceName + " in namespace : " + namespace);
return List.of();
}
}
}
}

View File

@@ -79,7 +79,7 @@ class KubernetesClientServiceInstanceMapperTests {
KubernetesServiceInstance serviceInstance = mapper.map(service);
DefaultKubernetesServiceInstance result = new DefaultKubernetesServiceInstance("0", "database",
"database.default.svc.cluster.local", 443, new HashMap(), true);
"database.default.svc.cluster.local", 443, Map.of(), true);
assertThat(serviceInstance).isEqualTo(result);
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2013-2020 the original author or authors.
* Copyright 2013-2024 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.
@@ -17,18 +17,19 @@
package org.springframework.cloud.kubernetes.client.loadbalancer;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import com.github.tomakehurst.wiremock.WireMockServer;
import com.github.tomakehurst.wiremock.client.WireMock;
import com.github.tomakehurst.wiremock.common.ConsoleNotifier;
import io.kubernetes.client.openapi.ApiClient;
import io.kubernetes.client.openapi.Configuration;
import io.kubernetes.client.openapi.JSON;
import io.kubernetes.client.openapi.apis.CoreV1Api;
import io.kubernetes.client.openapi.models.V1ObjectMetaBuilder;
import io.kubernetes.client.openapi.models.V1Service;
import io.kubernetes.client.openapi.models.V1ServiceBuilder;
import io.kubernetes.client.openapi.models.V1ServiceList;
import io.kubernetes.client.openapi.models.V1ServicePortBuilder;
@@ -36,141 +37,236 @@ import io.kubernetes.client.openapi.models.V1ServiceSpecBuilder;
import io.kubernetes.client.util.ClientBuilder;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import org.springframework.boot.test.system.CapturedOutput;
import org.springframework.boot.test.system.OutputCaptureExtension;
import org.springframework.cloud.client.ServiceInstance;
import org.springframework.cloud.kubernetes.commons.KubernetesNamespaceProvider;
import org.springframework.cloud.kubernetes.commons.discovery.DefaultKubernetesServiceInstance;
import org.springframework.cloud.kubernetes.commons.discovery.KubernetesDiscoveryProperties;
import org.springframework.cloud.kubernetes.commons.loadbalancer.KubernetesLoadBalancerProperties;
import org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory;
import org.springframework.mock.env.MockEnvironment;
import static com.github.tomakehurst.wiremock.client.WireMock.aResponse;
import static com.github.tomakehurst.wiremock.client.WireMock.get;
import static com.github.tomakehurst.wiremock.client.WireMock.stubFor;
import static com.github.tomakehurst.wiremock.client.WireMock.urlMatching;
import static com.github.tomakehurst.wiremock.client.WireMock.urlEqualTo;
import static com.github.tomakehurst.wiremock.core.WireMockConfiguration.options;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import static org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory.PROPERTY_NAME;
/**
* @author Ryan Baxter
*/
@ExtendWith(OutputCaptureExtension.class)
class KubernetesClientServicesListSupplierTests {
private static final V1ServiceList SERVICE_LIST = new V1ServiceList().addItemsItem(new V1ServiceBuilder()
.withMetadata(new V1ObjectMetaBuilder().withName("service1").withNamespace("default")
.withResourceVersion("1").addToLabels("beta", "true")
.addToAnnotations("org.springframework.cloud", "true").withUid("0").build())
private static final V1Service SERVICE_A_DEFAULT_NAMESPACE = new V1ServiceBuilder()
.withMetadata(new V1ObjectMetaBuilder().withName("service-a").withNamespace("default").withUid("0")
.addToLabels("beta", "true").addToAnnotations("org.springframework.cloud", "true").build())
.withSpec(new V1ServiceSpecBuilder()
.addToPorts(new V1ServicePortBuilder().withPort(80).withName("http").build()).build())
.build());
.build();
private static final V1Service SERVICE_A_TEST_NAMESPACE = new V1ServiceBuilder()
.withMetadata(new V1ObjectMetaBuilder().withName("service-a").withNamespace("test").withUid("1").build())
.withSpec(new V1ServiceSpecBuilder()
.addToPorts(new V1ServicePortBuilder().withPort(80).withName("http").build(),
new V1ServicePortBuilder().withPort(443).withName("https").build())
.build())
.build();
private static final V1ServiceList SINGLE_NAMESPACE_SERVICES = new V1ServiceList()
.addItemsItem(SERVICE_A_DEFAULT_NAMESPACE);
private static final V1ServiceList SERVICE_LIST_ALL_NAMESPACE = new V1ServiceList()
.addItemsItem(new V1ServiceBuilder()
.withMetadata(new V1ObjectMetaBuilder().withName("service1").withNamespace("default")
.withResourceVersion("1").addToLabels("beta", "true")
.addToAnnotations("org.springframework.cloud", "true").withUid("0").build())
.withSpec(new V1ServiceSpecBuilder()
.addToPorts(new V1ServicePortBuilder().withPort(80).withName("http").build()).build())
.build())
.addItemsItem(new V1ServiceBuilder()
.withMetadata(new V1ObjectMetaBuilder().withName("service1").withNamespace("test")
.withResourceVersion("1").withUid("1").build())
.withSpec(new V1ServiceSpecBuilder()
.addToPorts(new V1ServicePortBuilder().withPort(80).withName("http").build(),
new V1ServicePortBuilder().withPort(443).withName("https").build())
.build())
.build());
.addItemsItem(SERVICE_A_DEFAULT_NAMESPACE).addItemsItem(SERVICE_A_TEST_NAMESPACE);
private static final V1ServiceList SERVICE_A_DEFAULT_NAMESPACE_SELECTIVE_NAMESPACES = new V1ServiceList()
.addItemsItem(SERVICE_A_DEFAULT_NAMESPACE);
private static final V1ServiceList SERVICE_A_TEST_NAMESPACE_SELECTIVE_NAMESPACES = new V1ServiceList()
.addItemsItem(SERVICE_A_TEST_NAMESPACE);
private static WireMockServer wireMockServer;
@BeforeAll
public static void setup() {
wireMockServer = new WireMockServer(options().dynamicPort());
static void setup() {
wireMockServer = new WireMockServer(options().dynamicPort().notifier(new ConsoleNotifier(true)));
wireMockServer.start();
WireMock.configureFor("localhost", wireMockServer.port());
ApiClient client = new ClientBuilder().setBasePath("http://localhost:" + wireMockServer.port()).build();
client.setDebugging(true);
Configuration.setDefaultApiClient(client);
}
@AfterAll
public static void after() {
static void after() {
wireMockServer.stop();
}
@AfterEach
public void afterEach() {
void afterEach() {
WireMock.reset();
}
@Test
void getList() {
MockEnvironment env = new MockEnvironment();
env.setProperty(LoadBalancerClientFactory.PROPERTY_NAME, "service1");
void singleNamespaceTest(CapturedOutput output) {
MockEnvironment env = new MockEnvironment().withProperty(PROPERTY_NAME, "service-a");
KubernetesNamespaceProvider kubernetesNamespaceProvider = mock(KubernetesNamespaceProvider.class);
when(kubernetesNamespaceProvider.getNamespace()).thenReturn("default");
CoreV1Api coreV1Api = new CoreV1Api();
KubernetesClientServiceInstanceMapper mapper = new KubernetesClientServiceInstanceMapper(
new KubernetesLoadBalancerProperties(), KubernetesDiscoveryProperties.DEFAULT);
KubernetesClientServicesListSupplier listSupplier = new KubernetesClientServicesListSupplier(env, mapper,
KubernetesDiscoveryProperties.DEFAULT, coreV1Api, kubernetesNamespaceProvider);
stubFor(get(urlMatching("^/api/v1/namespaces/default/services.*"))
.willReturn(aResponse().withStatus(200).withBody(new JSON().serialize(SERVICE_LIST))));
boolean allNamespaces = false;
Set<String> selectiveNamespaces = Set.of();
KubernetesDiscoveryProperties discoveryProperties = new KubernetesDiscoveryProperties(true, allNamespaces,
selectiveNamespaces, true, 60, false, null, Set.of(443, 8443, 12345), Map.of(), null,
KubernetesDiscoveryProperties.Metadata.DEFAULT, 0, true);
KubernetesClientServicesListSupplier listSupplier = new KubernetesClientServicesListSupplier(env, mapper,
discoveryProperties, coreV1Api, kubernetesNamespaceProvider);
stubFor(get(urlEqualTo("/api/v1/namespaces/default/services?fieldSelector=metadata.name%3Dservice-a"))
.willReturn(aResponse().withStatus(200).withBody(new JSON().serialize(SINGLE_NAMESPACE_SERVICES))));
Flux<List<ServiceInstance>> instances = listSupplier.get();
Map<String, String> metadata = new HashMap<>();
metadata.put("org.springframework.cloud", "true");
metadata.put("beta", "true");
DefaultKubernetesServiceInstance service1 = new DefaultKubernetesServiceInstance("0", "service1",
"service1.default.svc.cluster.local", 80, metadata, false);
Map<String, String> metadata = Map.of("org.springframework.cloud", "true", "beta", "true");
DefaultKubernetesServiceInstance serviceA = new DefaultKubernetesServiceInstance("0", "service-a",
"service-a.default.svc.cluster.local", 80, metadata, false);
List<ServiceInstance> services = new ArrayList<>();
services.add(service1);
services.add(serviceA);
StepVerifier.create(instances).expectNext(services).verifyComplete();
Assertions.assertTrue(output.getOut().contains("serviceID : service-a"));
Assertions.assertTrue(output.getOut().contains("discovering services in namespace : default"));
}
@Test
void getListAllNamespaces() {
MockEnvironment env = new MockEnvironment();
env.setProperty(LoadBalancerClientFactory.PROPERTY_NAME, "service1");
void singleNamespaceNoServicePresentTest(CapturedOutput output) {
MockEnvironment env = new MockEnvironment().withProperty(PROPERTY_NAME, "service-a");
KubernetesNamespaceProvider kubernetesNamespaceProvider = mock(KubernetesNamespaceProvider.class);
when(kubernetesNamespaceProvider.getNamespace()).thenReturn("default");
KubernetesDiscoveryProperties kubernetesDiscoveryProperties = new KubernetesDiscoveryProperties(true, true,
Set.of(), true, 60, false, null, Set.of(), Map.of(), null,
KubernetesDiscoveryProperties.Metadata.DEFAULT, 0, false);
CoreV1Api coreV1Api = new CoreV1Api();
KubernetesClientServiceInstanceMapper mapper = new KubernetesClientServiceInstanceMapper(
new KubernetesLoadBalancerProperties(), kubernetesDiscoveryProperties);
KubernetesClientServicesListSupplier listSupplier = new KubernetesClientServicesListSupplier(env, mapper,
kubernetesDiscoveryProperties, coreV1Api, kubernetesNamespaceProvider);
new KubernetesLoadBalancerProperties(), KubernetesDiscoveryProperties.DEFAULT);
stubFor(get(urlMatching("^/api/v1/services.*"))
boolean allNamespaces = false;
Set<String> selectiveNamespaces = Set.of();
KubernetesDiscoveryProperties discoveryProperties = new KubernetesDiscoveryProperties(true, allNamespaces,
selectiveNamespaces, true, 60, false, null, Set.of(443, 8443, 12345), Map.of(), null,
KubernetesDiscoveryProperties.Metadata.DEFAULT, 0, true);
KubernetesClientServicesListSupplier listSupplier = new KubernetesClientServicesListSupplier(env, mapper,
discoveryProperties, coreV1Api, kubernetesNamespaceProvider);
stubFor(get(urlEqualTo("/api/v1/namespaces/default/services?fieldSelector=metadata.name%3Dservice-a"))
.willReturn(aResponse().withStatus(404)));
Flux<List<ServiceInstance>> instances = listSupplier.get();
List<ServiceInstance> services = List.of();
StepVerifier.create(instances).expectNext(services).verifyComplete();
Assertions.assertTrue(output.getOut().contains("serviceID : service-a"));
Assertions.assertTrue(output.getOut().contains("discovering services in namespace : default"));
Assertions.assertTrue(output.getOut().contains("Error retrieving service with name service-a"));
}
@Test
void allNamespacesTest(CapturedOutput output) {
MockEnvironment env = new MockEnvironment().withProperty(PROPERTY_NAME, "service-a");
KubernetesNamespaceProvider kubernetesNamespaceProvider = mock(KubernetesNamespaceProvider.class);
when(kubernetesNamespaceProvider.getNamespace()).thenReturn("default");
boolean allNamespaces = true;
Set<String> selectiveNamespaces = Set.of();
KubernetesDiscoveryProperties discoveryProperties = new KubernetesDiscoveryProperties(true, allNamespaces,
selectiveNamespaces, true, 60, false, null, Set.of(443, 8443, 12345), Map.of(), null,
KubernetesDiscoveryProperties.Metadata.DEFAULT, 0, true);
CoreV1Api coreV1Api = new CoreV1Api();
KubernetesClientServiceInstanceMapper mapper = new KubernetesClientServiceInstanceMapper(
new KubernetesLoadBalancerProperties(), discoveryProperties);
KubernetesClientServicesListSupplier listSupplier = new KubernetesClientServicesListSupplier(env, mapper,
discoveryProperties, coreV1Api, kubernetesNamespaceProvider);
stubFor(get(urlEqualTo("/api/v1/services?fieldSelector=metadata.name%3Dservice-a"))
.willReturn(aResponse().withStatus(200).withBody(new JSON().serialize(SERVICE_LIST_ALL_NAMESPACE))));
Flux<List<ServiceInstance>> instances = listSupplier.get();
Map<String, String> metadata = new HashMap<>();
metadata.put("org.springframework.cloud", "true");
metadata.put("beta", "true");
DefaultKubernetesServiceInstance service1 = new DefaultKubernetesServiceInstance("0", "service1",
"service1.default.svc.cluster.local", 80, metadata, false);
DefaultKubernetesServiceInstance service2 = new DefaultKubernetesServiceInstance("1", "service1",
"service1.test.svc.cluster.local", 80, new HashMap<>(), false);
Map<String, String> metadata = Map.of("org.springframework.cloud", "true", "beta", "true");
DefaultKubernetesServiceInstance serviceADefaultNamespace = new DefaultKubernetesServiceInstance("0",
"service-a", "service-a.default.svc.cluster.local", 80, metadata, false);
DefaultKubernetesServiceInstance serviceATestNamespace = new DefaultKubernetesServiceInstance("1", "service-a",
"service-a.test.svc.cluster.local", 80, Map.of(), false);
List<ServiceInstance> services = new ArrayList<>();
services.add(service1);
services.add(service2);
services.add(serviceADefaultNamespace);
services.add(serviceATestNamespace);
StepVerifier.create(instances).expectNext(services).verifyComplete();
Assertions.assertTrue(output.getOut().contains("discovering services in all namespaces"));
}
@Test
void selectiveNamespacesTest(CapturedOutput output) {
MockEnvironment env = new MockEnvironment().withProperty(PROPERTY_NAME, "service-a");
KubernetesNamespaceProvider kubernetesNamespaceProvider = mock(KubernetesNamespaceProvider.class);
boolean allNamespaces = false;
Set<String> selectiveNamespaces = Set.of("default", "test", "no-service");
KubernetesDiscoveryProperties discoveryProperties = new KubernetesDiscoveryProperties(true, allNamespaces,
selectiveNamespaces, true, 60, false, null, Set.of(443, 8443, 12345), Map.of(), null,
KubernetesDiscoveryProperties.Metadata.DEFAULT, 0, true);
CoreV1Api coreV1Api = new CoreV1Api();
KubernetesClientServiceInstanceMapper mapper = new KubernetesClientServiceInstanceMapper(
new KubernetesLoadBalancerProperties(), discoveryProperties);
KubernetesClientServicesListSupplier listSupplier = new KubernetesClientServicesListSupplier(env, mapper,
discoveryProperties, coreV1Api, kubernetesNamespaceProvider);
stubFor(get(urlEqualTo("/api/v1/namespaces/default/services?fieldSelector=metadata.name%3Dservice-a"))
.willReturn(aResponse().withStatus(200)
.withBody(new JSON().serialize(SERVICE_A_DEFAULT_NAMESPACE_SELECTIVE_NAMESPACES))));
stubFor(get(urlEqualTo("/api/v1/namespaces/test/services?fieldSelector=metadata.name%3Dservice-a"))
.willReturn(aResponse().withStatus(200)
.withBody(new JSON().serialize(SERVICE_A_TEST_NAMESPACE_SELECTIVE_NAMESPACES))));
stubFor(get(urlEqualTo("/api/v1/namespaces/no-service/services?fieldSelector=metadata.name%3Dservice-a"))
.willReturn(aResponse().withStatus(404)));
Flux<List<ServiceInstance>> instances = listSupplier.get();
Map<String, String> metadata = Map.of("org.springframework.cloud", "true", "beta", "true");
DefaultKubernetesServiceInstance serviceADefaultNamespace = new DefaultKubernetesServiceInstance("0",
"service-a", "service-a.default.svc.cluster.local", 80, metadata, false);
DefaultKubernetesServiceInstance serviceATestNamespace = new DefaultKubernetesServiceInstance("1", "service-a",
"service-a.test.svc.cluster.local", 80, Map.of(), false);
List<ServiceInstance> services = new ArrayList<>();
services.add(serviceADefaultNamespace);
services.add(serviceATestNamespace);
StepVerifier.create(instances).expectNext(services).verifyComplete();
Assertions.assertTrue(
output.getOut().contains("Error retrieving service with name service-a in namespace : no-service"));
Assertions.assertTrue(
output.getOut().contains("discovering services in selective namespaces : [default, no-service, test]"));
}
}

View File

@@ -0,0 +1,6 @@
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<include resource="org/springframework/boot/logging/logback/base.xml"/>
<!-- needed for CapturedOutput -->
<logger name="org.springframework.cloud.kubernetes.client.loadbalancer" level="DEBUG"/>
</configuration>