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

Reply via email to