Updates for merging into main
This commit is contained in:
@@ -72,9 +72,9 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
|
||||
private final String namespace;
|
||||
|
||||
public KubernetesInformerDiscoveryClient(String namespace, SharedInformerFactory sharedInformerFactory,
|
||||
Lister<V1Service> serviceLister, Lister<V1Endpoints> endpointsLister,
|
||||
SharedInformer<V1Service> serviceInformer, SharedInformer<V1Endpoints> endpointsInformer,
|
||||
KubernetesDiscoveryProperties properties) {
|
||||
Lister<V1Service> serviceLister, Lister<V1Endpoints> endpointsLister,
|
||||
SharedInformer<V1Service> serviceInformer, SharedInformer<V1Endpoints> endpointsInformer,
|
||||
KubernetesDiscoveryProperties properties) {
|
||||
this.namespace = namespace;
|
||||
this.sharedInformerFactory = sharedInformerFactory;
|
||||
|
||||
@@ -99,8 +99,8 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
|
||||
}
|
||||
|
||||
V1Service service = properties.isAllNamespaces() ? this.serviceLister.list().stream()
|
||||
.filter(svc -> serviceId.equals(svc.getMetadata().getName())).findFirst().orElse(null)
|
||||
: this.serviceLister.namespace(this.namespace).get(serviceId);
|
||||
.filter(svc -> serviceId.equals(svc.getMetadata().getName())).findFirst().orElse(null)
|
||||
: this.serviceLister.namespace(this.namespace).get(serviceId);
|
||||
if (service == null || !matchServiceLabels(service)) {
|
||||
// no such service present in the cluster
|
||||
return new ArrayList<>();
|
||||
@@ -111,25 +111,25 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
|
||||
if (this.properties.getMetadata().addLabels()) {
|
||||
if (service.getMetadata() != null && service.getMetadata().getLabels() != null) {
|
||||
String labelPrefix = this.properties.getMetadata().labelsPrefix() != null
|
||||
? this.properties.getMetadata().labelsPrefix() : "";
|
||||
? this.properties.getMetadata().labelsPrefix() : "";
|
||||
service.getMetadata().getLabels().entrySet().stream()
|
||||
.filter(e -> e.getKey().startsWith(labelPrefix))
|
||||
.forEach(e -> svcMetadata.put(e.getKey(), e.getValue()));
|
||||
.filter(e -> e.getKey().startsWith(labelPrefix))
|
||||
.forEach(e -> svcMetadata.put(e.getKey(), e.getValue()));
|
||||
}
|
||||
}
|
||||
if (this.properties.getMetadata().addAnnotations()) {
|
||||
if (service.getMetadata() != null && service.getMetadata().getAnnotations() != null) {
|
||||
String annotationPrefix = this.properties.getMetadata().annotationsPrefix() != null
|
||||
? this.properties.getMetadata().annotationsPrefix() : "";
|
||||
? this.properties.getMetadata().annotationsPrefix() : "";
|
||||
service.getMetadata().getAnnotations().entrySet().stream()
|
||||
.filter(e -> e.getKey().startsWith(annotationPrefix))
|
||||
.forEach(e -> svcMetadata.put(e.getKey(), e.getValue()));
|
||||
.filter(e -> e.getKey().startsWith(annotationPrefix))
|
||||
.forEach(e -> svcMetadata.put(e.getKey(), e.getValue()));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
V1Endpoints ep = this.endpointsLister.namespace(service.getMetadata().getNamespace())
|
||||
.get(service.getMetadata().getName());
|
||||
.get(service.getMetadata().getName());
|
||||
if (ep == null || ep.getSubsets() == null) {
|
||||
// no available endpoints in the cluster
|
||||
return new ArrayList<>();
|
||||
@@ -138,35 +138,35 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
|
||||
Optional<String> discoveredPrimaryPortName = Optional.empty();
|
||||
if (service.getMetadata() != null && service.getMetadata().getLabels() != null) {
|
||||
discoveredPrimaryPortName = Optional
|
||||
.ofNullable(service.getMetadata().getLabels().get(PRIMARY_PORT_NAME_LABEL_KEY));
|
||||
.ofNullable(service.getMetadata().getLabels().get(PRIMARY_PORT_NAME_LABEL_KEY));
|
||||
}
|
||||
final String primaryPortName = discoveredPrimaryPortName.orElse(this.properties.getPrimaryPortName());
|
||||
|
||||
return ep.getSubsets().stream()
|
||||
.filter(subset -> subset.getPorts() != null && subset.getPorts().size() > 0) // safeguard
|
||||
.flatMap(subset -> {
|
||||
Map<String, String> metadata = new HashMap<>(svcMetadata);
|
||||
List<V1EndpointPort> endpointPorts = subset.getPorts();
|
||||
if (this.properties.getMetadata() != null && this.properties.getMetadata().addPorts()) {
|
||||
endpointPorts.forEach(p -> metadata.put(StringUtils.hasText(p.getName()) ? p.getName() : UNSET_PORT_NAME,
|
||||
Integer.toString(p.getPort())));
|
||||
}
|
||||
List<V1EndpointAddress> addresses = subset.getAddresses();
|
||||
if (addresses == null) {
|
||||
addresses = new ArrayList<>();
|
||||
}
|
||||
if (this.properties.isIncludeNotReadyAddresses()
|
||||
&& !CollectionUtils.isEmpty(subset.getNotReadyAddresses())) {
|
||||
addresses.addAll(subset.getNotReadyAddresses());
|
||||
}
|
||||
return ep.getSubsets().stream().filter(subset -> subset.getPorts() != null && subset.getPorts().size() > 0) // safeguard
|
||||
.flatMap(subset -> {
|
||||
Map<String, String> metadata = new HashMap<>(svcMetadata);
|
||||
List<V1EndpointPort> endpointPorts = subset.getPorts();
|
||||
if (this.properties.getMetadata() != null && this.properties.getMetadata().addPorts()) {
|
||||
endpointPorts.forEach(
|
||||
p -> metadata.put(StringUtils.hasText(p.getName()) ? p.getName() : UNSET_PORT_NAME,
|
||||
Integer.toString(p.getPort())));
|
||||
}
|
||||
List<V1EndpointAddress> addresses = subset.getAddresses();
|
||||
if (addresses == null) {
|
||||
addresses = new ArrayList<>();
|
||||
}
|
||||
if (this.properties.isIncludeNotReadyAddresses()
|
||||
&& !CollectionUtils.isEmpty(subset.getNotReadyAddresses())) {
|
||||
addresses.addAll(subset.getNotReadyAddresses());
|
||||
}
|
||||
|
||||
final int port = findEndpointPort(endpointPorts, primaryPortName, serviceId);
|
||||
return addresses.stream()
|
||||
.map(addr -> new DefaultKubernetesServiceInstance(
|
||||
addr.getTargetRef() != null ? addr.getTargetRef().getUid() : "", serviceId,
|
||||
addr.getIp(), port, metadata, false, service.getMetadata().getNamespace(),
|
||||
service.getMetadata().getClusterName()));
|
||||
}).collect(Collectors.toList());
|
||||
final int port = findEndpointPort(endpointPorts, primaryPortName, serviceId);
|
||||
return addresses.stream()
|
||||
.map(addr -> new DefaultKubernetesServiceInstance(
|
||||
addr.getTargetRef() != null ? addr.getTargetRef().getUid() : "", serviceId,
|
||||
addr.getIp(), port, metadata, false, service.getMetadata().getNamespace(),
|
||||
service.getMetadata().getClusterName()));
|
||||
}).collect(Collectors.toList());
|
||||
}
|
||||
|
||||
private int findEndpointPort(List<V1EndpointPort> endpointPorts, String primaryPortName, String serviceId) {
|
||||
@@ -175,7 +175,7 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
|
||||
}
|
||||
else {
|
||||
Map<String, Integer> ports = endpointPorts.stream().filter(p -> StringUtils.hasText(p.getName()))
|
||||
.collect(Collectors.toMap(V1EndpointPort::getName, V1EndpointPort::getPort));
|
||||
.collect(Collectors.toMap(V1EndpointPort::getName, V1EndpointPort::getPort));
|
||||
// This oneliner is looking for a port with a name equal to the primary port
|
||||
// name specified in the service label
|
||||
// or in spring.cloud.kubernetes.discovery.primary-port-name, equal to https,
|
||||
@@ -183,18 +183,18 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
|
||||
// In case no port has been found return -1 to log a warning and fall back to
|
||||
// the first port in the list.
|
||||
int discoveredPort = ports.getOrDefault(primaryPortName,
|
||||
ports.getOrDefault(HTTPS, ports.getOrDefault(HTTP, -1)));
|
||||
ports.getOrDefault(HTTPS, ports.getOrDefault(HTTP, -1)));
|
||||
|
||||
if (discoveredPort == -1) {
|
||||
if (StringUtils.hasText(primaryPortName)) {
|
||||
log.warn("Could not find a port named '" + primaryPortName + "', 'https', or 'http' for service '"
|
||||
+ serviceId + "'.");
|
||||
+ serviceId + "'.");
|
||||
}
|
||||
else {
|
||||
log.warn("Could not find a port named 'https' or 'http' for service '" + serviceId + "'.");
|
||||
}
|
||||
log.warn(
|
||||
"Make sure that either the primary-port-name label has been added to the service, or that spring.cloud.kubernetes.discovery.primary-port-name has been configured.");
|
||||
"Make sure that either the primary-port-name label has been added to the service, or that spring.cloud.kubernetes.discovery.primary-port-name has been configured.");
|
||||
log.warn("Alternatively name the primary port 'https' or 'http'");
|
||||
log.warn("An incorrect configuration may result in non-deterministic behaviour.");
|
||||
discoveredPort = endpointPorts.get(0).getPort();
|
||||
@@ -206,30 +206,30 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
|
||||
@Override
|
||||
public List<String> getServices() {
|
||||
List<V1Service> services = this.properties.isAllNamespaces() ? this.serviceLister.list()
|
||||
: this.serviceLister.namespace(this.namespace).list();
|
||||
: this.serviceLister.namespace(this.namespace).list();
|
||||
return services.stream().filter(this::matchServiceLabels).map(s -> s.getMetadata().getName())
|
||||
.collect(Collectors.toList());
|
||||
.collect(Collectors.toList());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterPropertiesSet() throws Exception {
|
||||
this.sharedInformerFactory.startAllRegisteredInformers();
|
||||
if (!Wait.poll(Duration.ofSeconds(1), Duration.ofSeconds(this.properties.getCacheLoadingTimeoutSeconds()),
|
||||
() -> {
|
||||
log.info("Waiting for the cache of informers to be fully loaded..");
|
||||
return this.informersReadyFunc.get();
|
||||
})) {
|
||||
() -> {
|
||||
log.info("Waiting for the cache of informers to be fully loaded..");
|
||||
return this.informersReadyFunc.get();
|
||||
})) {
|
||||
if (this.properties.isWaitCacheReady()) {
|
||||
throw new IllegalStateException(
|
||||
"Timeout waiting for informers cache to be ready, is the kubernetes service up?");
|
||||
"Timeout waiting for informers cache to be ready, is the kubernetes service up?");
|
||||
}
|
||||
else {
|
||||
log.warn(
|
||||
"Timeout waiting for informers cache to be ready, ignoring the failure because waitForInformerCacheReady property is false");
|
||||
"Timeout waiting for informers cache to be ready, ignoring the failure because waitForInformerCacheReady property is false");
|
||||
}
|
||||
}
|
||||
log.info("Cache fully loaded (total " + serviceLister.list().size()
|
||||
+ " services) , discovery client is now available");
|
||||
+ " services) , discovery client is now available");
|
||||
}
|
||||
|
||||
private boolean matchServiceLabels(V1Service service) {
|
||||
@@ -251,9 +251,9 @@ public class KubernetesInformerDiscoveryClient implements DiscoveryClient, Initi
|
||||
return true;
|
||||
}
|
||||
return properties.getServiceLabels().keySet().stream()
|
||||
.allMatch(k -> service.getMetadata().getLabels() != null
|
||||
&& service.getMetadata().getLabels().containsKey(k)
|
||||
&& service.getMetadata().getLabels().get(k).equals(properties.getServiceLabels().get(k)));
|
||||
.allMatch(k -> service.getMetadata().getLabels() != null
|
||||
&& service.getMetadata().getLabels().containsKey(k)
|
||||
&& service.getMetadata().getLabels().get(k).equals(properties.getServiceLabels().get(k)));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -80,9 +80,9 @@ public class KubernetesInformerDiscoveryClientTests {
|
||||
.addSubsetsItem(new V1EndpointSubset().addAddressesItem(new V1EndpointAddress().ip("1.1.1.1")));
|
||||
|
||||
private static final V1Endpoints testEndpointWithUnsetPortName = new V1Endpoints()
|
||||
.metadata(new V1ObjectMeta().name("test-svc-1").namespace("namespace1"))
|
||||
.addSubsetsItem(new V1EndpointSubset().addPortsItem(new V1EndpointPort().port(80))
|
||||
.addAddressesItem(new V1EndpointAddress().ip("1.1.1.1")));
|
||||
.metadata(new V1ObjectMeta().name("test-svc-1").namespace("namespace1"))
|
||||
.addSubsetsItem(new V1EndpointSubset().addPortsItem(new V1EndpointPort().port(80))
|
||||
.addAddressesItem(new V1EndpointAddress().ip("1.1.1.1")));
|
||||
|
||||
private static final V1Endpoints testEndpointWithMultiplePorts = new V1Endpoints()
|
||||
.metadata(new V1ObjectMeta().name("test-svc-1").namespace("namespace1"))
|
||||
@@ -116,12 +116,13 @@ public class KubernetesInformerDiscoveryClientTests {
|
||||
when(kubernetesDiscoveryProperties.getMetadata()).thenReturn(new KubernetesDiscoveryProperties.Metadata());
|
||||
|
||||
KubernetesInformerDiscoveryClient discoveryClient = new KubernetesInformerDiscoveryClient("",
|
||||
sharedInformerFactory, serviceLister, endpointsLister, null, null, kubernetesDiscoveryProperties);
|
||||
sharedInformerFactory, serviceLister, endpointsLister, null, null, kubernetesDiscoveryProperties);
|
||||
|
||||
Map<String, String> ports = new HashMap<>();
|
||||
ports.put("<unset>", "80");
|
||||
assertThat(discoveryClient.getInstances("test-svc-1").toArray()).containsOnly(new KubernetesServiceInstance("",
|
||||
"test-svc-1", "1.1.1.1", 80, ports, false, "namespace1", null));
|
||||
assertThat(discoveryClient.getInstances("test-svc-1").toArray())
|
||||
.containsOnly(new DefaultKubernetesServiceInstance("", "test-svc-1", "1.1.1.1", 80, ports, false,
|
||||
"namespace1", null));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -304,7 +305,7 @@ public class KubernetesInformerDiscoveryClientTests {
|
||||
verify(kubernetesDiscoveryProperties, times(1)).isAllNamespaces();
|
||||
verify(kubernetesDiscoveryProperties, times(1)).getPrimaryPortName();
|
||||
verify(kubernetesDiscoveryProperties, times(1)).isIncludeNotReadyAddresses();
|
||||
//Reset metadata
|
||||
// Reset metadata
|
||||
testService1.metadata(oldMetadata);
|
||||
}
|
||||
|
||||
@@ -325,7 +326,7 @@ public class KubernetesInformerDiscoveryClientTests {
|
||||
"test-svc-1", "1.1.1.1", 80, new HashMap<>(), false, "namespace1", null));
|
||||
verify(kubernetesDiscoveryProperties, times(1)).isAllNamespaces();
|
||||
verify(kubernetesDiscoveryProperties, times(1)).getPrimaryPortName();
|
||||
//Reset testService1 metadata
|
||||
// Reset testService1 metadata
|
||||
testService1.metadata(oldMetadata);
|
||||
}
|
||||
|
||||
|
||||
@@ -47,4 +47,9 @@ public final class KubernetesDiscoveryConstants {
|
||||
*/
|
||||
public static final String NAMESPACE_METADATA_KEY = "k8s_namespace";
|
||||
|
||||
/**
|
||||
* Port name to use when there isn't one set.
|
||||
*/
|
||||
public static final String UNSET_PORT_NAME = "<unset>";
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user