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]

Reply via email to