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 6b85c6ddd6 fix(dubbo): offload reference initialization from request 
threads (#7273)
6b85c6ddd6 is described below

commit 6b85c6ddd6becf1262f0932bbe7a38bbf799fa2d
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 06:27:58 2026 +0800

    fix(dubbo): offload reference initialization from request threads (#7273)
---
 .../dubbo/proxy/ApacheDubboProxyService.java       | 19 ++++++
 .../dubbo/proxy/ApacheDubboProxyServiceTest.java   | 75 ++++++++++++++++++++++
 2 files changed, 94 insertions(+)

diff --git 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyService.java
 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyService.java
index 3dd5401f48..2fcab86f5b 100644
--- 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyService.java
+++ 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/main/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyService.java
@@ -43,9 +43,12 @@ import 
org.apache.shenyu.plugin.dubbo.common.param.DubboParamResolveService;
 import org.springframework.util.ObjectUtils;
 import org.springframework.web.server.ServerWebExchange;
 import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Schedulers;
 
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
 import java.util.concurrent.CompletableFuture;
@@ -79,6 +82,22 @@ public class ApacheDubboProxyService {
      * @throws ShenyuException the shenyu exception
      */
     public Mono<Object> genericInvoker(final String body, final MetaData 
metaData, final SelectorData selectorData, final RuleData ruleData, final 
ServerWebExchange exchange) throws ShenyuException {
+        Map<String, Object> attachments = new 
HashMap<>(RpcContext.getClientAttachment().getObjectAttachments());
+        return Mono.defer(() -> {
+            Map<String, Object> previous = new 
HashMap<>(RpcContext.getClientAttachment().getObjectAttachments());
+            try {
+                RpcContext.getClientAttachment().setObjectAttachments(new 
HashMap<>(attachments));
+                return invokeOnWorker(body, metaData, selectorData, ruleData, 
exchange);
+            } finally {
+                // Invocation and future lookup are synchronous; do not retain 
request data on a pooled worker.
+                RpcContext.getClientAttachment().clearAttachments();
+                
RpcContext.getClientAttachment().setObjectAttachments(previous);
+            }
+        })
+                .subscribeOn(Schedulers.boundedElastic());
+    }
+
+    private Mono<Object> invokeOnWorker(final String body, final MetaData 
metaData, final SelectorData selectorData, final RuleData ruleData, final 
ServerWebExchange exchange) {
         ReferenceConfig<GenericService> reference = 
this.getReferenceConfig(selectorData, ruleData, metaData, exchange);
         GenericService genericService = reference.get();
 
diff --git 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyServiceTest.java
 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyServiceTest.java
index f64c3d3408..cea85823ab 100644
--- 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyServiceTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-dubbo/shenyu-plugin-apache-dubbo/src/test/java/org/apache/shenyu/plugin/apache/dubbo/proxy/ApacheDubboProxyServiceTest.java
@@ -21,8 +21,11 @@ import com.google.common.cache.LoadingCache;
 import org.apache.commons.lang3.tuple.ImmutablePair;
 import org.apache.commons.lang3.tuple.Pair;
 import org.apache.dubbo.config.ReferenceConfig;
+import org.apache.dubbo.rpc.RpcContext;
+import org.apache.dubbo.rpc.RpcContextAttachment;
 import org.apache.dubbo.rpc.service.GenericService;
 import org.apache.shenyu.common.dto.MetaData;
+import org.apache.shenyu.common.constant.Constants;
 import org.apache.shenyu.common.dto.RuleData;
 import org.apache.shenyu.common.dto.SelectorData;
 import org.apache.shenyu.common.enums.RpcTypeEnum;
@@ -42,9 +45,18 @@ import org.springframework.web.server.ServerWebExchange;
 
 import java.lang.reflect.Field;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.atomic.AtomicReference;
+import reactor.core.publisher.Mono;
+import reactor.core.scheduler.Schedulers;
+import reactor.test.StepVerifier;
 
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 /**
  * The Test Case For ApacheDubboProxyService.
@@ -93,6 +105,7 @@ public final class ApacheDubboProxyServiceTest {
     @AfterEach
     public void after() {
         ApacheDubboConfigCache.getInstance().invalidateAll();
+        RpcContext.removeClientAttachment();
     }
 
     @Test
@@ -112,6 +125,68 @@ public final class ApacheDubboProxyServiceTest {
         future.complete("success");
     }
 
+    @Test
+    @SuppressWarnings("unchecked")
+    void defersReferenceAccessAndInvocationOffTheRequestThread() throws 
Exception {
+        GenericService genericService = mock(GenericService.class);
+        AtomicReference<Thread> worker = new AtomicReference<>();
+        AtomicReference<RpcContextAttachment> workerContext = new 
AtomicReference<>();
+        RpcContext.getClientAttachment().setAttachment("timeout", 1234);
+        
RpcContext.getClientAttachment().setAttachment(Constants.DUBBO_SELECTOR_ID, 
selectorData.getId());
+        
RpcContext.getClientAttachment().setAttachment(Constants.DUBBO_RULE_ID, 
ruleData.getId());
+        
RpcContext.getClientAttachment().setAttachment(Constants.DUBBO_REMOTE_ADDRESS, 
"127.0.0.1");
+        RpcContext.getClientAttachment().setAttachment("custom", "value");
+        when(referenceConfig.getInterface()).thenReturn(PATH);
+        when(referenceConfig.get()).thenAnswer(invocation -> {
+            assertFalse(Schedulers.isInNonBlockingThread());
+            worker.set(Thread.currentThread());
+            return genericService;
+        });
+        when(genericService.$invoke(METHOD_NAME, LEFT, 
RIGHT)).thenAnswer(invocation -> {
+            assertFalse(Schedulers.isInNonBlockingThread());
+            assertTrue(Thread.currentThread() == worker.get());
+            workerContext.set(RpcContext.getClientAttachment());
+            assertEquals(1234, 
workerContext.get().getObjectAttachments().get("timeout"));
+            assertEquals(selectorData.getId(), 
workerContext.get().getAttachment(Constants.DUBBO_SELECTOR_ID));
+            assertEquals(ruleData.getId(), 
workerContext.get().getAttachment(Constants.DUBBO_RULE_ID));
+            assertEquals("127.0.0.1", 
workerContext.get().getAttachment(Constants.DUBBO_REMOTE_ADDRESS));
+            assertEquals("value", workerContext.get().getAttachment("custom"));
+            return null;
+        });
+        Field field = ApacheDubboConfigCache.class.getDeclaredField("cache");
+        field.setAccessible(true);
+        ((LoadingCache<String, ReferenceConfig<GenericService>>) 
field.get(ApacheDubboConfigCache.getInstance())).put(PATH, referenceConfig);
+        ApacheDubboProxyService service = new ApacheDubboProxyService(new 
BodyParamResolveServiceImpl());
+        Mono<Object> result = service.genericInvoker("", metaData, 
selectorData, ruleData, exchange);
+        verifyNoInteractions(referenceConfig, genericService);
+
+        
StepVerifier.create(result.subscribeOn(Schedulers.parallel())).expectNext(Constants.DUBBO_RPC_RESULT_EMPTY).verifyComplete();
+
+        assertNotSame(Thread.currentThread(), worker.get());
+        assertTrue(workerContext.get().getObjectAttachments().isEmpty());
+        assertEquals("value", 
RpcContext.getClientAttachment().getAttachment("custom"));
+    }
+
+    @Test
+    @SuppressWarnings("unchecked")
+    void referenceInitializationErrorsAreReactive() throws Exception {
+        when(referenceConfig.getInterface()).thenReturn(PATH);
+        AtomicReference<RpcContextAttachment> workerContext = new 
AtomicReference<>();
+        RpcContext.getClientAttachment().setAttachment("custom", "value");
+        when(referenceConfig.get()).thenAnswer(invocation -> {
+            workerContext.set(RpcContext.getClientAttachment());
+            assertEquals("value", workerContext.get().getAttachment("custom"));
+            throw new IllegalStateException("registry unavailable");
+        });
+        Field field = ApacheDubboConfigCache.class.getDeclaredField("cache");
+        field.setAccessible(true);
+        ((LoadingCache<String, ReferenceConfig<GenericService>>) 
field.get(ApacheDubboConfigCache.getInstance())).put(PATH, referenceConfig);
+        ApacheDubboProxyService service = new ApacheDubboProxyService(new 
BodyParamResolveServiceImpl());
+        Mono<Object> result = service.genericInvoker("", metaData, 
selectorData, ruleData, exchange);
+        StepVerifier.create(result).expectErrorMessage("registry 
unavailable").verify();
+        assertTrue(workerContext.get().getObjectAttachments().isEmpty());
+    }
+
     static class BodyParamResolveServiceImpl implements 
DubboParamResolveService {
 
         @Override

Reply via email to