This is an automated email from the ASF dual-hosted git repository.

yuluo-yx 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 34ec5f73f8 fix(sync): process zookeeper delete events (#7102)
34ec5f73f8 is described below

commit 34ec5f73f8bc7113dfdec6f61cf4761eea323cab
Author: Liming Deng <[email protected]>
AuthorDate: Thu Sep 24 10:03:06 2026 +0800

    fix(sync): process zookeeper delete events (#7102)
    
    Co-authored-by: aias00 <[email protected]>
    Co-authored-by: shown <[email protected]>
---
 .../data/zookeeper/ZookeeperSyncDataService.java   | 73 +++++++++++++---------
 .../zookeeper/ZookeeperSyncDataServiceTest.java    | 25 ++++++++
 2 files changed, 68 insertions(+), 30 deletions(-)

diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/main/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataService.java
 
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/main/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataService.java
index 7bda998391..dab1b2c2aa 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/main/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataService.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/main/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataService.java
@@ -19,6 +19,8 @@ package org.apache.shenyu.sync.data.zookeeper;
 
 import com.google.common.base.Strings;
 import org.apache.commons.lang3.StringUtils;
+import org.apache.curator.framework.recipes.cache.ChildData;
+import org.apache.curator.framework.recipes.cache.CuratorCacheListener;
 import org.apache.shenyu.common.config.ShenyuConfig;
 import org.apache.shenyu.common.constant.Constants;
 import org.apache.shenyu.common.constant.DefaultPathConstants;
@@ -82,37 +84,48 @@ public class ZookeeperSyncDataService extends 
AbstractPathDataSyncService {
 
     private void watcherData0(final String registerPath) {
         String configNamespace = Constants.PATH_SEPARATOR + 
shenyuConfig.getNamespace();
-        zkClient.addCuratorCache(registerPath, (type, oldData, data) -> {
-            if (Objects.isNull(data) || Objects.isNull(data.getData())) {
-                return;
-            }
-            String path = data.getPath();
-            if (Strings.isNullOrEmpty(path)) {
-                return;
-            }
-            // if not uri register path, return.
-            if (!path.contains(registerPath)) {
-                return;
-            }
-            if (!StringUtils.containsIgnoreCase(path, configNamespace)) {
-                return;
-            }
+        zkClient.addCuratorCache(registerPath,
+                (type, oldData, data) -> handleEvent(type, oldData, data, 
registerPath, configNamespace));
+    }
 
-            EventType eventType = EventType.PUT;
-            switch (type) {
-                case NODE_DELETED:
-                    eventType = EventType.DELETE;
-                    break;
-                case NODE_CREATED:
-                case NODE_CHANGED:
-                    eventType = EventType.PUT;
-                    break;
-                default:
-                    break;
-            }
-            final String updateData = new String(data.getData(), 
StandardCharsets.UTF_8);
-            this.event(configNamespace, path, updateData, registerPath, 
eventType);
-        });
+    private void handleEvent(final CuratorCacheListener.Type type, final 
ChildData oldData, final ChildData data,
+                             final String registerPath, final String 
configNamespace) {
+        final ChildData eventData;
+        final EventType eventType;
+        final String updateData;
+        switch (type) {
+            case NODE_DELETED:
+                eventData = oldData;
+                eventType = EventType.DELETE;
+                updateData = null;
+                break;
+            case NODE_CREATED:
+            case NODE_CHANGED:
+                eventData = data;
+                eventType = EventType.PUT;
+                if (Objects.isNull(eventData) || 
Objects.isNull(eventData.getData())) {
+                    return;
+                }
+                updateData = new String(eventData.getData(), 
StandardCharsets.UTF_8);
+                break;
+            default:
+                return;
+        }
+        if (Objects.isNull(eventData)) {
+            return;
+        }
+        String path = eventData.getPath();
+        if (Strings.isNullOrEmpty(path)) {
+            return;
+        }
+        // if not uri register path, return.
+        if (!path.contains(registerPath)) {
+            return;
+        }
+        if (!StringUtils.containsIgnoreCase(path, configNamespace)) {
+            return;
+        }
+        this.event(configNamespace, path, updateData, registerPath, eventType);
     }
 
     @Override
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
index b572c99b1c..78d906d00f 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
@@ -19,6 +19,7 @@ package org.apache.shenyu.sync.data.zookeeper;
 
 import org.apache.curator.framework.CuratorFramework;
 import org.apache.curator.framework.recipes.cache.ChildData;
+import org.apache.curator.framework.recipes.cache.CuratorCacheListener;
 import org.apache.curator.framework.recipes.cache.TreeCacheEvent;
 import org.apache.curator.framework.recipes.cache.TreeCacheListener;
 import org.apache.shenyu.common.config.ShenyuConfig;
@@ -30,6 +31,7 @@ import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
 import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
 import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
 import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
 
 import java.util.ArrayList;
 import java.util.Collections;
@@ -37,12 +39,35 @@ import java.util.List;
 import java.util.Objects;
 
 import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.argThat;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 public final class ZookeeperSyncDataServiceTest {
 
+    @Test
+    public void testDeletedNodeUsesOldData() {
+        ZookeeperClient zkClient = mock(ZookeeperClient.class);
+        PluginDataSubscriber pluginDataSubscriber = 
mock(PluginDataSubscriber.class);
+        ShenyuConfig shenyuConfig = mock(ShenyuConfig.class);
+        
when(shenyuConfig.getNamespace()).thenReturn(Constants.SYS_DEFAULT_NAMESPACE_ID);
+        new ZookeeperSyncDataService(shenyuConfig, zkClient, 
pluginDataSubscriber, Collections.emptyList(),
+                Collections.emptyList(), Collections.emptyList(), 
Collections.emptyList());
+        ArgumentCaptor<CuratorCacheListener> listenerCaptor = 
ArgumentCaptor.forClass(CuratorCacheListener.class);
+        verify(zkClient, times(7)).addCuratorCache(anyString(), 
listenerCaptor.capture());
+
+        String pluginPath = Constants.PATH_SEPARATOR + 
Constants.SYS_DEFAULT_NAMESPACE_ID + "/shenyu/plugin/divide";
+        ChildData oldData = new ChildData(pluginPath, null, "{}".getBytes());
+        listenerCaptor.getAllValues().forEach(listener ->
+                listener.event(CuratorCacheListener.Type.NODE_DELETED, 
oldData, null));
+
+        verify(pluginDataSubscriber).unSubscribe(argThat(pluginData -> 
"divide".equals(pluginData.getName())));
+    }
+
     @Test
     public void testZookeeperInstanceRegisterRepository() throws Exception {
 

Reply via email to