This is an automated email from the ASF dual-hosted git repository.
Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git
The following commit(s) were added to refs/heads/master by this push:
new 5d3555b896 fix: keep the service port selected by an ingress on
endpoint updates (#6493) (#7287)
5d3555b896 is described below
commit 5d3555b89698994f784f9bd3b15e9ab06c8d351a
Author: wy471x <[email protected]>
AuthorDate: Thu Oct 1 05:54:15 2026 +0800
fix: keep the service port selected by an ingress on endpoint updates
(#6493) (#7287)
---
.../shenyu/k8s/cache/ServiceIngressCache.java | 46 ++---
.../shenyu/k8s/common/IngressBackendPort.java | 140 +++++++++++++++
.../shenyu/k8s/common/ServiceIngressRelation.java | 108 ++++++++++++
.../shenyu/k8s/reconciler/EndpointsReconciler.java | 110 ++++++------
.../shenyu/k8s/reconciler/IngressReconciler.java | 79 +++++----
.../apache/shenyu/k8s/EndpointsReconcilerTest.java | 192 +++++++++++++++++----
6 files changed, 536 insertions(+), 139 deletions(-)
diff --git
a/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/cache/ServiceIngressCache.java
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/cache/ServiceIngressCache.java
index 5328247811..5722ae0935 100644
---
a/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/cache/ServiceIngressCache.java
+++
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/cache/ServiceIngressCache.java
@@ -18,21 +18,22 @@
package org.apache.shenyu.k8s.cache;
import com.google.common.collect.Maps;
-import org.apache.commons.lang3.tuple.Pair;
+import org.apache.shenyu.k8s.common.ServiceIngressRelation;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Objects;
/**
- * The cache for mapping service name to ingress name.
+ * The cache for mapping service name to the ingress relations that reference
the service.
*/
public final class ServiceIngressCache {
private static final ServiceIngressCache INSTANCE = new
ServiceIngressCache();
- private static final Map<String, List<Pair<String, String>>> INGRESS_MAP =
Maps.newConcurrentMap();
+ private static final Map<String, List<ServiceIngressRelation>> INGRESS_MAP
= Maps.newConcurrentMap();
private ServiceIngressCache() {
}
@@ -47,42 +48,47 @@ public final class ServiceIngressCache {
}
/**
- * Get ingress namespace and name by service namespace and namespace.
+ * Get the ingress relations of the service, each relation keeps the
service port selected by the ingress.
*
- * @param namespace namespace
+ * @param namespace service namespace
* @param serviceName service name
- * @return ingress namespace and name
+ * @return the ingress relations of the service, empty if the service is
not referenced
*/
- public List<Pair<String, String>> getIngressName(final String namespace,
final String serviceName) {
- return INGRESS_MAP.get(getKey(namespace, serviceName));
+ public List<ServiceIngressRelation> getIngressName(final String namespace,
final String serviceName) {
+ List<ServiceIngressRelation> res = INGRESS_MAP.get(getKey(namespace,
serviceName));
+ return Objects.isNull(res) ? Collections.emptyList() : res;
}
/**
- * Put ingress by service namespace and name.
+ * Put the ingress that references the service, the previous relation of
the same ingress is
+ * replaced so that a changed backend service port does not leave a stale
relation behind.
*
* @param namespace service namespace
* @param serviceName service name
- * @param ingressNamespace ingress namespace
- * @param ingressName ingress name
+ * @param relation ingress relation of the service
*/
- public void putIngressName(final String namespace, final String
serviceName, final String ingressNamespace, final String ingressName) {
- List<Pair<String, String>> list =
INGRESS_MAP.computeIfAbsent(getKey(namespace, serviceName), k -> new
ArrayList<>());
- list.add(Pair.of(ingressNamespace, ingressName));
+ public void putIngressName(final String namespace, final String
serviceName, final ServiceIngressRelation relation) {
+ INGRESS_MAP.compute(getKey(namespace, serviceName), (key, relations)
-> {
+ List<ServiceIngressRelation> res = Objects.isNull(relations) ? new
ArrayList<>() : relations;
+ res.removeIf(item ->
item.isSameIngress(relation.getIngressNamespace(), relation.getIngressName()));
+ res.add(relation);
+ return res;
+ });
}
/**
- * Remove all ingress by service namespace and name.
+ * Remove all ingress relations by service namespace and name.
*
* @param namespace service namespace
* @param serviceName service name
- * @return the ingress list removed
+ * @return the ingress relation list removed
*/
- public List<Pair<String, String>> removeAllIngressName(final String
namespace, final String serviceName) {
+ public List<ServiceIngressRelation> removeAllIngressName(final String
namespace, final String serviceName) {
return INGRESS_MAP.remove(getKey(namespace, serviceName));
}
/**
- * Remove specified ingress by service and ingress.
+ * Remove specified ingress relation by service and ingress.
*
* @param namespace service namespace
* @param serviceName service name
@@ -90,9 +96,9 @@ public final class ServiceIngressCache {
* @param ingressName ingress name
*/
public void removeSpecifiedIngressName(final String namespace, final
String serviceName, final String ingressNamespace, final String ingressName) {
- List<Pair<String, String>> list = INGRESS_MAP.get(getKey(namespace,
serviceName));
+ List<ServiceIngressRelation> list = INGRESS_MAP.get(getKey(namespace,
serviceName));
if (Objects.nonNull(list)) {
- list.removeIf(item -> item.getLeft().equals(ingressNamespace) &&
item.getRight().equals(ingressName));
+ list.removeIf(item -> item.isSameIngress(ingressNamespace,
ingressName));
}
}
diff --git
a/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/common/IngressBackendPort.java
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/common/IngressBackendPort.java
new file mode 100644
index 0000000000..b54778ea0c
--- /dev/null
+++
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/common/IngressBackendPort.java
@@ -0,0 +1,140 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.k8s.common;
+
+import io.kubernetes.client.openapi.models.CoreV1EndpointPort;
+import io.kubernetes.client.openapi.models.V1ServiceBackendPort;
+import org.apache.commons.lang3.StringUtils;
+
+import java.util.List;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+/**
+ * The service port selected by an ingress backend, either a port name or a
port number.
+ */
+public final class IngressBackendPort {
+
+ private static final String TCP_PROTOCOL = "TCP";
+
+ private final String name;
+
+ private final Integer number;
+
+ private IngressBackendPort(final String name, final Integer number) {
+ this.name = name;
+ this.number = number;
+ }
+
+ /**
+ * Build the port from the service port of an ingress backend.
+ *
+ * @param servicePort service port of the ingress backend, may be null
+ * @return the selected service port, or null if the backend does not
select one
+ */
+ public static IngressBackendPort from(final V1ServiceBackendPort
servicePort) {
+ if (Objects.isNull(servicePort)) {
+ return null;
+ }
+ if (Objects.nonNull(servicePort.getNumber()) &&
servicePort.getNumber() > 0) {
+ return new IngressBackendPort(null, servicePort.getNumber());
+ }
+ if (StringUtils.isNotBlank(servicePort.getName())) {
+ return new IngressBackendPort(servicePort.getName().trim(), null);
+ }
+ return null;
+ }
+
+ /**
+ * Select the endpoint port that serves the service port, the first TCP
port is used when no
+ * endpoint port matches because a service may map the selected port to a
different target port.
+ *
+ * @param ports endpoint ports of an endpoint subset
+ * @param backendPort service port selected by the ingress backend, may be
null
+ * @return the endpoint port to route to, or null if the subset does not
expose a TCP port
+ */
+ public static CoreV1EndpointPort selectEndpointPort(final
List<CoreV1EndpointPort> ports, final IngressBackendPort backendPort) {
+ List<CoreV1EndpointPort> tcpPorts = ports.stream()
+ .filter(port -> TCP_PROTOCOL.equals(port.getProtocol()))
+ .collect(Collectors.toList());
+ if (tcpPorts.isEmpty()) {
+ return null;
+ }
+ if (Objects.nonNull(backendPort)) {
+ for (CoreV1EndpointPort tcpPort : tcpPorts) {
+ if (backendPort.matches(tcpPort)) {
+ return tcpPort;
+ }
+ }
+ }
+ return tcpPorts.get(0);
+ }
+
+ /**
+ * Whether the endpoint port serves this service port.
+ *
+ * @param endpointPort endpoint port
+ * @return true if the endpoint port matches this service port
+ */
+ public boolean matches(final CoreV1EndpointPort endpointPort) {
+ if (Objects.nonNull(number) && number.equals(endpointPort.getPort())) {
+ return true;
+ }
+ return Objects.nonNull(name) && name.equals(endpointPort.getName());
+ }
+
+ /**
+ * Get the port name.
+ *
+ * @return port name, null if the port is selected by number
+ */
+ public String getName() {
+ return name;
+ }
+
+ /**
+ * Get the port number.
+ *
+ * @return port number, null if the port is selected by name
+ */
+ public Integer getNumber() {
+ return number;
+ }
+
+ @Override
+ public boolean equals(final Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (Objects.isNull(o) || getClass() != o.getClass()) {
+ return false;
+ }
+ IngressBackendPort that = (IngressBackendPort) o;
+ return Objects.equals(name, that.name) && Objects.equals(number,
that.number);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(name, number);
+ }
+
+ @Override
+ public String toString() {
+ return "IngressBackendPort{name='" + name + "', number=" + number +
'}';
+ }
+}
diff --git
a/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/common/ServiceIngressRelation.java
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/common/ServiceIngressRelation.java
new file mode 100644
index 0000000000..c68d7edaa7
--- /dev/null
+++
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/common/ServiceIngressRelation.java
@@ -0,0 +1,108 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.k8s.common;
+
+import java.util.Objects;
+
+/**
+ * The relation between a service and the ingress that references it,
including the service port
+ * selected by the ingress backend.
+ */
+public final class ServiceIngressRelation {
+
+ private final String ingressNamespace;
+
+ private final String ingressName;
+
+ private final IngressBackendPort port;
+
+ /**
+ * Constructor of ServiceIngressRelation.
+ *
+ * @param ingressNamespace ingress namespace
+ * @param ingressName ingress name
+ * @param port service port selected by the ingress backend, may be null
+ */
+ public ServiceIngressRelation(final String ingressNamespace, final String
ingressName, final IngressBackendPort port) {
+ this.ingressNamespace = ingressNamespace;
+ this.ingressName = ingressName;
+ this.port = port;
+ }
+
+ /**
+ * Get the ingress namespace.
+ *
+ * @return ingress namespace
+ */
+ public String getIngressNamespace() {
+ return ingressNamespace;
+ }
+
+ /**
+ * Get the ingress name.
+ *
+ * @return ingress name
+ */
+ public String getIngressName() {
+ return ingressName;
+ }
+
+ /**
+ * Get the service port selected by the ingress backend.
+ *
+ * @return selected service port, null if the backend does not select one
+ */
+ public IngressBackendPort getPort() {
+ return port;
+ }
+
+ /**
+ * Whether the relation belongs to the given ingress.
+ *
+ * @param namespace ingress namespace
+ * @param name ingress name
+ * @return true if the relation belongs to the ingress
+ */
+ public boolean isSameIngress(final String namespace, final String name) {
+ return Objects.equals(ingressNamespace, namespace) &&
Objects.equals(ingressName, name);
+ }
+
+ @Override
+ public boolean equals(final Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (Objects.isNull(o) || getClass() != o.getClass()) {
+ return false;
+ }
+ ServiceIngressRelation that = (ServiceIngressRelation) o;
+ return Objects.equals(ingressNamespace, that.ingressNamespace)
+ && Objects.equals(ingressName, that.ingressName)
+ && Objects.equals(port, that.port);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(ingressNamespace, ingressName, port);
+ }
+
+ @Override
+ public String toString() {
+ return "ServiceIngressRelation{ingressNamespace='" + ingressNamespace
+ "', ingressName='" + ingressName + "', port=" + port + '}';
+ }
+}
diff --git
a/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/EndpointsReconciler.java
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/EndpointsReconciler.java
index b02ccac54d..6effdc83f6 100644
---
a/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/EndpointsReconciler.java
+++
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/EndpointsReconciler.java
@@ -34,10 +34,11 @@ import org.apache.shenyu.common.dto.SelectorData;
import org.apache.shenyu.common.dto.convert.selector.DivideUpstream;
import org.apache.shenyu.common.dto.convert.selector.WebSocketUpstream;
import org.apache.shenyu.common.enums.PluginEnum;
-import org.apache.shenyu.common.exception.ShenyuException;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.k8s.cache.IngressSelectorCache;
import org.apache.shenyu.k8s.cache.ServiceIngressCache;
+import org.apache.shenyu.k8s.common.IngressBackendPort;
+import org.apache.shenyu.k8s.common.ServiceIngressRelation;
import org.apache.shenyu.k8s.repository.ShenyuCacheRepository;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -92,7 +93,7 @@ public class EndpointsReconciler implements Reconciler {
*/
@Override
public Result reconcile(final Request request) {
- List<Pair<String, String>> ingressList =
ServiceIngressCache.getInstance().getIngressName(request.getNamespace(),
request.getName());
+ List<ServiceIngressRelation> ingressList =
ServiceIngressCache.getInstance().getIngressName(request.getNamespace(),
request.getName());
if (CollectionUtils.isEmpty(ingressList)) {
return new Result(false);
}
@@ -105,14 +106,14 @@ public class EndpointsReconciler implements Reconciler {
return new Result(false);
}
- updateSelectors(ingressList, PluginEnum.DIVIDE.getName(),
getDivideUpstreamFromEndpoints(v1Endpoints));
- updateSelectors(ingressList, PluginEnum.WEB_SOCKET.getName(),
getWebSocketUpstreamFromEndpoints(v1Endpoints));
+ updateSelectors(ingressList, v1Endpoints, PluginEnum.DIVIDE.getName());
+ updateSelectors(ingressList, v1Endpoints,
PluginEnum.WEB_SOCKET.getName());
LOG.info("Update selector for endpoint {}", request);
return new Result(false);
}
- private void updateSelectors(final List<Pair<String, String>> ingressList,
final String pluginName, final String handle) {
+ private void updateSelectors(final List<ServiceIngressRelation>
ingressList, final V1Endpoints v1Endpoints, final String pluginName) {
if (!ENDPOINT_UPSTREAM_PLUGINS.contains(pluginName)) {
return;
}
@@ -120,39 +121,59 @@ public class EndpointsReconciler implements Reconciler {
if (CollectionUtils.isEmpty(totalSelectors)) {
return;
}
- Set<String> needUpdateSelectorId = new HashSet<>();
- ingressList.forEach(item -> {
- List<String> selectorIdList =
IngressSelectorCache.getInstance().get(item.getLeft(), item.getRight(),
pluginName);
- if (CollectionUtils.isNotEmpty(selectorIdList)) {
- needUpdateSelectorId.addAll(selectorIdList);
+ for (ServiceIngressRelation relation : ingressList) {
+ List<String> selectorIdList = IngressSelectorCache.getInstance()
+ .get(relation.getIngressNamespace(),
relation.getIngressName(), pluginName);
+ if (CollectionUtils.isEmpty(selectorIdList)) {
+ continue;
}
- });
- if (needUpdateSelectorId.isEmpty()) {
- return;
- }
- totalSelectors.forEach(selectorData -> {
- if (needUpdateSelectorId.contains(selectorData.getId())) {
- SelectorData newSelectorData =
SelectorData.builder().id(selectorData.getId())
- .pluginId(selectorData.getPluginId())
- .pluginName(selectorData.getPluginName())
- .name(selectorData.getName())
- .matchMode(selectorData.getMatchMode())
- .type(selectorData.getType())
- .sort(selectorData.getSort())
- .enabled(selectorData.getEnabled())
- .logged(selectorData.getLogged())
- .continued(selectorData.getContinued())
- .handle(handle)
- .conditionList(selectorData.getConditionList())
- .matchRestful(selectorData.getMatchRestful()).build();
-
shenyuCacheRepository.saveOrUpdateSelectorData(newSelectorData);
+ // each ingress selects its own service port, so the upstream
handle of an ingress must
+ // be rebuilt with the endpoints of that port
+ String handle = getUpstreamHandle(endpointAddresses(v1Endpoints,
relation.getPort()), pluginName);
+ if (Objects.isNull(handle)) {
+ LOG.info("Cannot find endpoint addresses of the backend port
{} for ingress {}/{}",
+ relation.getPort(), relation.getIngressNamespace(),
relation.getIngressName());
+ continue;
}
- });
+ totalSelectors.forEach(selectorData -> {
+ if (selectorIdList.contains(selectorData.getId())) {
+ SelectorData newSelectorData =
SelectorData.builder().id(selectorData.getId())
+ .pluginId(selectorData.getPluginId())
+ .pluginName(selectorData.getPluginName())
+ .name(selectorData.getName())
+ .matchMode(selectorData.getMatchMode())
+ .type(selectorData.getType())
+ .sort(selectorData.getSort())
+ .enabled(selectorData.getEnabled())
+ .logged(selectorData.getLogged())
+ .continued(selectorData.getContinued())
+ .handle(handle)
+ .conditionList(selectorData.getConditionList())
+
.matchRestful(selectorData.getMatchRestful()).build();
+
shenyuCacheRepository.saveOrUpdateSelectorData(newSelectorData);
+ }
+ });
+ }
}
- private String getDivideUpstreamFromEndpoints(final V1Endpoints
v1Endpoints) {
+ private String getUpstreamHandle(final List<Pair<V1EndpointAddress,
String>> addresses, final String pluginName) {
+ if (CollectionUtils.isEmpty(addresses)) {
+ return null;
+ }
+ if (PluginEnum.WEB_SOCKET.getName().equals(pluginName)) {
+ List<WebSocketUpstream> res = new ArrayList<>();
+ addresses.forEach(pair -> res.add(WebSocketUpstream.builder()
+ .upstreamUrl(pair.getLeft().getIp() + ":" +
pair.getRight())
+ .weight(100)
+ .protocol("ws://")
+ .warmup(0)
+ .status(true)
+ .host("")
+ .build()));
+ return GsonUtils.getInstance().toJson(res);
+ }
List<DivideUpstream> res = new ArrayList<>();
- endpointAddresses(v1Endpoints).forEach(pair -> {
+ addresses.forEach(pair -> {
DivideUpstream upstream = new DivideUpstream();
upstream.setUpstreamUrl(pair.getLeft().getIp() + ":" +
pair.getRight());
upstream.setWeight(100);
@@ -166,20 +187,7 @@ public class EndpointsReconciler implements Reconciler {
return GsonUtils.getInstance().toJson(res);
}
- private String getWebSocketUpstreamFromEndpoints(final V1Endpoints
v1Endpoints) {
- List<WebSocketUpstream> res = new ArrayList<>();
- endpointAddresses(v1Endpoints).forEach(pair ->
res.add(WebSocketUpstream.builder()
- .upstreamUrl(pair.getLeft().getIp() + ":" + pair.getRight())
- .weight(100)
- .protocol("ws://")
- .warmup(0)
- .status(true)
- .host("")
- .build()));
- return GsonUtils.getInstance().toJson(res);
- }
-
- private List<Pair<V1EndpointAddress, String>> endpointAddresses(final
V1Endpoints v1Endpoints) {
+ private List<Pair<V1EndpointAddress, String>> endpointAddresses(final
V1Endpoints v1Endpoints, final IngressBackendPort backendPort) {
List<Pair<V1EndpointAddress, String>> res = new ArrayList<>();
List<V1EndpointSubset> subsets = v1Endpoints.getSubsets();
if (CollectionUtils.isNotEmpty(subsets)) {
@@ -189,10 +197,10 @@ public class EndpointsReconciler implements Reconciler {
if (CollectionUtils.isEmpty(ports) ||
CollectionUtils.isEmpty(addresses)) {
continue;
}
- CoreV1EndpointPort endpointPort = ports.stream()
- .filter(coreV1EndpointPort ->
"TCP".equals(coreV1EndpointPort.getProtocol()))
- .findFirst()
- .orElseThrow(() -> new ShenyuException("Can't find
port from endpoints"));
+ CoreV1EndpointPort endpointPort =
IngressBackendPort.selectEndpointPort(ports, backendPort);
+ if (Objects.isNull(endpointPort)) {
+ continue;
+ }
String port = null;
if (endpointPort.getPort() > 0) {
port = String.valueOf(endpointPort.getPort());
diff --git
a/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/IngressReconciler.java
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/IngressReconciler.java
index 9d44018074..bebba6fc94 100644
---
a/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/IngressReconciler.java
+++
b/shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/IngressReconciler.java
@@ -32,6 +32,7 @@ import io.kubernetes.client.openapi.models.V1HTTPIngressPath;
import io.kubernetes.client.openapi.models.V1Ingress;
import io.kubernetes.client.openapi.models.V1IngressBuilder;
import io.kubernetes.client.openapi.models.V1IngressRule;
+import io.kubernetes.client.openapi.models.V1IngressServiceBackend;
import io.kubernetes.client.openapi.models.V1Secret;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.commons.collections4.MapUtils;
@@ -53,8 +54,10 @@ import org.apache.shenyu.k8s.cache.IngressCache;
import org.apache.shenyu.k8s.cache.IngressSecretCache;
import org.apache.shenyu.k8s.cache.IngressSelectorCache;
import org.apache.shenyu.k8s.cache.ServiceIngressCache;
+import org.apache.shenyu.k8s.common.IngressBackendPort;
import org.apache.shenyu.k8s.common.IngressConfiguration;
import org.apache.shenyu.k8s.common.IngressConstants;
+import org.apache.shenyu.k8s.common.ServiceIngressRelation;
import org.apache.shenyu.k8s.common.ShenyuMemoryConfig;
import org.apache.shenyu.k8s.parser.IngressParser;
import org.apache.shenyu.k8s.repository.ShenyuCacheRepository;
@@ -64,6 +67,7 @@ import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
@@ -187,10 +191,12 @@ public class IngressReconciler implements Reconciler {
}
}
IngressCache.getInstance().put(request.getNamespace(),
request.getName(), v1Ingress);
- List<Pair<String, String>> serviceList =
parseServiceFromIngress(v1Ingress);
- Objects.requireNonNull(serviceList).forEach(pair -> {
- ServiceIngressCache.getInstance().putIngressName(pair.getLeft(),
pair.getRight(), request.getNamespace(), request.getName());
- LOG.info("Add service cache {} for ingress {}", pair.getLeft() +
"/" + pair.getRight(), request.getNamespace() + "/" + request.getName());
+ final String serviceNamespace =
Objects.requireNonNull(v1Ingress.getMetadata()).getNamespace();
+ parseServiceFromIngress(v1Ingress).forEach((serviceName, port) -> {
+ ServiceIngressCache.getInstance().putIngressName(serviceNamespace,
serviceName,
+ new ServiceIngressRelation(request.getNamespace(),
request.getName(), port));
+ LOG.info("Add service cache {} for ingress {}, the backend port is
{}", serviceNamespace + "/" + serviceName,
+ request.getNamespace() + "/" + request.getName(), port);
});
// Ensure upstream handles are populated from endpoints
@@ -244,10 +250,10 @@ public class IngressReconciler implements Reconciler {
IngressSelectorCache.getInstance().remove(request.getNamespace(),
request.getName(), PluginEnum.DIVIDE.getName());
}
}
- List<Pair<String, String>> serviceList =
parseServiceFromIngress(oldIngress);
- Objects.requireNonNull(serviceList).forEach(pair -> {
-
ServiceIngressCache.getInstance().removeSpecifiedIngressName(pair.getLeft(),
pair.getRight(), request.getNamespace(), request.getName());
- LOG.info("Delete service cache {} for ingress {}", pair.getLeft()
+ "/" + pair.getRight(), request.getNamespace() + "/" + request.getName());
+ final String serviceNamespace =
Objects.requireNonNull(oldIngress.getMetadata()).getNamespace();
+ parseServiceFromIngress(oldIngress).forEach((serviceName, port) -> {
+
ServiceIngressCache.getInstance().removeSpecifiedIngressName(serviceNamespace,
serviceName, request.getNamespace(), request.getName());
+ LOG.info("Delete service cache {} for ingress {}",
serviceNamespace + "/" + serviceName, request.getNamespace() + "/" +
request.getName());
});
deleteGlobalDefaultBackend(request.getNamespace(), request.getName());
}
@@ -382,29 +388,35 @@ public class IngressReconciler implements Reconciler {
return selectorList;
}
- private List<Pair<String, String>> parseServiceFromIngress(final V1Ingress
ingress) {
- List<Pair<String, String>> res = new ArrayList<>();
+ /**
+ * Parse the backend services referenced by the ingress, mapped to the
service port selected by the ingress.
+ *
+ * @param ingress ingress resource
+ * @return the backend service names mapped to the service port selected
by the ingress
+ */
+ private Map<String, IngressBackendPort> parseServiceFromIngress(final
V1Ingress ingress) {
+ Map<String, IngressBackendPort> res = new HashMap<>(4);
if (Objects.isNull(ingress) || Objects.isNull(ingress.getSpec())) {
return res;
}
String namespace =
Objects.requireNonNull(ingress.getMetadata()).getNamespace();
String name = ingress.getMetadata().getName();
String namespacedName = namespace + "/" + name;
- String defaultService = null;
+ V1IngressServiceBackend defaultBackendService = null;
if (Objects.nonNull(ingress.getSpec().getDefaultBackend()) &&
Objects.nonNull(ingress.getSpec().getDefaultBackend().getService())) {
- defaultService =
ingress.getSpec().getDefaultBackend().getService().getName();
+ defaultBackendService =
ingress.getSpec().getDefaultBackend().getService();
+ String defaultService = defaultBackendService.getName();
if (Objects.isNull(ingress.getSpec().getRules())) {
if (Objects.nonNull(globalDefaultBackend)) {
if
(globalDefaultBackend.getLeft().getLeft().equals(namespacedName)) {
- res.add(Pair.of(namespace, defaultService));
+ res.put(defaultService,
IngressBackendPort.from(defaultBackendService.getPort()));
}
} else {
- res.add(Pair.of(namespace, defaultService));
+ res.put(defaultService,
IngressBackendPort.from(defaultBackendService.getPort()));
}
return res;
}
}
- Set<String> deduplicateSet = new HashSet<>();
if (Objects.isNull(ingress.getSpec().getRules())) {
return res;
}
@@ -412,15 +424,10 @@ public class IngressReconciler implements Reconciler {
if (Objects.nonNull(rule.getHttp()) &&
Objects.nonNull(rule.getHttp().getPaths())) {
for (V1HTTPIngressPath path : rule.getHttp().getPaths()) {
if (Objects.nonNull(path.getBackend()) &&
Objects.nonNull(path.getBackend().getService())) {
- if
(!deduplicateSet.contains(path.getBackend().getService().getName())) {
- res.add(Pair.of(namespace,
path.getBackend().getService().getName()));
-
deduplicateSet.add(path.getBackend().getService().getName());
- }
- } else {
- if (Objects.nonNull(defaultService) &&
!deduplicateSet.contains(defaultService)) {
- res.add(Pair.of(namespace, defaultService));
- deduplicateSet.add(defaultService);
- }
+ V1IngressServiceBackend backendService =
path.getBackend().getService();
+ res.putIfAbsent(backendService.getName(),
IngressBackendPort.from(backendService.getPort()));
+ } else if (Objects.nonNull(defaultBackendService)) {
+ res.putIfAbsent(defaultBackendService.getName(),
IngressBackendPort.from(defaultBackendService.getPort()));
}
}
}
@@ -572,22 +579,21 @@ public class IngressReconciler implements Reconciler {
if (!PluginEnum.DIVIDE.getName().equals(pluginName) &&
!PluginEnum.WEB_SOCKET.getName().equals(pluginName)) {
return;
}
- List<Pair<String, String>> serviceList =
parseServiceFromIngress(v1Ingress);
- if (CollectionUtils.isEmpty(serviceList)) {
+ Map<String, IngressBackendPort> serviceList =
parseServiceFromIngress(v1Ingress);
+ if (serviceList.isEmpty()) {
return;
}
String namespace =
Objects.requireNonNull(v1Ingress.getMetadata()).getNamespace();
String ingressName = v1Ingress.getMetadata().getName();
Lister<V1Endpoints> endpointsLister =
ingressParser.getEndpointsLister();
- for (Pair<String, String> service : serviceList) {
- String serviceNamespace = service.getLeft();
- String serviceName = service.getRight();
- V1Endpoints v1Endpoints =
endpointsLister.namespace(serviceNamespace).get(serviceName);
+ for (Map.Entry<String, IngressBackendPort> service :
serviceList.entrySet()) {
+ String serviceName = service.getKey();
+ V1Endpoints v1Endpoints =
endpointsLister.namespace(namespace).get(serviceName);
if (Objects.isNull(v1Endpoints)) {
- LOG.info("Cannot find endpoints for service {}/{} when
updating upstream", serviceNamespace, serviceName);
+ LOG.info("Cannot find endpoints for service {}/{} when
updating upstream", namespace, serviceName);
continue;
}
- List<Pair<V1EndpointAddress, String>> addresses =
endpointAddresses(v1Endpoints);
+ List<Pair<V1EndpointAddress, String>> addresses =
endpointAddresses(v1Endpoints, service.getValue());
if (CollectionUtils.isEmpty(addresses)) {
continue;
}
@@ -623,7 +629,7 @@ public class IngressReconciler implements Reconciler {
.matchRestful(selectorData.getMatchRestful()).build();
shenyuCacheRepository.saveOrUpdateSelectorData(newSelectorData);
LOG.info("Updated upstream handle for selector {} of
plugin {} from endpoints {}/{}",
- selectorData.getId(), pluginName,
serviceNamespace, serviceName);
+ selectorData.getId(), pluginName, namespace,
serviceName);
}
}
}
@@ -657,7 +663,7 @@ public class IngressReconciler implements Reconciler {
return GsonUtils.getInstance().toJson(res);
}
- private List<Pair<V1EndpointAddress, String>> endpointAddresses(final
V1Endpoints v1Endpoints) {
+ private List<Pair<V1EndpointAddress, String>> endpointAddresses(final
V1Endpoints v1Endpoints, final IngressBackendPort backendPort) {
List<Pair<V1EndpointAddress, String>> res = new ArrayList<>();
List<V1EndpointSubset> subsets = v1Endpoints.getSubsets();
if (CollectionUtils.isNotEmpty(subsets)) {
@@ -667,10 +673,7 @@ public class IngressReconciler implements Reconciler {
if (CollectionUtils.isEmpty(ports) ||
CollectionUtils.isEmpty(addresses)) {
continue;
}
- CoreV1EndpointPort endpointPort = ports.stream()
- .filter(coreV1EndpointPort ->
"TCP".equals(coreV1EndpointPort.getProtocol()))
- .findFirst()
- .orElse(null);
+ CoreV1EndpointPort endpointPort =
IngressBackendPort.selectEndpointPort(ports, backendPort);
if (Objects.isNull(endpointPort)) {
continue;
}
diff --git
a/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/EndpointsReconcilerTest.java
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/EndpointsReconcilerTest.java
index 5d36eea3a9..516cbc43d9 100644
---
a/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/EndpointsReconcilerTest.java
+++
b/shenyu-kubernetes-controller/src/test/java/org/apache/shenyu/k8s/EndpointsReconcilerTest.java
@@ -28,21 +28,32 @@ import
io.kubernetes.client.openapi.models.V1EndpointSubsetBuilder;
import io.kubernetes.client.openapi.models.V1Endpoints;
import io.kubernetes.client.openapi.models.V1EndpointsBuilder;
import io.kubernetes.client.openapi.models.V1Ingress;
+import io.kubernetes.client.openapi.models.V1ServiceBackendPort;
import org.apache.shenyu.common.dto.SelectorData;
import org.apache.shenyu.common.enums.PluginEnum;
import org.apache.shenyu.k8s.cache.IngressSelectorCache;
import org.apache.shenyu.k8s.cache.ServiceIngressCache;
+import org.apache.shenyu.k8s.common.IngressBackendPort;
+import org.apache.shenyu.k8s.common.ServiceIngressRelation;
import org.apache.shenyu.k8s.reconciler.EndpointsReconciler;
import org.apache.shenyu.k8s.repository.ShenyuCacheRepository;
import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
+import java.util.Arrays;
import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.stream.Collectors;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.containsString;
+import static org.hamcrest.Matchers.not;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -51,55 +62,176 @@ import static org.mockito.Mockito.when;
*/
public final class EndpointsReconcilerTest {
+ private SharedIndexInformer<V1Ingress> ingressInformer;
+
+ private SharedIndexInformer<V1Endpoints> endpointsInformer;
+
+ private Indexer<V1Endpoints> endpointsIndexer;
+
+ private ShenyuCacheRepository shenyuCacheRepository;
+
+ @BeforeEach
+ public void init() {
+ ingressInformer = mock(SharedIndexInformer.class);
+ endpointsInformer = mock(SharedIndexInformer.class);
+ Indexer<V1Ingress> ingressIndexer = mock(Indexer.class);
+ endpointsIndexer = mock(Indexer.class);
+ when(ingressInformer.getIndexer()).thenReturn(ingressIndexer);
+ when(endpointsInformer.getIndexer()).thenReturn(endpointsIndexer);
+ shenyuCacheRepository = mock(ShenyuCacheRepository.class);
+
when(shenyuCacheRepository.findSelectorDataList(PluginEnum.DIVIDE.getName())).thenReturn(Collections.emptyList());
+ }
+
/**
* test websocket selector update.
*/
@Test
public void testUpdateWebSocketSelector() {
- SharedIndexInformer<V1Ingress> ingressInformer =
mock(SharedIndexInformer.class);
- SharedIndexInformer<V1Endpoints> endpointsInformer =
mock(SharedIndexInformer.class);
- Indexer<V1Ingress> ingressIndexer = mock(Indexer.class);
- Indexer<V1Endpoints> endpointsIndexer = mock(Indexer.class);
- when(ingressInformer.getIndexer()).thenReturn(ingressIndexer);
- when(endpointsInformer.getIndexer()).thenReturn(endpointsIndexer);
-
String namespace = "endpoint-websocket-ns";
String serviceName = "endpoint-websocket-service";
String ingressName = "endpoint-websocket-ingress";
String selectorId = "endpoint-websocket-selector";
+ mockEndpoints(namespace, serviceName, endpointPort(8001, null));
+ ServiceIngressCache.getInstance().putIngressName(namespace,
serviceName,
+ new ServiceIngressRelation(namespace, ingressName,
backendPort(8001)));
+ IngressSelectorCache.getInstance().put(namespace, ingressName,
PluginEnum.WEB_SOCKET.getName(), selectorId);
+ mockWebSocketSelectors(selectorId);
+
+ Result result = newReconciler().reconcile(new Request(namespace,
serviceName));
+
+ Assertions.assertEquals(new Result(false), result);
+ ArgumentCaptor<SelectorData> selectorCaptor =
ArgumentCaptor.forClass(SelectorData.class);
+
verify(shenyuCacheRepository).saveOrUpdateSelectorData(selectorCaptor.capture());
+ SelectorData updatedSelector = selectorCaptor.getValue();
+ Assertions.assertEquals(PluginEnum.WEB_SOCKET.getName(),
updatedSelector.getPluginName());
+ assertThat(updatedSelector.getHandle(),
containsString("\"protocol\":\"ws://\""));
+ assertThat(updatedSelector.getHandle(),
containsString("\"upstreamUrl\":\"127.0.0.1:8001\""));
+ }
+
+ /**
+ * test websocket selectors of two ingresses keep the service port each
ingress selects.
+ */
+ @Test
+ public void testUpdateMultiPortSelectors() {
+ String namespace = "endpoint-multi-port-ns";
+ String serviceName = "endpoint-multi-port-service";
+ String firstIngress = "endpoint-multi-port-ingress-1";
+ String secondIngress = "endpoint-multi-port-ingress-2";
+ mockEndpoints(namespace, serviceName, endpointPort(8001, null),
endpointPort(8002, null));
+ ServiceIngressCache.getInstance().putIngressName(namespace,
serviceName,
+ new ServiceIngressRelation(namespace, firstIngress,
backendPort(8001)));
+ ServiceIngressCache.getInstance().putIngressName(namespace,
serviceName,
+ new ServiceIngressRelation(namespace, secondIngress,
backendPort(8002)));
+ String firstSelectorId = "endpoint-multi-port-selector-1";
+ IngressSelectorCache.getInstance().put(namespace, firstIngress,
PluginEnum.WEB_SOCKET.getName(), firstSelectorId);
+ String secondSelectorId = "endpoint-multi-port-selector-2";
+ IngressSelectorCache.getInstance().put(namespace, secondIngress,
PluginEnum.WEB_SOCKET.getName(), secondSelectorId);
+ mockWebSocketSelectors(firstSelectorId, secondSelectorId);
+
+ Result result = newReconciler().reconcile(new Request(namespace,
serviceName));
+
+ Assertions.assertEquals(new Result(false), result);
+ Map<String, String> handleBySelectorId = captureUpdatedHandles(2);
+ assertThat(handleBySelectorId.get(firstSelectorId),
containsString("\"upstreamUrl\":\"127.0.0.1:8001\""));
+ assertThat(handleBySelectorId.get(firstSelectorId),
not(containsString(":8002")));
+ assertThat(handleBySelectorId.get(secondSelectorId),
containsString("\"upstreamUrl\":\"127.0.0.1:8002\""));
+ assertThat(handleBySelectorId.get(secondSelectorId),
not(containsString(":8001")));
+ }
+
+ /**
+ * test the service port can be selected by name.
+ */
+ @Test
+ public void testUpdateSelectorWithNamedBackendPort() {
+ String namespace = "endpoint-named-port-ns";
+ String serviceName = "endpoint-named-port-service";
+ String ingressName = "endpoint-named-port-ingress";
+ String selectorId = "endpoint-named-port-selector";
+ mockEndpoints(namespace, serviceName, endpointPort(8001, null),
endpointPort(8003, "ws"));
+ ServiceIngressCache.getInstance().putIngressName(namespace,
serviceName,
+ new ServiceIngressRelation(namespace, ingressName,
namedBackendPort("ws")));
+ IngressSelectorCache.getInstance().put(namespace, ingressName,
PluginEnum.WEB_SOCKET.getName(), selectorId);
+ mockWebSocketSelectors(selectorId);
+
+ Result result = newReconciler().reconcile(new Request(namespace,
serviceName));
+
+ Assertions.assertEquals(new Result(false), result);
+ Map<String, String> handleBySelectorId = captureUpdatedHandles(1);
+ assertThat(handleBySelectorId.get(selectorId),
containsString("\"upstreamUrl\":\"127.0.0.1:8003\""));
+ }
+
+ /**
+ * test the first TCP port is used when no endpoint port matches the
selected service port,
+ * a service may map the selected port to a different target port.
+ */
+ @Test
+ public void testUpdateSelectorWithUnmatchedBackendPort() {
+ String namespace = "endpoint-unmatched-port-ns";
+ String serviceName = "endpoint-unmatched-port-service";
+ String ingressName = "endpoint-unmatched-port-ingress";
+ String selectorId = "endpoint-unmatched-port-selector";
+ mockEndpoints(namespace, serviceName, endpointPort(8001, null),
endpointPort(8002, null));
+ ServiceIngressCache.getInstance().putIngressName(namespace,
serviceName,
+ new ServiceIngressRelation(namespace, ingressName,
backendPort(9000)));
+ IngressSelectorCache.getInstance().put(namespace, ingressName,
PluginEnum.WEB_SOCKET.getName(), selectorId);
+ mockWebSocketSelectors(selectorId);
+
+ Result result = newReconciler().reconcile(new Request(namespace,
serviceName));
+
+ Assertions.assertEquals(new Result(false), result);
+ Map<String, String> handleBySelectorId = captureUpdatedHandles(1);
+ assertThat(handleBySelectorId.get(selectorId),
containsString("\"upstreamUrl\":\"127.0.0.1:8001\""));
+ }
+
+ private EndpointsReconciler newReconciler() {
+ return new EndpointsReconciler(ingressInformer, endpointsInformer,
shenyuCacheRepository, mock(ApiClient.class));
+ }
+
+ private void mockEndpoints(final String namespace, final String
serviceName, final CoreV1EndpointPort... ports) {
V1Endpoints endpoints = new V1EndpointsBuilder().withKind("Endpoints")
.withNewMetadata().withNamespace(namespace).withName(serviceName).endMetadata()
.withSubsets(new V1EndpointSubsetBuilder()
.withAddresses(new V1EndpointAddress().ip("127.0.0.1"))
- .withPorts(new
CoreV1EndpointPort().port(8001).protocol("TCP"))
+ .withPorts(ports)
.build())
.build();
when(endpointsIndexer.getByKey(namespace + "/" +
serviceName)).thenReturn(endpoints);
+ }
- ServiceIngressCache.getInstance().putIngressName(namespace,
serviceName, namespace, ingressName);
- IngressSelectorCache.getInstance().put(namespace, ingressName,
PluginEnum.WEB_SOCKET.getName(), selectorId);
+ private void mockWebSocketSelectors(final String... selectorIds) {
+ List<SelectorData> selectorDataList = Arrays.stream(selectorIds)
+ .map(selectorId -> SelectorData.builder()
+ .id(selectorId)
+
.pluginId(String.valueOf(PluginEnum.WEB_SOCKET.getCode()))
+ .pluginName(PluginEnum.WEB_SOCKET.getName())
+ .name("/**")
+ .handle("[]")
+ .enabled(true)
+ .build())
+ .collect(Collectors.toList());
+
when(shenyuCacheRepository.findSelectorDataList(PluginEnum.WEB_SOCKET.getName())).thenReturn(selectorDataList);
+ }
- ShenyuCacheRepository shenyuCacheRepository =
mock(ShenyuCacheRepository.class);
- SelectorData selectorData = SelectorData.builder()
- .id(selectorId)
- .pluginId(String.valueOf(PluginEnum.WEB_SOCKET.getCode()))
- .pluginName(PluginEnum.WEB_SOCKET.getName())
- .name("/**")
- .handle("[]")
- .enabled(true)
- .build();
-
when(shenyuCacheRepository.findSelectorDataList(PluginEnum.DIVIDE.getName())).thenReturn(Collections.emptyList());
-
when(shenyuCacheRepository.findSelectorDataList(PluginEnum.WEB_SOCKET.getName())).thenReturn(Collections.singletonList(selectorData));
+ private Map<String, String> captureUpdatedHandles(final int
expectedUpdates) {
+ ArgumentCaptor<SelectorData> selectorCaptor =
ArgumentCaptor.forClass(SelectorData.class);
+ verify(shenyuCacheRepository,
times(expectedUpdates)).saveOrUpdateSelectorData(selectorCaptor.capture());
+ return selectorCaptor.getAllValues().stream()
+ .collect(Collectors.toMap(SelectorData::getId,
SelectorData::getHandle, (first, second) -> second));
+ }
+
+ private CoreV1EndpointPort endpointPort(final int port, final String name)
{
+ CoreV1EndpointPort endpointPort = new
CoreV1EndpointPort().port(port).protocol("TCP");
+ if (Objects.nonNull(name)) {
+ endpointPort.setName(name);
+ }
+ return endpointPort;
+ }
- EndpointsReconciler endpointsReconciler = new
EndpointsReconciler(ingressInformer, endpointsInformer, shenyuCacheRepository,
mock(ApiClient.class));
- Result result = endpointsReconciler.reconcile(new Request(namespace,
serviceName));
+ private IngressBackendPort backendPort(final int port) {
+ return IngressBackendPort.from(new
V1ServiceBackendPort().number(port));
+ }
- Assertions.assertEquals(new Result(false), result);
- ArgumentCaptor<SelectorData> selectorCaptor =
ArgumentCaptor.forClass(SelectorData.class);
-
verify(shenyuCacheRepository).saveOrUpdateSelectorData(selectorCaptor.capture());
- SelectorData updatedSelector = selectorCaptor.getValue();
- Assertions.assertEquals(PluginEnum.WEB_SOCKET.getName(),
updatedSelector.getPluginName());
- assertThat(updatedSelector.getHandle(),
containsString("\"protocol\":\"ws://\""));
- assertThat(updatedSelector.getHandle(),
containsString("\"upstreamUrl\":\"127.0.0.1:8001\""));
+ private IngressBackendPort namedBackendPort(final String name) {
+ return IngressBackendPort.from(new V1ServiceBackendPort().name(name));
}
}