From f5633ee3b1b80d92877f6d92da4e9d67bf430c81 Mon Sep 17 00:00:00 2001 From: wy471x Date: Mon, 14 Sep 2026 21:29:33 +0800 Subject: [PATCH] fix: handle discovery upstream DELETE events in path data sync (#6661) --- .../core/AbstractPathDataSyncService.java | 22 ++++++++++++++----- .../core/AbstractPathDataSyncServiceTest.java | 22 +++++++++++++++++++ 2 files changed, 39 insertions(+), 5 deletions(-) diff --git a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncService.java b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncService.java index f9a3c2a0e9b4..10c86a1aa102 100644 --- a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncService.java +++ b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncService.java @@ -130,14 +130,21 @@ private void proxyHandlerEvent(final String updatePath, final String updateData, } private void discoveryUpstreamHandlerEvent(final String updatePath, final String updateData, final EventType eventType) { - String[] pathInfoArray2 = updatePath.split("/"); - if (pathInfoArray2.length != 5) { + String[] pathInfoArray = updatePath.split("/"); + if (pathInfoArray.length != 5) { return; } - if (!EventType.DELETE.equals(eventType)) { - Optional.ofNullable(updateData) - .ifPresent(e -> cacheDiscoveryUpstreamData(GsonUtils.getInstance().fromJson(updateData, DiscoverySyncData.class))); + String pluginName = pathInfoArray[pathInfoArray.length - 2]; + String selectorId = pathInfoArray[pathInfoArray.length - 1]; + if (EventType.DELETE.equals(eventType)) { + DiscoverySyncData discoverySyncData = new DiscoverySyncData(); + discoverySyncData.setPluginName(pluginName); + discoverySyncData.setSelectorId(selectorId); + unCacheDiscoveryUpstreamData(discoverySyncData); + return; } + Optional.ofNullable(updateData) + .ifPresent(e -> cacheDiscoveryUpstreamData(GsonUtils.getInstance().fromJson(updateData, DiscoverySyncData.class))); } private void ruleHandlerEvent(final String updatePath, final String updateData, final EventType eventType) { @@ -280,6 +287,11 @@ protected void cacheDiscoveryUpstreamData(final DiscoverySyncData upstreamDataLi .ifPresent(data -> discoveryUpstreamDataSubscribers.forEach(e -> e.onSubscribe(upstreamDataList))); } + protected void unCacheDiscoveryUpstreamData(final DiscoverySyncData discoverySyncData) { + Optional.ofNullable(discoverySyncData) + .ifPresent(data -> discoveryUpstreamDataSubscribers.forEach(e -> e.unSubscribe(data))); + } + protected void unCacheMetaData(final MetaData metaData) { Optional.ofNullable(metaData) .ifPresent(data -> metaDataSubscribers.forEach(e -> e.unSubscribe(metaData))); diff --git a/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncServiceTest.java b/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncServiceTest.java index bfb80e0fcc5f..c1c25e9deb8e 100644 --- a/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncServiceTest.java +++ b/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/core/AbstractPathDataSyncServiceTest.java @@ -17,7 +17,9 @@ package org.apache.shenyu.sync.data.core; +import org.apache.shenyu.common.constant.DefaultPathConstants; import org.apache.shenyu.common.dto.AppAuthData; +import org.apache.shenyu.common.dto.DiscoverySyncData; import org.apache.shenyu.common.utils.GsonUtils; import org.apache.shenyu.sync.data.api.AuthDataSubscriber; import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber; @@ -27,6 +29,7 @@ import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.MockitoAnnotations; import org.mockito.junit.MockitoJUnitRunner; @@ -34,6 +37,7 @@ import java.util.ArrayList; import java.util.List; +import static org.junit.Assert.assertEquals; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.verify; @@ -85,6 +89,24 @@ public void testUnCacheAuthData() { verify(authDataSubscriber).unSubscribe(any()); } + @Test + public void testDiscoveryUpstreamHandlerEvent() { + + String namespaceId = "/namespace"; + String registerPath = namespaceId + DefaultPathConstants.DISCOVERY_UPSTREAM; + String updatePath = registerPath + "/divide/testSelectorId"; + String jsonData = "{\"pluginName\":\"divide\",\"selectorId\":\"testSelectorId\",\"selectorName\":\"testSelector\"}"; + + pathDataSyncService.event(namespaceId, updatePath, jsonData, registerPath, AbstractPathDataSyncService.EventType.PUT); + verify(discoveryUpstreamDataSubscriber).onSubscribe(any()); + + pathDataSyncService.event(namespaceId, updatePath, null, registerPath, AbstractPathDataSyncService.EventType.DELETE); + ArgumentCaptor captor = ArgumentCaptor.forClass(DiscoverySyncData.class); + verify(discoveryUpstreamDataSubscriber).unSubscribe(captor.capture()); + assertEquals("divide", captor.getValue().getPluginName()); + assertEquals("testSelectorId", captor.getValue().getSelectorId()); + } + // Mock implementation static class AbstractPathDataSyncServiceImpl extends AbstractPathDataSyncService {