This is an automated email from the ASF dual-hosted git repository.
albumenj pushed a commit to branch 3.0
in repository https://gitbox.apache.org/repos/asf/dubbo.git
The following commit(s) were added to refs/heads/3.0 by this push:
new bfef39b [3.0] release CPU to run StreamObserver methods (#9226)
bfef39b is described below
commit bfef39bbc183fd7ceab638aa7b0c244a65b0210d
Author: zrlw <[email protected]>
AuthorDate: Sun Nov 7 15:07:39 2021 +0800
[3.0] release CPU to run StreamObserver methods (#9226)
* release CPU to run StreamObserver methods
* release CPU to run StreamObserver methods of GrpcProtocolTest
---
.../java/org/apache/dubbo/rpc/protocol/grpc/GrpcProtocolTest.java | 8 +++++++-
.../org/apache/dubbo/rpc/protocol/tri/TripleProtocolTest.java | 8 ++++++++
2 files changed, 15 insertions(+), 1 deletion(-)
diff --git
a/dubbo-rpc/dubbo-rpc-grpc/src/test/java/org/apache/dubbo/rpc/protocol/grpc/GrpcProtocolTest.java
b/dubbo-rpc/dubbo-rpc-grpc/src/test/java/org/apache/dubbo/rpc/protocol/grpc/GrpcProtocolTest.java
index 4aab8e8..3dc1e38 100644
---
a/dubbo-rpc/dubbo-rpc-grpc/src/test/java/org/apache/dubbo/rpc/protocol/grpc/GrpcProtocolTest.java
+++
b/dubbo-rpc/dubbo-rpc-grpc/src/test/java/org/apache/dubbo/rpc/protocol/grpc/GrpcProtocolTest.java
@@ -40,6 +40,8 @@ import org.junit.jupiter.api.Test;
import java.util.HashMap;
import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
public class GrpcProtocolTest {
private Protocol protocol =
ExtensionLoader.getExtensionLoader(Protocol.class).getAdaptiveExtension();
@@ -85,7 +87,7 @@ public class GrpcProtocolTest {
ListenableFuture<HelloReply> future =
serviceImpl.sayHelloAsync(HelloRequest.newBuilder().setName("World").build());
Assertions.assertEquals("Hello World", future.get().getMessage());
-
+ CountDownLatch latch = new CountDownLatch(1);
serviceImpl.sayHello(HelloRequest.newBuilder().setName("World").build(), new
StreamObserver<HelloReply>() {
@Override
@@ -100,11 +102,15 @@ public class GrpcProtocolTest {
@Override
public void onCompleted() {
+ latch.countDown();
System.out.println("onCompleted");
}
});
+ // release CPU to run StreamObserver methods.
+ latch.await(1000, TimeUnit.MILLISECONDS);
// resource recycle.
serviceRepository.destroy();
+ System.out.println("serviceRepository destroyed");
}
class MockReferenceConfig extends ReferenceConfigBase {
diff --git
a/dubbo-rpc/dubbo-rpc-triple/src/test/java/org/apache/dubbo/rpc/protocol/tri/TripleProtocolTest.java
b/dubbo-rpc/dubbo-rpc-triple/src/test/java/org/apache/dubbo/rpc/protocol/tri/TripleProtocolTest.java
index 490f7f3..821d6e9 100644
---
a/dubbo-rpc/dubbo-rpc-triple/src/test/java/org/apache/dubbo/rpc/protocol/tri/TripleProtocolTest.java
+++
b/dubbo-rpc/dubbo-rpc-triple/src/test/java/org/apache/dubbo/rpc/protocol/tri/TripleProtocolTest.java
@@ -17,6 +17,9 @@
package org.apache.dubbo.rpc.protocol.tri;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
import org.apache.dubbo.common.URL;
import org.apache.dubbo.common.extension.ExtensionLoader;
import org.apache.dubbo.common.stream.StreamObserver;
@@ -70,6 +73,7 @@ public class TripleProtocolTest {
Assertions.assertEquals("hello world", serviceImpl.echo("hello
world"));
// fixme will throw exception
// Assertions.assertEquals("hello world", serviceImpl.echoAsync("hello
world").get());
+ CountDownLatch latch = new CountDownLatch(1);
serviceImpl.serverStream("hello world", new StreamObserver<String>() {
@Override
public void onNext(String data) {
@@ -83,11 +87,15 @@ public class TripleProtocolTest {
@Override
public void onCompleted() {
+ latch.countDown();
System.out.println("onCompleted");
}
});
+ // release CPU to run StreamObserver methods.
+ latch.await(1000, TimeUnit.MILLISECONDS);
// resource recycle.
serviceRepository.destroy();
+ System.out.println("serviceRepository destroyed");
}
}