mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] add test cases
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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;
|
||||
|
||||
+1
-1
@@ -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),
|
||||
|
||||
+1
@@ -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;
|
||||
|
||||
|
||||
-1
@@ -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;
|
||||
|
||||
|
||||
+1
@@ -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 {
|
||||
|
||||
|
||||
+7
-7
@@ -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<SubscriptionEntry> subscriptionEntryList = new ArrayList<>();
|
||||
subscriptionEntryList.add(SubscriptionEntry.newBuilder()
|
||||
-1
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user