Aias00 commented on code in PR #7287:
URL: https://github.com/apache/shenyu/pull/7287#discussion_r4109930837
##########
shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/cache/ServiceIngressCache.java:
##########
@@ -47,52 +48,57 @@ public static ServiceIngressCache getInstance() {
}
/**
- * 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;
Review Comment:
Blocking: getIngressName() returns this exact list instance to the caller,
and it is mutated here (removeIf + add) while the endpoints controller is
iterating it. IngressionControllerConfiguration builds the ingress controller
with withWorkerCount(2) and the endpoints controller with withWorkerCount(2)
(IngressControllerConfiguration.java:102 and :141), so
EndpointsReconciler#reconcile iterating getIngressName(...) races with
IngressReconciler#reconcile calling putIngressName.
I reproduced it locally with the same structure (one thread iterating the
list returned by the getter, another doing removeIf + add inside
ConcurrentHashMap#compute):
instance identity: list returned by get() == list inside map -> true
reader -> java.util.ConcurrentModificationException
Please switch the per-key list to CopyOnWriteArrayList:
List<ServiceIngressRelation> res = Objects.isNull(relations) ? new
CopyOnWriteArrayList<>() : relations;
The list is read much more often than written, so this is the cheapest way
to close it - and it also fixes the same hazard that exists on master today,
while you are rewriting these accessors anyway.
##########
shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/IngressReconciler.java:
##########
@@ -379,45 +385,46 @@ private List<String> deleteSelectorByIngressName(final
String namespace, final S
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;
}
for (V1IngressRule rule : ingress.getSpec().getRules()) {
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()));
Review Comment:
Question (please decide before merge): keying this map by service name and
using putIfAbsent means one ingress that routes two paths to ports 8001 and
8002 of the same service still collapses into a single relation, and
IngressSelectorCache#get(namespace, ingressName, plugin) returns every selector
id of the ingress - so both selectors still receive the same upstream handle,
i.e. the exact misrouting of #6493, just inside one ingress instead of across
two. Your new tests use two separate ingresses, so if the reporter had a single
ingress with two paths this patch does not change their outcome.
Either extend ServiceIngressRelation with something that identifies the path
and filter the selectors accordingly (divide/websocket selectors are named
after their path, so it looks feasible), or keep the behaviour and state the
limitation explicitly in this javadoc plus a test that pins the current result.
Both are fine for me, I only want the claim in the PR title/description to
match what happens.
##########
shenyu-kubernetes-controller/src/main/java/org/apache/shenyu/k8s/reconciler/EndpointsReconciler.java:
##########
@@ -105,54 +106,74 @@ public Result reconcile(final Request request) {
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;
}
List<SelectorData> totalSelectors =
shenyuCacheRepository.findSelectorDataList(pluginName);
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)) {
Review Comment:
Non-blocking: when no address can be resolved you skip the selector
entirely, whereas on master the handle was always rewritten - a scale-to-zero
used to produce an empty upstream list and now keeps the previous, already dead
upstreams. Given the intentional choice documented in reconcile (deleting all
upstreams is treated as equivalent to leaving them), this is acceptable, but
please log it at WARN instead of INFO so a permanently unresolved backend port
does not stay invisible in the logs.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]