From 79fcc9aea2ffdd7b3b5c63ea3f1301ceeb048653 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Fri, 25 Mar 2022 17:00:54 +0800 Subject: [PATCH] [ISSUE #3949] add test cases --- .../rocketmq/proxy/connector/ForwardProducer.java | 1 - .../proxy/grpc/GrpcMessagingProcessor.java | 5 ++--- .../grpc/adapter/channel/GrpcClientChannel.java | 2 +- .../grpc/service/cluster/ConsumerService.java | 1 + .../grpc/service/cluster/PullMessageService.java | 1 - .../grpc/service/cluster/TransactionService.java | 1 + ...viceTest.java => ForwardClientServiceTest.java} | 14 +++++++------- .../grpc/service/cluster/RouteServiceTest.java | 1 - 8 files changed, 12 insertions(+), 14 deletions(-) rename proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/{ClientServiceTest.java => ForwardClientServiceTest.java} (93%) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java index 18b33d1799..3ba34a6406 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java @@ -29,7 +29,6 @@ import org.apache.rocketmq.common.sysflag.MessageSysFlag; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; -import org.apache.rocketmq.remoting.common.RemotingHelper; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class ForwardProducer extends AbstractForwardClient { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java index 08414c5fbf..bb947e9320 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java @@ -59,10 +59,9 @@ import io.grpc.stub.StreamObserver; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; import org.apache.rocketmq.common.constant.LoggerName; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyException; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.adapter.ResponseWriter; -import org.apache.rocketmq.proxy.grpc.common.ProxyException; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.common.ResponseWriter; import org.apache.rocketmq.proxy.grpc.service.GrpcForwardService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java index b42d10728a..6a3fcb045d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java @@ -71,7 +71,7 @@ public class GrpcClientChannel extends SimpleChannel { ChannelManager channelManager, String group, String clientId, - PollCommandResponseManager manager + PollResponseManager manager ) { GrpcClientChannel channel = channelManager.createChannel( buildKey(group, clientId), diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java index 7d8c58d7ad..78c5d77736 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java @@ -46,6 +46,7 @@ import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.DelayPolicy; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyException; import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java index caa870efd6..4193609f94 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java @@ -40,7 +40,6 @@ import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.DefaultForwardClient; import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; -import org.apache.rocketmq.proxy.grpc.adapter.ProxyException; import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java index 1317490bc3..f0f730f17e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java @@ -39,6 +39,7 @@ import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; +import org.apache.rocketmq.remoting.common.RemotingHelper; public class TransactionService extends BaseService implements TransactionStateChecker { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientServiceTest.java similarity index 93% rename from proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientServiceTest.java rename to proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientServiceTest.java index a099e25923..aadad7a29e 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientServiceTest.java @@ -23,18 +23,18 @@ import org.apache.rocketmq.common.consumer.ConsumeFromWhere; import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType; import org.apache.rocketmq.common.protocol.heartbeat.MessageModel; import org.apache.rocketmq.proxy.channel.ChannelManager; +import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; -import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.remoting.protocol.LanguageCode; import org.junit.Test; import static org.junit.Assert.*; -public class ClientServiceTest extends BaseServiceTest { +public class ForwardClientServiceTest extends BaseServiceTest { private ChannelManager channelManager = new ChannelManager(); - private PollCommandResponseManager pollCommandResponseManager = new PollCommandResponseManager(); + private PollResponseManager pollResponseManager = new PollResponseManager(); @Override public void beforeEach() throws Throwable { @@ -43,11 +43,11 @@ public class ClientServiceTest extends BaseServiceTest { @Test public void testProducerHeartbeat() { - ClientService clientService = new ClientService( + ForwardClientService clientService = new ForwardClientService( this.connectorManager, Executors.newSingleThreadScheduledExecutor(), this.channelManager, - this.pollCommandResponseManager); + this.pollResponseManager); Metadata metadata = new Metadata(); metadata.put(InterceptorConstants.LANGUAGE, "JAVA"); @@ -79,11 +79,11 @@ public class ClientServiceTest extends BaseServiceTest { @Test public void testConsumerHeartbeat() { - ClientService clientService = new ClientService( + ForwardClientService clientService = new ForwardClientService( this.connectorManager, Executors.newSingleThreadScheduledExecutor(), this.channelManager, - this.pollCommandResponseManager); + this.pollResponseManager); List subscriptionEntryList = new ArrayList<>(); subscriptionEntryList.add(SubscriptionEntry.newBuilder() diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java index 65a89b6888..ba059db95f 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteServiceTest.java @@ -43,7 +43,6 @@ import org.apache.rocketmq.common.protocol.route.QueueData; import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; -import org.apache.rocketmq.proxy.grpc.common.ProxyMode; import org.junit.Test; import static org.assertj.core.api.Assertions.assertThat;