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 60c23ff3f8 fix(grpc): complete client calls asynchronously (#7103)
60c23ff3f8 is described below

commit 60c23ff3f8b5003617e4d3b7c5c2de31ea3447e4
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 06:29:30 2026 +0800

    fix(grpc): complete client calls asynchronously (#7103)
---
 .../plugin/grpc/client/ShenyuGrpcClient.java       | 39 ++++++++++--------
 .../plugin/grpc/client/ShenyuGrpcClientTest.java   | 46 ++++++++++++++++++++++
 2 files changed, 69 insertions(+), 16 deletions(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/main/java/org/apache/shenyu/plugin/grpc/client/ShenyuGrpcClient.java
 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/main/java/org/apache/shenyu/plugin/grpc/client/ShenyuGrpcClient.java
index 15bd89763b..52fb41a8c5 100644
--- 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/main/java/org/apache/shenyu/plugin/grpc/client/ShenyuGrpcClient.java
+++ 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/main/java/org/apache/shenyu/plugin/grpc/client/ShenyuGrpcClient.java
@@ -17,7 +17,10 @@
 
 package org.apache.shenyu.plugin.grpc.client;
 
+import com.google.common.util.concurrent.FutureCallback;
+import com.google.common.util.concurrent.Futures;
 import com.google.common.util.concurrent.ListenableFuture;
+import com.google.common.util.concurrent.MoreExecutors;
 import com.google.protobuf.DynamicMessage;
 import io.grpc.CallOptions;
 import io.grpc.ClientCall;
@@ -42,7 +45,6 @@ import java.io.Closeable;
 import java.util.Collections;
 import java.util.List;
 import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.ExecutionException;
 
 import static io.grpc.stub.ClientCalls.asyncServerStreamingCall;
 import static io.grpc.stub.ClientCalls.asyncUnaryCall;
@@ -99,21 +101,26 @@ public class ShenyuGrpcClient implements Closeable {
         callParams.setResponseObserver(streamObserver);
         callParams.setRequests(jsonRequestList);
         
-        try {
-            this.invoke(callParams).get();
-        } catch (InterruptedException e) {
-            // InterruptedExceptions should never be ignored in the code.
-            // InterruptedExceptions should either be rethrown - immediately 
or after cleaning up the method’s state -
-            // or the thread should be re-interrupted by calling 
Thread.interrupt() even if this is supposed to be a single-threaded application.
-            // Any other course of action risks delaying thread shutdown and 
loses the information
-            // that the thread was interrupted - probably without finishing 
its task.
-            LOG.error("Grpc plugin invoke method is exception, Will cause the 
thread to be interrupted");
-            Thread.currentThread().interrupt();
-            throw new ShenyuGrpcException("Caught exception while waiting for 
rpc :{ " + e.getMessage() + "}", e);
-        } catch (ExecutionException e) {
-            throw new ShenyuGrpcException("Caught exception while waiting for 
rpc :{ " + e.getMessage() + "}", e);
-        }
-        return CompletableFuture.completedFuture(shenyuGrpcResponse);
+        ListenableFuture<Void> invocation = this.invoke(callParams);
+        CompletableFuture<ShenyuGrpcResponse> result = new 
CompletableFuture<>();
+        Futures.addCallback(invocation, new FutureCallback<>() {
+            @Override
+            public void onSuccess(final Void ignored) {
+                result.complete(shenyuGrpcResponse);
+            }
+
+            @Override
+            public void onFailure(final Throwable throwable) {
+                result.completeExceptionally(new ShenyuGrpcException(
+                        "Caught exception while waiting for rpc :{ " + 
throwable.getMessage() + "}", throwable));
+            }
+        }, MoreExecutors.directExecutor());
+        result.whenComplete((ignored, throwable) -> {
+            if (result.isCancelled()) {
+                invocation.cancel(true);
+            }
+        });
+        return result;
     }
     
     /**
diff --git 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/test/java/org/apache/shenyu/plugin/grpc/client/ShenyuGrpcClientTest.java
 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/test/java/org/apache/shenyu/plugin/grpc/client/ShenyuGrpcClientTest.java
index 3ece0d3d36..e20c3bbeba 100644
--- 
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/test/java/org/apache/shenyu/plugin/grpc/client/ShenyuGrpcClientTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-rpc/shenyu-plugin-grpc/src/test/java/org/apache/shenyu/plugin/grpc/client/ShenyuGrpcClientTest.java
@@ -17,15 +17,23 @@
 
 package org.apache.shenyu.plugin.grpc.client;
 
+import com.google.common.util.concurrent.SettableFuture;
 import com.google.common.util.concurrent.Futures;
 import io.grpc.CallOptions;
 import io.grpc.ManagedChannel;
 import io.grpc.MethodDescriptor;
 import org.apache.shenyu.common.dto.MetaData;
 import org.apache.shenyu.plugin.grpc.proto.ShenyuGrpcCallRequest;
+import org.apache.shenyu.plugin.grpc.proto.ShenyuGrpcResponse;
 import org.junit.jupiter.api.Test;
 import org.mockito.ArgumentCaptor;
 
+import java.time.Duration;
+import java.util.concurrent.CompletableFuture;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.doReturn;
@@ -38,6 +46,44 @@ import static org.mockito.Mockito.verify;
  */
 public final class ShenyuGrpcClientTest {
 
+    @Test
+    public void testCallCompletesAsynchronously() {
+        ShenyuGrpcClient client = spy(new 
ShenyuGrpcClient(mock(ManagedChannel.class)));
+        SettableFuture<Void> invocation = SettableFuture.create();
+        
doReturn(invocation).when(client).invoke(any(ShenyuGrpcCallRequest.class));
+        MetaData metaData = MetaData.builder()
+                .serviceName("echo.EchoService")
+                .methodName("echo")
+                .build();
+
+        CompletableFuture<ShenyuGrpcResponse> result = 
assertTimeoutPreemptively(Duration.ofSeconds(1),
+                () -> client.call(metaData, CallOptions.DEFAULT,
+                        "{\"data\":[{}]}", MethodDescriptor.MethodType.UNARY));
+
+        assertFalse(result.isDone());
+        invocation.set(null);
+        assertTrue(result.isDone());
+        assertFalse(result.isCompletedExceptionally());
+    }
+
+    @Test
+    public void testCancellationPropagatesToInvocation() {
+        ShenyuGrpcClient client = spy(new 
ShenyuGrpcClient(mock(ManagedChannel.class)));
+        SettableFuture<Void> invocation = SettableFuture.create();
+        
doReturn(invocation).when(client).invoke(any(ShenyuGrpcCallRequest.class));
+        MetaData metaData = MetaData.builder()
+                .serviceName("echo.EchoService")
+                .methodName("echo")
+                .build();
+
+        CompletableFuture<ShenyuGrpcResponse> result = 
assertTimeoutPreemptively(Duration.ofSeconds(1),
+                () -> client.call(metaData, CallOptions.DEFAULT,
+                        "{\"data\":[{}]}", MethodDescriptor.MethodType.UNARY));
+        result.cancel(true);
+
+        assertTrue(invocation.isCancelled());
+    }
+
     @Test
     public void testCallWithNullRequestCreatesDefaultMessage() {
         assertDefaultRequest(null);

Reply via email to