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

Reply via email to