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 fe6463da5f fix(tars): publish complete proxy snapshots on refresh 
(#7225)
fe6463da5f is described below

commit fe6463da5f5563954218d5c2a5fdb98a24f959cf
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 11:58:44 2026 +0800

    fix(tars): publish complete proxy snapshots on refresh (#7225)
---
 .../plugin/tars/cache/ApplicationConfigCache.java  | 34 +++++---
 .../tars/cache/ApplicationConfigCacheTest.java     | 90 ++++++++++++++++++++++
 2 files changed, 113 insertions(+), 11 deletions(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-tars/src/main/java/org/apache/shenyu/plugin/tars/cache/ApplicationConfigCache.java
 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-tars/src/main/java/org/apache/shenyu/plugin/tars/cache/ApplicationConfigCache.java
index aba526b390..3d8a6f8379 100644
--- 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-tars/src/main/java/org/apache/shenyu/plugin/tars/cache/ApplicationConfigCache.java
+++ 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-tars/src/main/java/org/apache/shenyu/plugin/tars/cache/ApplicationConfigCache.java
@@ -186,6 +186,7 @@ public final class ApplicationConfigCache {
      * @param metaData metaData
      */
     public void initPrx(final MetaData metaData) {
+        LOCK.lock();
         try {
             if (Objects.isNull(prxClassCache.get(metaData.getPath()))) {
                 lockedLoadMetaData(metaData);
@@ -195,6 +196,8 @@ public final class ApplicationConfigCache {
             }
         } catch (Exception e) {
             LOG.error("ShenyuTarsPluginInitializeException: init tars ref 
ex:{}", e.getMessage());
+        } finally {
+            LOCK.unlock();
         }
     }
 
@@ -291,6 +294,7 @@ public final class ApplicationConfigCache {
      * @param selectorData selectorData
      */
     public void initPrxClass(final SelectorData selectorData) {
+        LOCK.lock();
         try {
             final List<TarsUpstream> upstreamList = 
GsonUtils.getInstance().fromList(selectorData.getHandle(), TarsUpstream.class);
             if (CollectionUtils.isEmpty(upstreamList)) {
@@ -304,6 +308,8 @@ public final class ApplicationConfigCache {
             }
         } catch (ExecutionException | NoSuchMethodException e) {
             throw new ShenyuException(e.getCause());
+        } finally {
+            LOCK.unlock();
         }
     }
     
@@ -318,8 +324,8 @@ public final class ApplicationConfigCache {
         if (Objects.isNull(prxClass)) {
             return;
         }
-        TarsInvokePrxList tarsInvokePrxList = cache.get(metaData.getPath());
-        tarsInvokePrxList.getTarsInvokePrxList().clear();
+        TarsInvokePrxList previous = cache.get(metaData.getPath());
+        TarsInvokePrxList tarsInvokePrxList = new 
TarsInvokePrxList(previous.getMethod(), previous.getParamTypes(), 
previous.getParamNames());
         if (Objects.isNull(tarsInvokePrxList.getMethod())) {
             TarsParamInfo tarsParamInfo = 
prxParamCache.get(getClassMethodKey(prxClass.getName(), 
metaData.getMethodName()));
             Object prx = communicator.stringToProxy(prxClass, 
PrxInfoUtil.getObjectName(upstreamList.get(0).getUpstreamUrl(), 
metaData.getServiceName()));
@@ -333,6 +339,7 @@ public final class ApplicationConfigCache {
             Object strProxy = communicator.stringToProxy(prxClass, 
PrxInfoUtil.getObjectName(upstream.getUpstreamUrl(), 
metaData.getServiceName()));
             return new TarsInvokePrx(strProxy, upstream.getUpstreamUrl());
         }).collect(Collectors.toList()));
+        cache.put(metaData.getPath(), tarsInvokePrxList);
     }
     
     /**
@@ -341,16 +348,21 @@ public final class ApplicationConfigCache {
      * @param contextPath context path
      */
     public void invalidate(final String contextPath) {
-        List<MetaData> metaDataList = ctxPathCache.remove(contextPath);
-        if (CollectionUtils.isNotEmpty(metaDataList)) {
-            metaDataList.forEach(metaData -> {
-                cache.invalidate(metaData.getPath());
-                prxClassCache.remove(metaData.getPath());
-                String paramKeyPrefix = PrxInfoUtil.getPrxName(metaData) + "_";
-                prxParamCache.keySet().removeIf(key -> 
key.startsWith(paramKeyPrefix));
-            });
+        LOCK.lock();
+        try {
+            refreshUpstreamCache.remove(contextPath);
+            List<MetaData> metaDataList = ctxPathCache.remove(contextPath);
+            if (CollectionUtils.isNotEmpty(metaDataList)) {
+                metaDataList.forEach(metaData -> {
+                    cache.invalidate(metaData.getPath());
+                    prxClassCache.remove(metaData.getPath());
+                    String paramKeyPrefix = PrxInfoUtil.getPrxName(metaData) + 
"_";
+                    prxParamCache.keySet().removeIf(key -> 
key.startsWith(paramKeyPrefix));
+                });
+            }
+        } finally {
+            LOCK.unlock();
         }
-        refreshUpstreamCache.remove(contextPath);
     }
     
     /**
diff --git 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-tars/src/test/java/org/apache/shenyu/plugin/tars/cache/ApplicationConfigCacheTest.java
 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-tars/src/test/java/org/apache/shenyu/plugin/tars/cache/ApplicationConfigCacheTest.java
index d4ee86e2b5..2ccc99dda4 100644
--- 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-tars/src/test/java/org/apache/shenyu/plugin/tars/cache/ApplicationConfigCacheTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-tars/src/test/java/org/apache/shenyu/plugin/tars/cache/ApplicationConfigCacheTest.java
@@ -18,17 +18,25 @@
 package org.apache.shenyu.plugin.tars.cache;
 
 import com.qq.tars.protocol.annotation.Servant;
+import com.qq.tars.client.Communicator;
 import org.apache.shenyu.common.concurrent.ShenyuThreadFactory;
 import org.apache.shenyu.common.constant.Constants;
 import org.apache.shenyu.common.dto.MetaData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.plugin.tars.handler.TarsPluginDataHandler;
 import org.apache.shenyu.common.enums.RpcTypeEnum;
+import org.apache.shenyu.common.dto.convert.selector.TarsUpstream;
+import org.apache.shenyu.plugin.tars.proxy.TarsInvokePrx;
 import org.apache.shenyu.plugin.tars.proxy.TarsInvokePrxList;
 import org.apache.shenyu.plugin.tars.util.PrxInfoUtil;
 import org.assertj.core.util.Lists;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.test.util.ReflectionTestUtils;
 
 import java.lang.reflect.Field;
 import java.util.Arrays;
@@ -46,6 +54,12 @@ import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNotSame;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
 
 /**
  * Test case for {@link ApplicationConfigCache}.
@@ -175,6 +189,45 @@ public final class ApplicationConfigCacheTest {
         assertNotNull(result);
     }
 
+    @Test
+    @SuppressWarnings("unchecked")
+    public void testRefreshPublishesCompleteSnapshot() throws Exception {
+        final String path = "snapshot-refresh";
+        final Map<String, Class<?>> classes = (Map<String, Class<?>>) 
ReflectionTestUtils.getField(applicationConfigCacheUnderTest, "prxClassCache");
+        final Communicator original = (Communicator) 
ReflectionTestUtils.getField(applicationConfigCacheUnderTest, "communicator");
+        final Communicator communicator = mock(Communicator.class);
+        final TarsInvokePrxList previous = 
applicationConfigCacheUnderTest.get(path);
+        previous.setMethod(Object.class.getMethod("toString"));
+        previous.addTarsInvokePrxList(Collections.singletonList(new 
TarsInvokePrx(new Object(), "old")));
+        final MetaData metadata = new MetaData();
+        metadata.setPath(path);
+        metadata.setServiceName("service");
+        final TarsUpstream upstream = 
TarsUpstream.builder().upstreamUrl("127.0.0.1:8080").build();
+        classes.put(path, Object.class);
+        ReflectionTestUtils.setField(applicationConfigCacheUnderTest, 
"communicator", communicator);
+        try {
+            when(communicator.stringToProxy(eq(Object.class), 
anyString())).thenAnswer(invocation -> {
+                assertSame(previous, 
applicationConfigCacheUnderTest.get(path));
+                assertEquals(1, previous.getTarsInvokePrxList().size());
+                return new Object();
+            });
+            ReflectionTestUtils.invokeMethod(applicationConfigCacheUnderTest, 
"refreshTarsInvokePrxList", metadata, Collections.singletonList(upstream));
+            assertNotSame(previous, applicationConfigCacheUnderTest.get(path));
+            assertEquals(1, previous.getTarsInvokePrxList().size());
+            assertEquals("old", 
previous.getTarsInvokePrxList().get(0).getHost());
+            assertEquals("127.0.0.1:8080", 
applicationConfigCacheUnderTest.get(path).getTarsInvokePrxList().get(0).getHost());
+            final TarsInvokePrxList current = 
applicationConfigCacheUnderTest.get(path);
+            org.mockito.Mockito.doThrow(new IllegalStateException("proxy 
unavailable")).when(communicator).stringToProxy(eq(Object.class), anyString());
+            assertThrows(IllegalStateException.class, () -> 
ReflectionTestUtils.invokeMethod(applicationConfigCacheUnderTest,
+                    "refreshTarsInvokePrxList", metadata, 
Collections.singletonList(upstream)));
+            assertSame(current, applicationConfigCacheUnderTest.get(path));
+            assertEquals(1, current.getTarsInvokePrxList().size());
+        } finally {
+            classes.remove(path);
+            ReflectionTestUtils.setField(applicationConfigCacheUnderTest, 
"communicator", original);
+        }
+    }
+
     @Test
     @SuppressWarnings("unchecked")
     public void testInvalidateRemovesCompanionCaches() throws Exception {
@@ -205,6 +258,43 @@ public final class ApplicationConfigCacheTest {
         assertNotSame(cached, 
applicationConfigCacheUnderTest.get(metaData.getPath()));
     }
 
+    @ParameterizedTest
+    @ValueSource(booleans = {true, false})
+    @SuppressWarnings("unchecked")
+    void deletedUpstreamsCannotBeRepublishedByMetadata(final boolean 
emptyUpdate) throws Exception {
+        final String context = "/deleted" + emptyUpdate;
+        final MetaData metadata = new MetaData("id", "app", context, context + 
"/path", RpcTypeEnum.TARS.getName(),
+                "deletedService" + emptyUpdate, "method", "", 
"{\"methodInfo\":[]}", false, Constants.SYS_DEFAULT_NAMESPACE_ID);
+        final Map<String, List<MetaData>> contexts = (Map<String, 
List<MetaData>>) getField("ctxPathCache");
+        final Map<String, Class<?>> classes = (Map<String, Class<?>>) 
getField("prxClassCache");
+        final Map<String, List<TarsUpstream>> upstreams = (Map<String, 
List<TarsUpstream>>) getField("refreshUpstreamCache");
+        final Communicator original = (Communicator) getField("communicator");
+        final Communicator communicator = mock(Communicator.class);
+        contexts.put(context, List.of(metadata));
+        classes.put(metadata.getPath(), Object.class);
+        upstreams.put(context, 
List.of(TarsUpstream.builder().upstreamUrl("127.0.0.1:8080").build()));
+        
applicationConfigCacheUnderTest.get(metadata.getPath()).addTarsInvokePrxList(List.of(new
 TarsInvokePrx(new Object(), "old")));
+        ReflectionTestUtils.setField(applicationConfigCacheUnderTest, 
"communicator", communicator);
+        try {
+            SelectorData selector = new SelectorData();
+            selector.setName(context);
+            selector.setHandle("[]");
+            if (emptyUpdate) {
+                applicationConfigCacheUnderTest.initPrxClass(selector);
+            } else {
+                new TarsPluginDataHandler().removeSelector(selector);
+            }
+            assertFalse(upstreams.containsKey(context));
+            applicationConfigCacheUnderTest.initPrx(metadata);
+            assertTrue(classes.containsKey(metadata.getPath()), "Metadata must 
initialize successfully after deletion");
+            
assertTrue(applicationConfigCacheUnderTest.get(metadata.getPath()).getTarsInvokePrxList().isEmpty());
+            verifyNoInteractions(communicator);
+        } finally {
+            applicationConfigCacheUnderTest.invalidate(context);
+            ReflectionTestUtils.setField(applicationConfigCacheUnderTest, 
"communicator", original);
+        }
+    }
+
     private Object getField(final String fieldName) throws 
NoSuchFieldException, IllegalAccessException {
         java.lang.reflect.Field field = 
ApplicationConfigCache.class.getDeclaredField(fieldName);
         field.setAccessible(true);

Reply via email to