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 1da23ea9fa feat: server stream supports requests without parameters 
(#15026)
1da23ea9fa is described below

commit 1da23ea9faaba6a66cc10f9e6dd8b65a38b1c385
Author: funkye <[email protected]>
AuthorDate: Tue Dec 31 14:52:02 2024 +0800

    feat: server stream supports requests without parameters (#15026)
    
    * 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
---
 .../rpc/model/ReflectionMethodDescriptor.java      | 37 ++++++++++++++++++++--
 .../dubbo/springboot/demo/servlet/ApiConsumer.java | 19 +++++++++++
 .../springboot/demo/servlet/GreeterService.java    |  2 ++
 .../demo/servlet/GreeterServiceImpl.java           | 11 +++++++
 .../remoting/http12/message/MethodMetadata.java    | 24 ++------------
 .../tri/h12/ServerStreamServerCallListener.java    |  7 ++++
 6 files changed, 75 insertions(+), 25 deletions(-)

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 6bb9ddcb33..7d4fd1208a 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
@@ -22,8 +22,10 @@ import org.apache.dubbo.common.stream.StreamObserver;
 import org.apache.dubbo.common.utils.ReflectUtils;
 
 import java.lang.reflect.Method;
+import java.lang.reflect.ParameterizedType;
 import java.lang.reflect.Type;
 import java.util.Arrays;
+import java.util.Collections;
 import java.util.Objects;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentMap;
@@ -48,6 +50,8 @@ public class ReflectionMethodDescriptor implements 
MethodDescriptor {
     private final Method method;
     private final boolean generic;
     private final RpcType rpcType;
+    private Class<?>[] actualRequestTypes;
+    private Class<?> actualResponseType;
 
     public ReflectionMethodDescriptor(Method method) {
         this.method = method;
@@ -82,13 +86,25 @@ public class ReflectionMethodDescriptor implements 
MethodDescriptor {
         if (parameterClasses.length > 2) {
             return RpcType.UNARY;
         }
+        Type[] genericParameterTypes = method.getGenericParameterTypes();
         if (parameterClasses.length == 1 && isStreamType(parameterClasses[0]) 
&& isStreamType(returnClass)) {
+            this.actualRequestTypes = new Class<?>[] {
+                obtainActualTypeInStreamObserver(
+                        ((ParameterizedType) 
method.getGenericReturnType()).getActualTypeArguments()[0])
+            };
+            actualResponseType = obtainActualTypeInStreamObserver(
+                    ((ParameterizedType) 
genericParameterTypes[0]).getActualTypeArguments()[0]);
             return RpcType.BI_STREAM;
         }
-        if (parameterClasses.length == 2
+        boolean returnIsVoid = 
returnClass.getName().equals(void.class.getName());
+        if (returnIsVoid && parameterClasses.length == 1 && 
isStreamType(parameterClasses[0])) {
+            actualRequestTypes = Collections.emptyList().toArray(new 
Class<?>[0]);
+            return RpcType.SERVER_STREAM;
+        }
+        if (returnIsVoid
+                && parameterClasses.length == 2
                 && !isStreamType(parameterClasses[0])
-                && isStreamType(parameterClasses[1])
-                && returnClass.getName().equals(void.class.getName())) {
+                && isStreamType(parameterClasses[1])) {
             return RpcType.SERVER_STREAM;
         }
         if (Arrays.stream(parameterClasses).anyMatch(this::isStreamType) || 
isStreamType(returnClass)) {
@@ -103,6 +119,13 @@ 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;
@@ -156,6 +179,14 @@ public class ReflectionMethodDescriptor implements 
MethodDescriptor {
         return this.attributeMap.get(key);
     }
 
+    public Class<?>[] getActualRequestTypes() {
+        return actualRequestTypes;
+    }
+
+    public Class<?> getActualResponseType() {
+        return actualResponseType;
+    }
+
     @Override
     public boolean equals(Object o) {
         if (this == o) {
diff --git 
a/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/ApiConsumer.java
 
b/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/ApiConsumer.java
index a16fdd0722..d4519a3c1a 100644
--- 
a/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/ApiConsumer.java
+++ 
b/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/ApiConsumer.java
@@ -73,6 +73,25 @@ public class ApiConsumer {
         System.out.println("Call sayHelloServerStream");
         greeterService.sayHelloServerStream(buildRequest("triple"), 
responseObserver);
 
+        StreamObserver<HelloReply> 
sayHelloServerStreamNoParameterResponseObserver = new 
StreamObserver<HelloReply>() {
+            @Override
+            public void onNext(HelloReply reply) {
+                System.out.println("sayHelloServerStreamNoParameter onNext: " 
+ reply.getMessage());
+            }
+
+            @Override
+            public void onError(Throwable t) {
+                System.out.println("sayHelloServerStreamNoParameter onError: " 
+ t.getMessage());
+            }
+
+            @Override
+            public void onCompleted() {
+                System.out.println("sayHelloServerStreamNoParameter 
onCompleted");
+            }
+        };
+
+        
greeterService.sayHelloServerStreamNoParameter(sayHelloServerStreamNoParameterResponseObserver);
+
         StreamObserver<HelloReply> biResponseObserver = new 
StreamObserver<HelloReply>() {
             @Override
             public void onNext(HelloReply reply) {
diff --git 
a/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/GreeterService.java
 
b/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/GreeterService.java
index e777429e71..1c3345ee0b 100644
--- 
a/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/GreeterService.java
+++ 
b/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/GreeterService.java
@@ -37,6 +37,8 @@ public interface GreeterService {
      */
     void sayHelloServerStream(HelloRequest request, StreamObserver<HelloReply> 
responseObserver);
 
+    void sayHelloServerStreamNoParameter(StreamObserver<HelloReply> 
responseObserver);
+
     /**
      * Sends greetings with bi streaming
      */
diff --git 
a/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/GreeterServiceImpl.java
 
b/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/GreeterServiceImpl.java
index ef699dd69d..8626e9aa83 100644
--- 
a/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/GreeterServiceImpl.java
+++ 
b/dubbo-demo/dubbo-demo-spring-boot/dubbo-demo-spring-boot-servlet/src/main/java/org/apache/dubbo/springboot/demo/servlet/GreeterServiceImpl.java
@@ -52,6 +52,17 @@ public class GreeterServiceImpl implements GreeterService {
         responseObserver.onCompleted();
     }
 
+    @Override
+    public void sayHelloServerStreamNoParameter(StreamObserver<HelloReply> 
responseObserver) {
+        LOGGER.info("Received sayHelloServerStreamNoParameter request");
+        for (int i = 1; i < 6; i++) {
+            LOGGER.info("sayHelloServerStreamNoParameter onNext:  {} times", 
i);
+            responseObserver.onNext(toReply("Hello " + ' ' + i + " times"));
+        }
+        LOGGER.info("sayHelloServerStreamNoParameter onCompleted");
+        responseObserver.onCompleted();
+    }
+
     @Override
     public StreamObserver<HelloRequest> 
sayHelloBiStream(StreamObserver<HelloReply> responseObserver) {
         LOGGER.info("Received sayHelloBiStream request");
diff --git 
a/dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/MethodMetadata.java
 
b/dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/MethodMetadata.java
index 7f87edafbc..2c20af654c 100644
--- 
a/dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/MethodMetadata.java
+++ 
b/dubbo-remoting/dubbo-remoting-http12/src/main/java/org/apache/dubbo/remoting/http12/message/MethodMetadata.java
@@ -20,9 +20,6 @@ import org.apache.dubbo.rpc.model.MethodDescriptor;
 import org.apache.dubbo.rpc.model.ReflectionMethodDescriptor;
 import org.apache.dubbo.rpc.model.StubMethodDescriptor;
 
-import java.lang.reflect.ParameterizedType;
-import java.lang.reflect.Type;
-
 public class MethodMetadata {
 
     private final Class<?>[] actualRequestTypes;
@@ -64,19 +61,9 @@ public class MethodMetadata {
         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]);
-                return new MethodMetadata(actualRequestTypes, 
actualResponseType);
             case SERVER_STREAM:
-                actualRequestTypes = new Class[] 
{method.getMethod().getParameterTypes()[0]};
-                actualResponseType = obtainActualTypeInStreamObserver(
-                        ((ParameterizedType) 
method.getMethod().getGenericParameterTypes()[1])
-                                .getActualTypeArguments()[0]);
+                actualRequestTypes = method.getActualRequestTypes();
+                actualResponseType = method.getActualResponseType();
                 return new MethodMetadata(actualRequestTypes, 
actualResponseType);
             case UNARY:
                 actualRequestTypes = method.getParameterClasses();
@@ -85,11 +72,4 @@ public class MethodMetadata {
         }
         throw new IllegalStateException("Can not reach here");
     }
-
-    static Class<?> obtainActualTypeInStreamObserver(Type 
typeInStreamObserver) {
-        return (Class<?>)
-                (typeInStreamObserver instanceof ParameterizedType
-                        ? ((ParameterizedType) 
typeInStreamObserver).getRawType()
-                        : typeInStreamObserver);
-    }
 }
diff --git 
a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/ServerStreamServerCallListener.java
 
b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/ServerStreamServerCallListener.java
index 62fe3d7405..8c6a83f37a 100644
--- 
a/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/ServerStreamServerCallListener.java
+++ 
b/dubbo-rpc/dubbo-rpc-triple/src/main/java/org/apache/dubbo/rpc/protocol/tri/h12/ServerStreamServerCallListener.java
@@ -33,6 +33,13 @@ public class ServerStreamServerCallListener extends 
AbstractServerCallListener {
 
     @Override
     public void onMessage(Object message) {
+        Class<?>[] params = invocation.getParameterTypes();
+        if (params.length == 1) {
+            if (params[0].isInstance(responseObserver)) {
+                invocation.setArguments(new Object[] {responseObserver});
+                return;
+            }
+        }
         if (message instanceof Object[]) {
             message = ((Object[]) message)[0];
         }

Reply via email to