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 02af907610 fix (registry-nacos) : distinguish instance updates from 
additions and deletions. (#7270)
02af907610 is described below

commit 02af9076105b53c74829a4a18a390f9897ae84b9
Author: Jerry聊AI <[email protected]>
AuthorDate: Sat Sep 26 12:29:46 2026 +0800

    fix (registry-nacos) : distinguish instance updates from additions and 
deletions. (#7270)
    
    Co-authored-by: aias00 <[email protected]>
---
 .../nacos/NacosInstanceRegisterRepository.java     |  19 ++--
 .../nacos/NacosInstanceRegisterRepositoryTest.java | 104 ++++++++++++++++++++-
 2 files changed, 116 insertions(+), 7 deletions(-)

diff --git 
a/shenyu-registry/shenyu-registry-nacos/src/main/java/org/apache/shenyu/registry/nacos/NacosInstanceRegisterRepository.java
 
b/shenyu-registry/shenyu-registry-nacos/src/main/java/org/apache/shenyu/registry/nacos/NacosInstanceRegisterRepository.java
index d643dc3da3..d129521860 100644
--- 
a/shenyu-registry/shenyu-registry-nacos/src/main/java/org/apache/shenyu/registry/nacos/NacosInstanceRegisterRepository.java
+++ 
b/shenyu-registry/shenyu-registry-nacos/src/main/java/org/apache/shenyu/registry/nacos/NacosInstanceRegisterRepository.java
@@ -168,7 +168,7 @@ public class NacosInstanceRegisterRepository implements 
ShenyuInstanceRegisterRe
 
     private void compareInstances(final Set<Instance> previousInstances, final 
Set<Instance> currentInstances, final ChangedEventListener listener) {
         Set<Instance> addedInstances = currentInstances.stream()
-                .filter(item -> !previousInstances.contains(item))
+                .filter(item -> previousInstances.stream().noneMatch(previous 
-> isSameInstance(item, previous)))
                 .collect(Collectors.toSet());
         if (!addedInstances.isEmpty()) {
             for (Instance instance: addedInstances) {
@@ -177,7 +177,7 @@ public class NacosInstanceRegisterRepository implements 
ShenyuInstanceRegisterRe
         }
 
         Set<Instance> deletedInstances = previousInstances.stream()
-                .filter(item -> !currentInstances.contains(item))
+                .filter(item -> currentInstances.stream().noneMatch(current -> 
isSameInstance(item, current)))
                 .collect(Collectors.toSet());
         if (!deletedInstances.isEmpty()) {
             for (Instance instance: deletedInstances) {
@@ -188,10 +188,8 @@ public class NacosInstanceRegisterRepository implements 
ShenyuInstanceRegisterRe
 
         Set<Instance> updatedInstances = currentInstances.stream()
             .filter(
-                currentInstance -> 
Objects.nonNull(currentInstance.getInstanceId())
-                    && previousInstances.stream().anyMatch(
-                        previousInstance -> 
StringUtils.isNotBlank(previousInstance.getInstanceId())
-                        && 
currentInstance.getInstanceId().equals(previousInstance.getInstanceId())
+                currentInstance -> previousInstances.stream().anyMatch(
+                    previousInstance -> isSameInstance(currentInstance, 
previousInstance)
                         && !currentInstance.equals(previousInstance)))
             .collect(Collectors.toSet());
 
@@ -202,6 +200,15 @@ public class NacosInstanceRegisterRepository implements 
ShenyuInstanceRegisterRe
         }
     }
 
+    private boolean isSameInstance(final Instance current, final Instance 
previous) {
+        if (StringUtils.isNotBlank(current.getInstanceId()) && 
StringUtils.isNotBlank(previous.getInstanceId())) {
+            return current.getInstanceId().equals(previous.getInstanceId());
+        }
+        return Objects.equals(current.getIp(), previous.getIp())
+                && current.getPort() == previous.getPort()
+                && Objects.equals(current.getClusterName(), 
previous.getClusterName());
+    }
+
     private String buildUpstreamJsonFromInstance(final Instance instance) {
         JsonObject upstreamJson = new JsonObject();
         upstreamJson.addProperty("url", instance.getIp() + ":" + 
instance.getPort());
diff --git 
a/shenyu-registry/shenyu-registry-nacos/src/test/java/org/apache/shenyu/registry/nacos/NacosInstanceRegisterRepositoryTest.java
 
b/shenyu-registry/shenyu-registry-nacos/src/test/java/org/apache/shenyu/registry/nacos/NacosInstanceRegisterRepositoryTest.java
index 07c13bbdb5..6f7c137161 100644
--- 
a/shenyu-registry/shenyu-registry-nacos/src/test/java/org/apache/shenyu/registry/nacos/NacosInstanceRegisterRepositoryTest.java
+++ 
b/shenyu-registry/shenyu-registry-nacos/src/test/java/org/apache/shenyu/registry/nacos/NacosInstanceRegisterRepositoryTest.java
@@ -19,26 +19,44 @@ package org.apache.shenyu.registry.nacos;
 
 import com.alibaba.nacos.api.exception.NacosException;
 import com.alibaba.nacos.api.naming.NamingService;
+import com.alibaba.nacos.api.naming.listener.EventListener;
+import com.alibaba.nacos.api.naming.listener.NamingEvent;
 import com.alibaba.nacos.api.naming.pojo.Instance;
+import org.apache.shenyu.common.utils.GsonUtils;
 import org.apache.shenyu.registry.api.entity.InstanceEntity;
+import org.apache.shenyu.registry.api.event.ChangedEventListener;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.NullAndEmptySource;
+import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.ArgumentCaptor;
 
 import java.lang.reflect.Field;
+import java.util.Arrays;
+import java.util.Collections;
 import java.util.HashMap;
+import java.util.List;
 import java.util.Map;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.clearInvocations;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
+import static org.mockito.Mockito.when;
 
 public final class NacosInstanceRegisterRepositoryTest {
 
     private NacosInstanceRegisterRepository repository;
 
+    private NamingService namingService;
+
     private final Map<String, Instance> storage = new HashMap<>();
 
     @BeforeEach
@@ -48,7 +66,8 @@ public final class NacosInstanceRegisterRepositoryTest {
 
         Field field = clazz.getDeclaredField("namingService");
         field.setAccessible(true);
-        field.set(repository, mockNamingService());
+        namingService = mockNamingService();
+        field.set(repository, namingService);
 
         field = clazz.getDeclaredField("groupName");
         field.setAccessible(true);
@@ -100,4 +119,87 @@ public final class NacosInstanceRegisterRepositoryTest {
         repository.selectInstances(selectKey);
         repository.close();
     }
+
+    @ParameterizedTest
+    @NullAndEmptySource
+    @ValueSource(strings = {"instance-1"})
+    public void testWeightChangeOnlyPublishesUpdate(final String instanceId) 
throws NacosException {
+        Instance previous = newInstance(instanceId, "127.0.0.1");
+        Instance current = newInstance(instanceId, "127.0.0.1");
+        current.setWeight(100);
+        ChangedEventListener listener = 
notifyChange(Collections.singletonList(previous), 
Collections.singletonList(current));
+        ArgumentCaptor<String> payload = ArgumentCaptor.forClass(String.class);
+        verify(listener).onEvent(eq("service"), payload.capture(), 
eq(ChangedEventListener.Event.UPDATED));
+        assertEquals(100.0, 
GsonUtils.getInstance().fromJson(payload.getValue(), Map.class).get("weight"));
+        verifyNoMoreInteractions(listener);
+        assertTrue(previous.isHealthy());
+    }
+
+    @Test
+    public void testMetadataChangeOnlyPublishesUpdate() throws NacosException {
+        Instance previous = newInstance("instance-1", "127.0.0.1");
+        Instance current = newInstance("instance-1", "127.0.0.1");
+        current.setMetadata(Collections.singletonMap("props", "version-v2"));
+        ChangedEventListener listener = 
notifyChange(Collections.singletonList(previous), 
Collections.singletonList(current));
+        ArgumentCaptor<String> payload = ArgumentCaptor.forClass(String.class);
+        verify(listener).onEvent(eq("service"), payload.capture(), 
eq(ChangedEventListener.Event.UPDATED));
+        assertEquals("version-v2", 
GsonUtils.getInstance().fromJson(payload.getValue(), Map.class).get("props"));
+        verifyNoMoreInteractions(listener);
+    }
+
+    @Test
+    public void testUnchangedInstanceDoesNotPublishEvent() throws 
NacosException {
+        ChangedEventListener listener = 
notifyChange(Collections.singletonList(newInstance("instance-1", "127.0.0.1")),
+                Collections.singletonList(newInstance("instance-1", 
"127.0.0.1")));
+        verifyNoMoreInteractions(listener);
+    }
+
+    @Test
+    public void testAddedAndDeletedInstances() throws NacosException {
+        Instance unchanged = newInstance("instance-1", "127.0.0.1");
+        ChangedEventListener listener = notifyChange(Arrays.asList(unchanged, 
newInstance("instance-2", "127.0.0.2")),
+                Arrays.asList(unchanged, newInstance("instance-3", 
"127.0.0.3")));
+        verify(listener).onEvent(eq("service"), anyString(), 
eq(ChangedEventListener.Event.ADDED));
+        verify(listener).onEvent(eq("service"), anyString(), 
eq(ChangedEventListener.Event.DELETED));
+        verifyNoMoreInteractions(listener);
+    }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"ip", "port", "cluster"})
+    public void testDistinctInstancesWithoutIds(final String changedField) 
throws NacosException {
+        Instance previous = newInstance(null, "127.0.0.1");
+        Instance current = newInstance(null, "127.0.0.1");
+        if ("ip".equals(changedField)) {
+            current.setIp("127.0.0.2");
+        } else if ("port".equals(changedField)) {
+            current.setPort(8081);
+        } else {
+            current.setClusterName("another-cluster");
+        }
+        ChangedEventListener listener = 
notifyChange(Collections.singletonList(previous), 
Collections.singletonList(current));
+        verify(listener).onEvent(eq("service"), anyString(), 
eq(ChangedEventListener.Event.ADDED));
+        verify(listener).onEvent(eq("service"), anyString(), 
eq(ChangedEventListener.Event.DELETED));
+        verifyNoMoreInteractions(listener);
+    }
+
+    private ChangedEventListener notifyChange(final List<Instance> previous, 
final List<Instance> current) throws NacosException {
+        when(namingService.selectInstances("service", "group", 
true)).thenReturn(previous, current);
+        ChangedEventListener listener = mock(ChangedEventListener.class);
+        repository.watchInstances("service", listener);
+        ArgumentCaptor<EventListener> callback = 
ArgumentCaptor.forClass(EventListener.class);
+        verify(namingService).subscribe(eq("service"), eq("group"), 
callback.capture());
+        clearInvocations(listener);
+        callback.getValue().onEvent(new NamingEvent("service", current));
+        return listener;
+    }
+
+    private Instance newInstance(final String instanceId, final String ip) {
+        Instance instance = new Instance();
+        instance.setInstanceId(instanceId);
+        instance.setIp(ip);
+        instance.setPort(8080);
+        instance.setServiceName("service");
+        instance.setWeight(50);
+        return instance;
+    }
 }

Reply via email to