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