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;
+ }
}