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];
}