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 7964f36a92 fix(consul): stop cross-watch dedup dropping per-root
change events (#7329)
7964f36a92 is described below
commit 7964f36a92114d459e07959063720d4e9bfa06f8
Author: Sean-Walker0 <[email protected]>
AuthorDate: Sun Sep 27 12:34:02 2026 +0800
fix(consul): stop cross-watch dedup dropping per-root change events (#7329)
watchConfigKeyValues gated dispatch on
!consulIndexes.containsValue(newIndex), but consulIndexes holds one
entry per watch root and Consul's X-Consul-Index is a cluster-wide raft
index, so roots polled after the same write burst legitimately observe
the same new index. The first root stored it and every other root was
then silently skipped and had its index advanced anyway - its changed
keys were never dispatched and its removed keys never deleted, leaving
gateways stale until a full reload. Per-root advancement is already
checked immediately above, so the containsValue clause adds nothing
except the cross-root drop; the per-key modifyIndex/md5 comparisons
inside the loop already filter unchanged data.
The new test seeds two roots where one already stored the shared new
index, advances the other root, and fails on current master with an
empty dispatch list.
Co-authored-by: Sean-Walker0
<[email protected]>
Co-authored-by: aias00 <[email protected]>
---
.../sync/data/consul/ConsulSyncDataService.java | 3 +-
.../data/consul/ConsulSyncDataServiceTest.java | 55 ++++++++++++++++++++++
2 files changed, 56 insertions(+), 2 deletions(-)
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-consul/src/main/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataService.java
b/shenyu-sync-data-center/shenyu-sync-data-consul/src/main/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataService.java
index 8b9d1d7058..4bbd10db27 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-consul/src/main/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataService.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-consul/src/main/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataService.java
@@ -151,8 +151,7 @@ public class ConsulSyncDataService extends
AbstractPathDataSyncService {
-1, TimeUnit.MILLISECONDS);
return;
}
- if (!this.consulIndexes.containsValue(newIndex)
- &&
!currentIndex.equals(ConsulConstants.INIT_CONFIG_VERSION_INDEX)) {
+ if
(!currentIndex.equals(ConsulConstants.INIT_CONFIG_VERSION_INDEX)) {
if (LOG.isTraceEnabled()) {
LOG.trace("watchPathRoot {} has new index {}",
watchPathRoot, newIndex);
}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-consul/src/test/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataServiceTest.java
b/shenyu-sync-data-center/shenyu-sync-data-consul/src/test/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataServiceTest.java
index 606e78df39..ebdfdd67ad 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-consul/src/test/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataServiceTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-consul/src/test/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataServiceTest.java
@@ -34,15 +34,21 @@ import org.mockito.quality.Strictness;
import java.lang.reflect.Field;
import java.lang.reflect.Method;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
+import java.util.Base64;
+import java.util.Collections;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -147,4 +153,53 @@ public final class ConsulSyncDataServiceTest {
Assertions.assertTrue(consulIndexesField.get(consulSyncDataService)
instanceof ConcurrentMap);
Assertions.assertTrue(cacheDataField.get(consulSyncDataService)
instanceof ConcurrentMap);
}
+
+ @Test
+ public void
watchConfigKeyValuesShouldDispatchWhenAnotherRootSharesTheNewIndex() throws
Exception {
+ final Method watchConfigKeyValues =
ConsulSyncDataService.class.getDeclaredMethod("watchConfigKeyValues",
+ String.class, BiConsumer.class, Consumer.class);
+ watchConfigKeyValues.setAccessible(true);
+
+ final GetValue changed = new GetValue();
+ changed.setKey("/shenyu/selector/sel-1");
+
changed.setValue(Base64.getEncoder().encodeToString("changed".getBytes(StandardCharsets.UTF_8)));
+ changed.setModifyIndex(20L);
+ changed.setCreateIndex(10L);
+
+ final ConsulClient consulClientMock = mock(ConsulClient.class);
+ when(consulClientMock.getKVValues(anyString(), any(), any()))
+ .thenReturn(new Response<>(Collections.singletonList(changed),
20L, true, 1L));
+ final Field consulField =
ConsulSyncDataService.class.getDeclaredField("consulClient");
+ consulField.setAccessible(true);
+ consulField.set(consulSyncDataService, consulClientMock);
+
+ final ConsulConfig consulConfigMock = mock(ConsulConfig.class);
+ when(consulConfigMock.getWatchDelay()).thenReturn(1000);
+ when(consulConfigMock.getWaitTime()).thenReturn(1000);
+ final Field configField =
ConsulSyncDataService.class.getDeclaredField("consulConfig");
+ configField.setAccessible(true);
+ configField.set(consulSyncDataService, consulConfigMock);
+
+ final Field executorField =
ConsulSyncDataService.class.getDeclaredField("executor");
+ executorField.setAccessible(true);
+ executorField.set(consulSyncDataService,
mock(ScheduledThreadPoolExecutor.class));
+
+ @SuppressWarnings("unchecked")
+ final Map<String, Long> consulIndexes = (Map<String, Long>)
getConsulIndexesField().get(consulSyncDataService);
+ consulIndexes.put("/shenyu/plugin", 20L);
+ consulIndexes.put("/shenyu/selector", 10L);
+
+ final List<String> dispatched = new CopyOnWriteArrayList<>();
+ final BiConsumer<String, String> updateHandler = (key, value) ->
dispatched.add(key);
+ final Consumer<String> deleteHandler = key ->
dispatched.add("deleted:" + key);
+ watchConfigKeyValues.invoke(consulSyncDataService, "/shenyu/selector",
updateHandler, deleteHandler);
+ assertFalse(dispatched.isEmpty(), "a root whose own index advanced
must dispatch even when another root already stored that index");
+ }
+
+ private Field getConsulIndexesField() throws NoSuchFieldException {
+ final Field field =
ConsulSyncDataService.class.getDeclaredField("consulIndexes");
+ field.setAccessible(true);
+ return field;
+ }
+
}