This is an automated email from the ASF dual-hosted git repository.

earthchen pushed a commit to branch 3.3
in repository https://gitbox.apache.org/repos/asf/dubbo.git


The following commit(s) were added to refs/heads/3.3 by this push:
     new 82274d9864 opt: client supports requests to server-streaming endpoints 
without parameters (#15029)
82274d9864 is described below

commit 82274d98647726fc8294bb18cd3861ac368ccc66
Author: funkye <[email protected]>
AuthorDate: Mon Jan 6 10:55:17 2025 +0800

    opt: client supports requests to server-streaming endpoints without 
parameters (#15029)
    
    * optimize:  support http1 automatic keepalive setting
    
    * optimize:  support http1 automatic keepalive setting
    
    * optimize:  support http1 automatic keepalive setting
    
    * optimize:  support http1 automatic keepalive setting
    
    * feat:  server stream supports requests without parameters
    
    * feat:  server stream supports requests without parameters
    
    * fix ut
    
    * fix ut
    
    * fix ut
    
    * fix ut
    
    * update
    
    * update
    
    * update
---
 .../apache/dubbo/rpc/model/MethodDescriptor.java    |  4 ++++
 .../dubbo/rpc/model/ReflectionMethodDescriptor.java | 21 ++++++++++++++-------
 .../dubbo/rpc/model/StubMethodDescriptor.java       | 10 ++++++++++
 .../rpc/protocol/tri/ReflectionPackableMethod.java  | 17 +++++------------
 .../dubbo/rpc/protocol/tri/TripleInvoker.java       | 18 ++++++++++++++----
 5 files changed, 47 insertions(+), 23 deletions(-)

diff --git 
a/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/MethodDescriptor.java 
b/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/MethodDescriptor.java
index 2496d2dea3..77f87ede4a 100644
--- 
a/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/MethodDescriptor.java
+++ 
b/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/MethodDescriptor.java
@@ -51,6 +51,10 @@ public interface MethodDescriptor {
 
     Object getAttribute(String key);
 
+    Class<?>[] getActualRequestTypes();
+
+    Class<?> getActualResponseType();
+
     enum RpcType {
         UNARY,
         CLIENT_STREAM,
diff --git 
a/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/ReflectionMethodDescriptor.java
 
b/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/ReflectionMethodDescriptor.java
index 7d4fd1208a..ac0778d601 100644
--- 
a/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/ReflectionMethodDescriptor.java
+++ 
b/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/ReflectionMethodDescriptor.java
@@ -99,12 +99,17 @@ public class ReflectionMethodDescriptor implements 
MethodDescriptor {
         boolean returnIsVoid = 
returnClass.getName().equals(void.class.getName());
         if (returnIsVoid && parameterClasses.length == 1 && 
isStreamType(parameterClasses[0])) {
             actualRequestTypes = Collections.emptyList().toArray(new 
Class<?>[0]);
+            actualResponseType = obtainActualTypeInStreamObserver(
+                    ((ParameterizedType) 
method.getGenericParameterTypes()[0]).getActualTypeArguments()[0]);
             return RpcType.SERVER_STREAM;
         }
         if (returnIsVoid
                 && parameterClasses.length == 2
                 && !isStreamType(parameterClasses[0])
                 && isStreamType(parameterClasses[1])) {
+            actualRequestTypes = parameterClasses;
+            actualResponseType = obtainActualTypeInStreamObserver(
+                    ((ParameterizedType) 
method.getGenericParameterTypes()[1]).getActualTypeArguments()[0]);
             return RpcType.SERVER_STREAM;
         }
         if (Arrays.stream(parameterClasses).anyMatch(this::isStreamType) || 
isStreamType(returnClass)) {
@@ -119,13 +124,6 @@ public class ReflectionMethodDescriptor implements 
MethodDescriptor {
         return StreamObserver.class.isAssignableFrom(classType);
     }
 
-    private static Class<?> obtainActualTypeInStreamObserver(Type 
typeInStreamObserver) {
-        return (Class<?>)
-                (typeInStreamObserver instanceof ParameterizedType
-                        ? ((ParameterizedType) 
typeInStreamObserver).getRawType()
-                        : typeInStreamObserver);
-    }
-
     @Override
     public String getMethodName() {
         return methodName;
@@ -179,14 +177,23 @@ public class ReflectionMethodDescriptor implements 
MethodDescriptor {
         return this.attributeMap.get(key);
     }
 
+    @Override
     public Class<?>[] getActualRequestTypes() {
         return actualRequestTypes;
     }
 
+    @Override
     public Class<?> getActualResponseType() {
         return actualResponseType;
     }
 
+    private Class<?> obtainActualTypeInStreamObserver(Type 
typeInStreamObserver) {
+        return (Class<?>)
+                (typeInStreamObserver instanceof ParameterizedType
+                        ? ((ParameterizedType) 
typeInStreamObserver).getRawType()
+                        : typeInStreamObserver);
+    }
+
     @Override
     public boolean equals(Object o) {
         if (this == o) {
diff --git 
a/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/StubMethodDescriptor.java
 
b/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/StubMethodDescriptor.java
index 784c93e12b..2e0c3f84a7 100644
--- 
a/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/StubMethodDescriptor.java
+++ 
b/dubbo-common/src/main/java/org/apache/dubbo/rpc/model/StubMethodDescriptor.java
@@ -119,6 +119,16 @@ public class StubMethodDescriptor implements 
MethodDescriptor, PackableMethod {
         return this.attributeMap.get(key);
     }
 
+    @Override
+    public Class<?>[] getActualRequestTypes() {
+        return this.parameterClasses;
+    }
+
+    @Override
+    public Class<?> getActualResponseType() {
+        return this.returnClass;
+    }
+
     @Override
     public Pack getRequestPack() {
         return requestPack;
diff --git 
a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/ReflectionPackableMethod.java
 
b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/ReflectionPackableMethod.java
index 8fc8867a43..4fd388d336 100644
--- 
a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/ReflectionPackableMethod.java
+++ 
b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/ReflectionPackableMethod.java
@@ -73,19 +73,9 @@ public class ReflectionPackableMethod implements 
PackableMethod {
         switch (method.getRpcType()) {
             case CLIENT_STREAM:
             case BI_STREAM:
-                actualRequestTypes = new Class<?>[] {
-                    obtainActualTypeInStreamObserver(
-                            ((ParameterizedType) 
method.getMethod().getGenericReturnType()).getActualTypeArguments()[0])
-                };
-                actualResponseType = obtainActualTypeInStreamObserver(
-                        ((ParameterizedType) 
method.getMethod().getGenericParameterTypes()[0])
-                                .getActualTypeArguments()[0]);
-                break;
             case SERVER_STREAM:
-                actualRequestTypes = method.getMethod().getParameterTypes();
-                actualResponseType = obtainActualTypeInStreamObserver(
-                        ((ParameterizedType) 
method.getMethod().getGenericParameterTypes()[1])
-                                .getActualTypeArguments()[0]);
+                actualRequestTypes = method.getActualRequestTypes();
+                actualResponseType = method.getActualResponseType();
                 break;
             case UNARY:
                 actualRequestTypes = method.getParameterClasses();
@@ -411,6 +401,9 @@ public class ReflectionPackableMethod implements 
PackableMethod {
             for (String type : argumentsType) {
                 builder.addArgTypes(type);
             }
+            if (actualRequestTypes == null || actualRequestTypes.length == 0) {
+                return builder.build().toByteArray();
+            }
             ByteArrayOutputStream bos = new ByteArrayOutputStream();
             for (int i = 0; i < arguments.length; i++) {
                 Object argument = arguments[i];
diff --git 
a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java
 
b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java
index 6a26fbc32f..278524e4ac 100644
--- 
a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java
+++ 
b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/TripleInvoker.java
@@ -207,10 +207,20 @@ public class TripleInvoker<T> extends AbstractInvoker<T> {
 
     AsyncRpcResult invokeServerStream(MethodDescriptor methodDescriptor, 
Invocation invocation, ClientCall call) {
         RequestMetadata request = createRequest(methodDescriptor, invocation, 
null);
-        StreamObserver<Object> responseObserver =
-                (StreamObserver<Object>) invocation.getArguments()[1];
-        final StreamObserver<Object> requestObserver = streamCall(call, 
request, responseObserver);
-        requestObserver.onNext(invocation.getArguments()[0]);
+        Object[] arguments = invocation.getArguments();
+        final StreamObserver<Object> requestObserver;
+        if (arguments.length == 2) {
+            StreamObserver<Object> responseObserver = (StreamObserver<Object>) 
arguments[1];
+            requestObserver = streamCall(call, request, responseObserver);
+            requestObserver.onNext(invocation.getArguments()[0]);
+        } else if (arguments.length == 1) {
+            StreamObserver<Object> responseObserver = (StreamObserver<Object>) 
arguments[0];
+            requestObserver = streamCall(call, request, responseObserver);
+            requestObserver.onNext(null);
+        } else {
+            throw new IllegalStateException(
+                    "The first parameter must be a StreamObserver when there 
are no parameters, or the second parameter must be a StreamObserver when there 
are parameters");
+        }
         requestObserver.onCompleted();
         return new AsyncRpcResult(CompletableFuture.completedFuture(new 
AppResponse()), invocation);
     }

Reply via email to