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);