From 36b1e9fa2aac47af5972740f894ed92e4d0715ef Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Fri, 25 Mar 2022 12:44:02 +0800 Subject: [PATCH] [ISSUE #3949] Move common package, using adapter. --- .../apache/rocketmq/proxy/ProxyStartup.java | 2 +- .../rocketmq/proxy/config/Configuration.java | 10 ++- .../rocketmq/proxy/config/ProxyConfig.java | 2 +- .../connector/AbstractForwardClient.java | 14 ++-- .../proxy/connector/DefaultForwardClient.java | 19 ++---- .../proxy/connector/ForwardProducer.java | 19 +++--- .../proxy/connector/ForwardReadConsumer.java | 7 +- .../proxy/connector/ForwardWriteConsumer.java | 7 +- .../factory/AbstractClientFactory.java | 2 +- .../factory/AbstractMQClientFactory.java | 8 +-- .../proxy/grpc/GrpcMessagingProcessor.java | 2 +- .../grpc/{common => adapter}/DelayPolicy.java | 2 +- .../GrpcConverter.java} | 64 +++++++++---------- .../ParameterConverter.java | 2 +- .../PollResponseFuture.java} | 8 +-- .../PollResponseManager.java} | 10 +-- .../{common => adapter}/ProxyException.java | 2 +- .../grpc/{common => adapter}/ProxyMode.java | 2 +- .../ProxyResponseCode.java | 2 +- .../{common => adapter}/ResponseBuilder.java | 2 +- .../{common => adapter}/ResponseHook.java | 2 +- .../{common => adapter}/ResponseWriter.java | 2 +- .../adapter/channel/GrpcClientChannel.java | 14 ++-- .../handler/PullMessageResponseHandler.java | 9 +-- .../ReceiveMessageResponseHandler.java | 8 +-- .../handler/SendMessageResponseHandler.java | 2 +- .../grpc/service/ClusterGrpcService.java | 16 ++--- .../proxy/grpc/service/LocalGrpcService.java | 52 +++++++-------- .../grpc/service/cluster/BaseService.java | 2 +- .../grpc/service/cluster/ConsumerService.java | 26 ++++---- .../DefaultAssignmentQueueSelector.java | 4 +- ...Service.java => ForwardClientService.java} | 56 +++++++++------- .../grpc/service/cluster/ProducerService.java | 16 ++--- .../service/cluster/PullMessageService.java | 20 +++--- .../grpc/service/cluster/RouteService.java | 14 ++-- .../service/cluster/TransactionService.java | 10 +-- .../proxy/common/utils/FilterUtilTest.java | 9 ++- .../config/ConfigurationManagerTest.java | 2 +- .../connector/ForwardClientManagerTest.java | 22 ++++--- .../grpc/service/LocalGrpcServiceTest.java | 6 +- .../DefaultProducerQueueSelectorTest.java | 18 +++--- .../service/cluster/ProducerServiceTest.java | 2 +- .../service/cluster/RouteServiceTest.java | 2 +- 43 files changed, 252 insertions(+), 248 deletions(-) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => adapter}/DelayPolicy.java (98%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common/Converter.java => adapter/GrpcConverter.java} (92%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => adapter}/ParameterConverter.java (95%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common/PollCommandResponseFuture.java => adapter/PollResponseFuture.java} (84%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common/PollCommandResponseManager.java => adapter/PollResponseManager.java} (77%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => adapter}/ProxyException.java (96%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => adapter}/ProxyMode.java (97%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => adapter}/ProxyResponseCode.java (95%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => adapter}/ResponseBuilder.java (99%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => adapter}/ResponseHook.java (94%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/{common => adapter}/ResponseWriter.java (98%) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/{ClientService.java => ForwardClientService.java} (75%) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java index 6648992efc..27a21eb626 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java @@ -30,7 +30,7 @@ import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; import org.apache.rocketmq.proxy.grpc.GrpcServer; -import org.apache.rocketmq.proxy.grpc.common.ProxyMode; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; import org.apache.rocketmq.proxy.grpc.service.ClusterGrpcService; import org.apache.rocketmq.proxy.grpc.service.GrpcForwardService; import org.apache.rocketmq.proxy.grpc.service.LocalGrpcService; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/Configuration.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/Configuration.java index 84a41a3c9f..9e79532e32 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/Configuration.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/Configuration.java @@ -25,11 +25,15 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class Configuration { - private final static Logger log = LoggerFactory.getLogger(Configuration.class); + private final static Logger LOGGER = LoggerFactory.getLogger(Configuration.class); private final AtomicReference proxyConfigReference = new AtomicReference<>(); public void init() throws Exception { String proxyConfigData = loadJsonConfig(ProxyConfig.CONFIG_FILE_NAME); + if (null == proxyConfigData) { + throw new RuntimeException(String.format("load configuration from file: %s error.", ProxyConfig.CONFIG_FILE_NAME)); + } + ProxyConfig proxyConfig = JSON.parseObject(proxyConfigData, ProxyConfig.class); setProxyConfig(proxyConfig); } @@ -39,12 +43,12 @@ public class Configuration { File file = new File(filePath); if (!file.exists()) { - log.warn("the config file {} not exist", filePath); + LOGGER.warn("the config file {} not exist", filePath); return null; } long fileLength = file.length(); if (fileLength <= 0) { - log.warn("the config file {} length is zero", filePath); + LOGGER.warn("the config file {} length is zero", filePath); return null; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java index 494bf75177..9c09b8bb33 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java @@ -17,7 +17,7 @@ package org.apache.rocketmq.proxy.config; -import org.apache.rocketmq.proxy.grpc.common.ProxyMode; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; public class ProxyConfig { public final static String CONFIG_FILE_NAME = "rmq-proxy.json"; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java index 2f08e7e67c..1ea6539a92 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java @@ -23,18 +23,22 @@ import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; public abstract class AbstractForwardClient implements StartAndShutdown { - private final ForwardClientFactory forwardClientFactory; + private final ForwardClientFactory clientFactory; private MQClientAPIExt[] clients; + private final String gidPrefix; - public AbstractForwardClient(ForwardClientFactory forwardClientFactory) { - this.forwardClientFactory = forwardClientFactory; + public AbstractForwardClient(ForwardClientFactory clientFactory, String gidPrefix) { + this.clientFactory = clientFactory; + this.gidPrefix = gidPrefix; } protected abstract int getClientNum(); protected abstract MQClientAPIExt createNewClient(ForwardClientFactory forwardClientFactory, String name); - protected abstract String getNamePrefix(); + protected String getNamePrefix() { + return this.gidPrefix; + } protected MQClientAPIExt getClient() { if (clients.length == 1) { @@ -50,7 +54,7 @@ public abstract class AbstractForwardClient implements StartAndShutdown { for (int i = 0; i < clientCount; i++) { String name = getNamePrefix() + "N_" + i; - clients[i] = createNewClient(forwardClientFactory, name); + clients[i] = createNewClient(clientFactory, name); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java index 81d180e66e..43d7bbafd9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java @@ -29,8 +29,8 @@ import org.apache.rocketmq.remoting.exception.RemotingException; public class DefaultForwardClient extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_DEFAULT_"; - public DefaultForwardClient(ForwardClientFactory forwardClientFactory) { - super(forwardClientFactory); + public DefaultForwardClient(ForwardClientFactory clientFactory) { + super(clientFactory, CID_PREFIX); } @Override @@ -41,27 +41,22 @@ public class DefaultForwardClient extends AbstractForwardClient { @Override protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { double workerFactor = ConfigurationManager.getProxyConfig().getDefaultForwardClientWorkerFactor(); - final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); + int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); return clientFactory.getMQClient(name, threadCount); } - @Override - protected String getNamePrefix() { - return CID_PREFIX; - } - public CompletableFuture> getConsumerListByGroup( String brokerAddr, GetConsumerListByGroupRequestHeader requestHeader, long timeoutMillis ) { - return getClient().getConsumerListByGroup(brokerAddr, requestHeader, timeoutMillis); + return this.getClient().getConsumerListByGroup(brokerAddr, requestHeader, timeoutMillis); } public TopicRouteData getTopicRouteInfoFromNameServer(String topic, long timeoutMillis) throws RemotingException, InterruptedException, MQClientException { - return getClient().getTopicRouteInfoFromNameServer(topic, timeoutMillis); + return this.getClient().getTopicRouteInfoFromNameServer(topic, timeoutMillis); } public CompletableFuture getMaxOffset( @@ -70,7 +65,7 @@ public class DefaultForwardClient extends AbstractForwardClient { int queueId, long timeoutMillis ) { - return getClient().getMaxOffset(brokerAddr, topic, queueId, timeoutMillis); + return this.getClient().getMaxOffset(brokerAddr, topic, queueId, timeoutMillis); } public CompletableFuture searchOffset( @@ -80,6 +75,6 @@ public class DefaultForwardClient extends AbstractForwardClient { long timestamp, long timeoutMillis ) { - return getClient().searchOffset(brokerAddr, topic, queueId, timestamp, timeoutMillis); + return this.getClient().searchOffset(brokerAddr, topic, queueId, timestamp, timeoutMillis); } } 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 2f6d504c72..835c8dbbd2 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 @@ -37,7 +37,7 @@ public class ForwardProducer extends AbstractForwardClient { private static final String PID_PREFIX = "PID_RMQ_PROXY_PUBLISH_MESSAGE_"; public ForwardProducer(ForwardClientFactory clientFactory) { - super(clientFactory); + super(clientFactory, PID_PREFIX); } @Override @@ -47,16 +47,12 @@ public class ForwardProducer extends AbstractForwardClient { @Override protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { - double sendClientWorkerFactor = ConfigurationManager.getProxyConfig().getForwardProducerWorkerFactor(); - final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * sendClientWorkerFactor); + double workerFactor = ConfigurationManager.getProxyConfig().getForwardProducerWorkerFactor(); + final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); return clientFactory.getTransactionalProducer(name, threadCount); } - @Override - protected String getNamePrefix() { - return PID_PREFIX; - } public CompletableFuture heartBeat(String heartbeatAddr, HeartbeatData heartbeatData, long timeout) throws Exception { return this.getClient().sendHeartbeat(heartbeatAddr, heartbeatData, timeout); @@ -83,8 +79,13 @@ public class ForwardProducer extends AbstractForwardClient { ); } - public CompletableFuture sendMessage(String address, String brokerName, Message msg, - SendMessageRequestHeader requestHeader, long timeoutMillis) { + public CompletableFuture sendMessage( + String address, + String brokerName, + Message msg, + SendMessageRequestHeader requestHeader, + long timeoutMillis + ) { CompletableFuture future = this.getClient().sendMessage(address, brokerName, msg, requestHeader, timeoutMillis); return future.thenApply(sendResult -> { int tranType = MessageSysFlag.getTransactionValue(requestHeader.getSysFlag()); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java index a8faa303b3..869aeffc14 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java @@ -30,7 +30,7 @@ public class ForwardReadConsumer extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_CONSUME_MESSAGE_"; public ForwardReadConsumer(ForwardClientFactory clientFactory) { - super(clientFactory); + super(clientFactory, CID_PREFIX); } @Override @@ -46,11 +46,6 @@ public class ForwardReadConsumer extends AbstractForwardClient { return clientFactory.getMQClient(name, threadCount); } - @Override - protected String getNamePrefix() { - return CID_PREFIX; - } - public CompletableFuture popMessage(String address, String brokerName, PopMessageRequestHeader requestHeader, long timeoutMillis) { return getClient().popMessage(address, brokerName, requestHeader, timeoutMillis); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java index fc67796796..1c0d721cc0 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java @@ -31,7 +31,7 @@ public class ForwardWriteConsumer extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_DELETE_MESSAGE_"; public ForwardWriteConsumer(ForwardClientFactory clientFactory) { - super(clientFactory); + super(clientFactory, CID_PREFIX); } @Override @@ -47,11 +47,6 @@ public class ForwardWriteConsumer extends AbstractForwardClient { return clientFactory.getMQClient(name, threadCount); } - @Override - protected String getNamePrefix() { - return CID_PREFIX; - } - public CompletableFuture ackMessage( String address, AckMessageRequestHeader requestHeader, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java index daae5bc2bd..f5d3aa379a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java @@ -78,7 +78,7 @@ public abstract class AbstractClientFactory { try { this.shutdown(v); } catch (Exception e) { - LOGGER.warn("RocketMQClientConstructor shutdown all err.", e); + LOGGER.warn("try to shutdown client err.", e); } }); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java index 6f126e9200..1a772dc04d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java @@ -16,6 +16,7 @@ */ package org.apache.rocketmq.proxy.connector.factory; +import java.time.Duration; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.client.ClientConfig; @@ -25,8 +26,7 @@ import org.apache.rocketmq.remoting.RPCHook; public abstract class AbstractMQClientFactory extends AbstractClientFactory { - public AbstractMQClientFactory(ScheduledExecutorService scheduledExecutorService, - RPCHook rpcHook) { + public AbstractMQClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) { super(scheduledExecutorService, rpcHook); } @@ -50,8 +50,8 @@ public abstract class AbstractMQClientFactory extends AbstractClientFactory property = buildMessageProperty(message); - requestHeader.setProducerGroup(getResourceNameWithNamespace(systemAttribute.getProducerGroup())); - requestHeader.setTopic(getResourceNameWithNamespace(message.getTopic())); + requestHeader.setProducerGroup(wrapResourceWithNamespace(systemAttribute.getProducerGroup())); + requestHeader.setTopic(wrapResourceWithNamespace(message.getTopic())); requestHeader.setDefaultTopic(""); requestHeader.setDefaultTopicQueueNums(0); requestHeader.setQueueId(systemAttribute.getPartitionId()); @@ -135,10 +135,10 @@ public class Converter { public static PopMessageRequestHeader buildPopMessageRequestHeader(ReceiveMessageRequest request, long pollTime) { Resource group = request.getGroup(); - String groupName = Converter.getResourceNameWithNamespace(group); + String groupName = GrpcConverter.wrapResourceWithNamespace(group); Partition partition = request.getPartition(); Resource topic = partition.getTopic(); - String topicName = Converter.getResourceNameWithNamespace(topic); + String topicName = GrpcConverter.wrapResourceWithNamespace(topic); int queueId = partition.getId(); int maxMessageNumbers = request.getBatchSize(); if (maxMessageNumbers > ProxyUtils.MAX_MSG_NUMS_FOR_POP_REQUEST) { @@ -149,11 +149,11 @@ public class Converter { long invisibleTime = Durations.toMillis(request.getInvisibleDuration()); long bornTime = Timestamps.toMillis(request.getInitializationTimestamp()); ConsumePolicy policy = request.getConsumePolicy(); - int initMode = Converter.buildConsumeInitMode(policy); + int initMode = GrpcConverter.buildConsumeInitMode(policy); FilterExpression filterExpression = request.getFilterExpression(); String expression = filterExpression.getExpression(); - String expressionType = Converter.buildExpressionType(filterExpression.getType()); + String expressionType = GrpcConverter.buildExpressionType(filterExpression.getType()); PopMessageRequestHeader requestHeader = new PopMessageRequestHeader(); requestHeader.setConsumerGroup(groupName); @@ -172,8 +172,8 @@ public class Converter { } public static AckMessageRequestHeader buildAckMessageRequestHeader(AckMessageRequest request) { - String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); - String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); + String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); String receiptHandleStr = request.getReceiptHandle(); ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); @@ -188,8 +188,8 @@ public class Converter { public static ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(NackMessageRequest request, DelayPolicy delayPolicy) { - String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); - String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); + String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); String receiptHandleStr = request.getReceiptHandle(); ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); @@ -206,8 +206,8 @@ public class Converter { public static ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader( ChangeInvisibleDurationRequest request) { - String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); - String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); + String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); String receiptHandleStr = request.getReceiptHandle(); ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); @@ -223,8 +223,8 @@ public class Converter { public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackRequestHeader( ForwardMessageToDeadLetterQueueRequest request) { - String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); - String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); + String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); String receiptHandleStr = request.getReceiptHandle(); ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); @@ -240,8 +240,8 @@ public class Converter { public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader( NackMessageRequest request) { - String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); - String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); + String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); String receiptHandleStr = request.getReceiptHandle(); ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); @@ -256,7 +256,7 @@ public class Converter { } public static EndTransactionRequestHeader buildEndTransactionRequestHeader(EndTransactionRequest request) { - String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); + String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); String messageId = request.getMessageId(); String transactionId = request.getTransactionId(); TransactionId handle; @@ -268,7 +268,7 @@ public class Converter { long transactionStateTableOffset = handle.getTranStateTableOffset(); long commitLogOffset = handle.getCommitLogOffset(); boolean fromTransactionCheck = request.getSource() == EndTransactionRequest.Source.SERVER_CHECK; - int commitOrRollback = Converter.buildTransactionCommitOrRollback(request.getResolution()); + int commitOrRollback = GrpcConverter.buildTransactionCommitOrRollback(request.getResolution()); EndTransactionRequestHeader endTransactionRequestHeader = new EndTransactionRequestHeader(); endTransactionRequestHeader.setProducerGroup(groupName); @@ -284,13 +284,13 @@ public class Converter { public static PullMessageRequestHeader buildPullMessageRequestHeader(PullMessageRequest request, long pollTimeoutInMillis) { Partition partition = request.getPartition(); - String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); - String topicName = Converter.getResourceNameWithNamespace(partition.getTopic()); + String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); + String topicName = GrpcConverter.wrapResourceWithNamespace(partition.getTopic()); int queueId = partition.getId(); int sysFlag = PullSysFlag.buildSysFlag(false, true, true, false, false); String expression = request.getFilterExpression().getExpression(); - String expressionType = Converter.buildExpressionType(request.getFilterExpression().getType()); + String expressionType = GrpcConverter.buildExpressionType(request.getFilterExpression().getType()); PullMessageRequestHeader requestHeader = new PullMessageRequestHeader(); requestHeader.setConsumerGroup(groupName); @@ -370,7 +370,7 @@ public class Converter { MessageAccessor.setReconsumeTime(messageWithHeader, String.valueOf(reconsumeTimes)); // set producer group Resource producerGroup = message.getSystemAttribute().getProducerGroup(); - String producerGroupName = getResourceNameWithNamespace(producerGroup); + String producerGroupName = wrapResourceWithNamespace(producerGroup); MessageAccessor.putProperty(messageWithHeader, MessageConst.PROPERTY_PRODUCER_GROUP, producerGroupName); // set message group String messageGroup = message.getSystemAttribute().getMessageGroup(); @@ -386,7 +386,7 @@ public class Converter { } public static org.apache.rocketmq.common.message.Message buildMessage(Message protoMessage) { - String topic = getResourceNameWithNamespace(protoMessage.getTopic()); + String topic = wrapResourceWithNamespace(protoMessage.getTopic()); org.apache.rocketmq.common.message.Message message = new org.apache.rocketmq.common.message.Message(topic, protoMessage.getBody().toByteArray()); @@ -420,13 +420,13 @@ public class Converter { public static org.apache.rocketmq.common.protocol.heartbeat.ProducerData buildProducerData(ProducerData producerData) { org.apache.rocketmq.common.protocol.heartbeat.ProducerData buildProducerData = new org.apache.rocketmq.common.protocol.heartbeat.ProducerData(); - buildProducerData.setGroupName(getResourceNameWithNamespace(producerData.getGroup())); + buildProducerData.setGroupName(wrapResourceWithNamespace(producerData.getGroup())); return buildProducerData; } public static org.apache.rocketmq.common.protocol.heartbeat.ConsumerData buildConsumerData(ConsumerData consumerData) { org.apache.rocketmq.common.protocol.heartbeat.ConsumerData buildConsumerData = new org.apache.rocketmq.common.protocol.heartbeat.ConsumerData(); - buildConsumerData.setGroupName(getResourceNameWithNamespace(consumerData.getGroup())); + buildConsumerData.setGroupName(wrapResourceWithNamespace(consumerData.getGroup())); buildConsumerData.setConsumeType(buildConsumeType(consumerData.getConsumeType())); buildConsumerData.setMessageModel(buildMessageModel(consumerData.getConsumeModel())); buildConsumerData.setConsumeFromWhere(buildConsumeFromWhere(consumerData.getConsumePolicy())); @@ -472,7 +472,7 @@ public class Converter { public static Set buildSubscriptionDataSet(List subscriptionEntryList) { Set subscriptionDataSet = new HashSet<>(); for (SubscriptionEntry sub : subscriptionEntryList) { - String topicName = Converter.getResourceNameWithNamespace(sub.getTopic()); + String topicName = GrpcConverter.wrapResourceWithNamespace(sub.getTopic()); FilterExpression filterExpression = sub.getExpression(); subscriptionDataSet.add(buildSubscriptionData(topicName, filterExpression)); } @@ -481,7 +481,7 @@ public class Converter { public static SubscriptionData buildSubscriptionData(String topicName, FilterExpression filterExpression) { String expression = filterExpression.getExpression(); - String expressionType = Converter.buildExpressionType(filterExpression.getType()); + String expressionType = GrpcConverter.buildExpressionType(filterExpression.getType()); try { return FilterAPI.build(topicName, expression, expressionType); } catch (Exception e) { @@ -693,10 +693,10 @@ public class Converter { UnregisterClientRequestHeader header = new UnregisterClientRequestHeader(); header.setClientID(request.getClientId()); if (request.hasProducerGroup()) { - header.setProducerGroup(getResourceNameWithNamespace(request.getProducerGroup())); + header.setProducerGroup(wrapResourceWithNamespace(request.getProducerGroup())); } if (request.hasConsumerGroup()) { - header.setConsumerGroup(getResourceNameWithNamespace(request.getConsumerGroup())); + header.setConsumerGroup(wrapResourceWithNamespace(request.getConsumerGroup())); } return header; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ParameterConverter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ParameterConverter.java similarity index 95% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ParameterConverter.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ParameterConverter.java index 7561a21bc4..47641cf700 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ParameterConverter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ParameterConverter.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; import io.grpc.Context; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/PollCommandResponseFuture.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/PollResponseFuture.java similarity index 84% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/PollCommandResponseFuture.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/PollResponseFuture.java index 122077b542..0bd79de2b6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/PollCommandResponseFuture.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/PollResponseFuture.java @@ -15,18 +15,18 @@ * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; -public class PollCommandResponseFuture { +public class PollResponseFuture { private final String commandId; private final Integer opaque; - public PollCommandResponseFuture(String commandId, int opaque) { + public PollResponseFuture(String commandId, int opaque) { this.commandId = commandId; this.opaque = opaque; } - public PollCommandResponseFuture(String commandId) { + public PollResponseFuture(String commandId) { this.commandId = commandId; this.opaque = null; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/PollCommandResponseManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/PollResponseManager.java similarity index 77% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/PollCommandResponseManager.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/PollResponseManager.java index 0311788756..8379fdf4c8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/PollCommandResponseManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/PollResponseManager.java @@ -15,23 +15,23 @@ * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.atomic.AtomicLong; -public class PollCommandResponseManager { - private final ConcurrentMap futureTable = new ConcurrentHashMap<>(); +public class PollResponseManager { + private final ConcurrentMap futureTable = new ConcurrentHashMap<>(); private final AtomicLong commandIdGenerator = new AtomicLong(0); public String putResponse(int opaque) { String commandId = String.valueOf(commandIdGenerator.incrementAndGet()); - futureTable.put(commandId, new PollCommandResponseFuture(commandId, opaque)); + futureTable.put(commandId, new PollResponseFuture(commandId, opaque)); return commandId; } - public PollCommandResponseFuture getResponse(String commandId) { + public PollResponseFuture getResponse(String commandId) { return futureTable.get(commandId); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyException.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyException.java similarity index 96% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyException.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyException.java index f476044383..da19dafc4e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyException.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyException.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; import com.google.rpc.Code; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyMode.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyMode.java similarity index 97% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyMode.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyMode.java index 25ac8665f6..73856c64b5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyMode.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyMode.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; public enum ProxyMode { LOCAL("LOCAL"), diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyResponseCode.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyResponseCode.java similarity index 95% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyResponseCode.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyResponseCode.java index 2134cb6eb3..0c1fa6ee00 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ProxyResponseCode.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ProxyResponseCode.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; public enum ProxyResponseCode { SYS_ERR, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseBuilder.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseBuilder.java similarity index 99% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseBuilder.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseBuilder.java index 7902014e02..d475d11fa1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseBuilder.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseBuilder.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; import apache.rocketmq.v1.HeartbeatResponse; import apache.rocketmq.v1.ResponseCommon; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseHook.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseHook.java similarity index 94% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseHook.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseHook.java index 6af3052142..2a0a2bac08 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseHook.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseHook.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; public interface ResponseHook { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseWriter.java similarity index 98% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseWriter.java index 5d8fff2998..01965c8399 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/ResponseWriter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/ResponseWriter.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.rocketmq.proxy.grpc.common; +package org.apache.rocketmq.proxy.grpc.adapter; import io.grpc.stub.ServerCallStreamObserver; import io.grpc.stub.StreamObserver; 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 65bd2ad5eb..5e7f79f6a9 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 @@ -30,8 +30,8 @@ import org.apache.rocketmq.common.protocol.header.CheckTransactionStateRequestHe import org.apache.rocketmq.common.protocol.header.GetConsumerRunningInfoRequestHeader; import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.channel.SimpleChannel; -import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class GrpcClientChannel extends SimpleChannel { @@ -39,9 +39,9 @@ public class GrpcClientChannel extends SimpleChannel { private final String group; private final String clientId; - private final PollCommandResponseManager manager; + private final PollResponseManager manager; - private GrpcClientChannel(String group, String clientId, PollCommandResponseManager manager) { + private GrpcClientChannel(String group, String clientId, PollResponseManager manager) { super(ChannelManager.createSimpleChannelDirectly()); this.group = group; this.clientId = clientId; @@ -56,7 +56,7 @@ public class GrpcClientChannel extends SimpleChannel { ChannelManager channelManager, String group, String clientId, - PollCommandResponseManager manager + PollResponseManager manager ) { GrpcClientChannel channel = channelManager.createChannel( buildKey(group, clientId), @@ -103,7 +103,7 @@ public class GrpcClientChannel extends SimpleChannel { future.complete(PollCommandResponse.newBuilder() .setRecoverOrphanedTransactionCommand(RecoverOrphanedTransactionCommand.newBuilder() .setTransactionId(requestHeader.getTransactionId()) - .setOrphanedTransactionalMessage(Converter.buildMessage(messageExt)) + .setOrphanedTransactionalMessage(GrpcConverter.buildMessage(messageExt)) .build()) .build()); break; @@ -123,7 +123,7 @@ public class GrpcClientChannel extends SimpleChannel { break; } } - } catch (Exception e) { + } catch (Exception ignore) { } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/PullMessageResponseHandler.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/PullMessageResponseHandler.java index ce85e2e4fe..7344d5f451 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/PullMessageResponseHandler.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/PullMessageResponseHandler.java @@ -26,12 +26,13 @@ import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.ResponseCode; import org.apache.rocketmq.common.protocol.header.PullMessageResponseHeader; import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext; -import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class PullMessageResponseHandler implements ResponseHandler { - @Override public void handle(RemotingCommand responseCommand, + @Override + public void handle(RemotingCommand responseCommand, InvocationContext context) { try { PullMessageResponseHeader responseHeader = (PullMessageResponseHeader) responseCommand.readCustomHeader(); @@ -40,7 +41,7 @@ public class PullMessageResponseHandler implements ResponseHandler msgFoundList = MessageDecoder.decodes(byteBuffer); for (MessageExt messageExt : msgFoundList) { - builder.addMessages(Converter.buildMessage(messageExt)); + builder.addMessages(GrpcConverter.buildMessage(messageExt)); } } PullMessageResponse response = builder.setCommon(ResponseBuilder.buildCommon(responseCommand.getCode(), responseCommand.getRemark())) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/ReceiveMessageResponseHandler.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/ReceiveMessageResponseHandler.java index f85cf4cc41..82b5127067 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/ReceiveMessageResponseHandler.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/handler/ReceiveMessageResponseHandler.java @@ -37,8 +37,8 @@ import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.header.ExtraInfoUtil; import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader; import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext; -import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.apache.rocketmq.remoting.protocol.RemotingSysResponseCode; import org.slf4j.Logger; @@ -122,7 +122,7 @@ public class ReceiveMessageResponseHandler implements ResponseHandler { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index 63bee24262..d74f83220b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -64,10 +64,10 @@ import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.transaction.TransactionStateCheckRequest; import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker; -import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; -import org.apache.rocketmq.proxy.grpc.common.ProxyMode; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.service.cluster.ClientService; +import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.service.cluster.ForwardClientService; import org.apache.rocketmq.proxy.grpc.service.cluster.ConsumerService; import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService; import org.apache.rocketmq.proxy.grpc.service.cluster.PullMessageService; @@ -87,19 +87,19 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc private final ProducerService producerService; private final ConsumerService receiveMessageService; private final RouteService routeService; - private final ClientService clientService; + private final ForwardClientService clientService; private final PullMessageService pullMessageService; private final TransactionService transactionService; - private final PollCommandResponseManager pollCommandResponseManager; + private final PollResponseManager pollCommandResponseManager; public ClusterGrpcService() { this.channelManager = new ChannelManager(); - this.pollCommandResponseManager = new PollCommandResponseManager(); + this.pollCommandResponseManager = new PollResponseManager(); this.connectorManager = new ConnectorManager(new GrpcTransactionStateChecker()); this.receiveMessageService = new ConsumerService(connectorManager); this.producerService = new ProducerService(connectorManager); this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager); - this.clientService = new ClientService(connectorManager, scheduledExecutorService, channelManager, pollCommandResponseManager); + this.clientService = new ForwardClientService(connectorManager, scheduledExecutorService, channelManager, pollCommandResponseManager); this.pullMessageService = new PullMessageService(connectorManager); this.transactionService = new TransactionService(connectorManager, channelManager); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index d4c2cce67b..3565be03c8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -98,12 +98,12 @@ import org.apache.rocketmq.proxy.grpc.adapter.channel.SendMessageChannel; import org.apache.rocketmq.proxy.grpc.adapter.handler.PullMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.adapter.handler.ReceiveMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler; -import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.DelayPolicy; -import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseFuture; -import org.apache.rocketmq.proxy.grpc.common.PollCommandResponseManager; -import org.apache.rocketmq.proxy.grpc.common.ProxyMode; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.adapter.DelayPolicy; +import org.apache.rocketmq.proxy.grpc.adapter.PollResponseFuture; +import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService; import org.apache.rocketmq.remoting.RemotingServer; @@ -120,7 +120,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor( new ThreadFactoryImpl("LocalGrpcServiceScheduledThread")); private final ChannelManager channelManager; - private final PollCommandResponseManager pollCommandResponseManager; + private final PollResponseManager pollCommandResponseManager; private final RouteService routeService; private final DelayPolicy delayPolicy; @@ -129,7 +129,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo this.channelManager = new ChannelManager(); // TransactionStateChecker is not used in Local mode. ConnectorManager connectorManager = new ConnectorManager(null); - this.pollCommandResponseManager = new PollCommandResponseManager(); + this.pollCommandResponseManager = new PollResponseManager(); this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager); this.delayPolicy = DelayPolicy.build(brokerController.getMessageStoreConfig().getMessageDelayLevel()); this.appendStartAndShutdown(connectorManager); @@ -146,17 +146,17 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo LanguageCode languageCode; String language = InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.LANGUAGE); languageCode = LanguageCode.valueOf(language); - HeartbeatData heartbeatData = Converter.buildHeartbeatData(request); + HeartbeatData heartbeatData = GrpcConverter.buildHeartbeatData(request); CompletableFuture future = new CompletableFuture<>(); String groupName; switch (request.getClientDataCase()) { case PRODUCER_DATA: { - groupName = Converter.getResourceNameWithNamespace(request.getProducerData().getGroup()); + groupName = GrpcConverter.wrapResourceWithNamespace(request.getProducerData().getGroup()); break; } case CONSUMER_DATA: { - groupName = Converter.getResourceNameWithNamespace(request.getConsumerData().getGroup()); + groupName = GrpcConverter.wrapResourceWithNamespace(request.getConsumerData().getGroup()); break; } default: { @@ -191,7 +191,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo @Override public CompletableFuture sendMessage(Context ctx, SendMessageRequest request) { - SendMessageRequestHeader requestHeader = Converter.buildSendMessageRequestHeader(request); + SendMessageRequestHeader requestHeader = GrpcConverter.buildSendMessageRequestHeader(request); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.SEND_MESSAGE, requestHeader); Message message = request.getMessage(); command.setBody(message.getBody().toByteArray()); @@ -233,7 +233,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo if (pollTime <= 0) { pollTime = timeRemaining; } - PopMessageRequestHeader requestHeader = Converter.buildPopMessageRequestHeader(request, pollTime); + PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader); command.makeCustomHeaderToNet(); @@ -262,7 +262,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); - AckMessageRequestHeader requestHeader = Converter.buildAckMessageRequestHeader(request); + AckMessageRequestHeader requestHeader = GrpcConverter.buildAckMessageRequestHeader(request); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.ACK_MESSAGE, requestHeader); command.makeCustomHeaderToNet(); @@ -289,7 +289,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); - ChangeInvisibleTimeRequestHeader requestHeader = Converter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); + ChangeInvisibleTimeRequestHeader requestHeader = GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, requestHeader); command.makeCustomHeaderToNet(); @@ -314,7 +314,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo SimpleChannel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); - ConsumerSendMsgBackRequestHeader requestHeader = Converter.buildConsumerSendMsgBackRequestHeader(request); + ConsumerSendMsgBackRequestHeader requestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader(request); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader); command.makeCustomHeaderToNet(); @@ -344,7 +344,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); - EndTransactionRequestHeader requestHeader = Converter.buildEndTransactionRequestHeader(request); + EndTransactionRequestHeader requestHeader = GrpcConverter.buildEndTransactionRequestHeader(request); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.END_TRANSACTION, requestHeader); command.makeCustomHeaderToNet(); @@ -370,7 +370,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo @Override public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) { Partition partition = request.getPartition(); - String topicName = Converter.getResourceNameWithNamespace(partition.getTopic()); + String topicName = GrpcConverter.wrapResourceWithNamespace(partition.getTopic()); int queueId = partition.getId(); long offset; @@ -399,7 +399,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo if (pollTime <= 0) { pollTime = timeRemaining; } - PullMessageRequestHeader requestHeader = Converter.buildPullMessageRequestHeader(request, pollTime); + PullMessageRequestHeader requestHeader = GrpcConverter.buildPullMessageRequestHeader(request, pollTime); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.PULL_MESSAGE, requestHeader); command.makeCustomHeaderToNet(); @@ -431,7 +431,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo switch (request.getGroupCase()) { case PRODUCER_GROUP: Resource producerGroup = request.getProducerGroup(); - String producerGroupName = Converter.getResourceNameWithNamespace(producerGroup); + String producerGroupName = GrpcConverter.wrapResourceWithNamespace(producerGroup); GrpcClientChannel producerChannel = GrpcClientChannel.getChannel(channelManager, producerGroupName, clientId); if (producerChannel == null) { future.complete(PollCommandResponse.newBuilder() @@ -443,7 +443,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo break; case CONSUMER_GROUP: Resource consumerGroup = request.getConsumerGroup(); - String consumerGroupName = Converter.getResourceNameWithNamespace(consumerGroup); + String consumerGroupName = GrpcConverter.wrapResourceWithNamespace(consumerGroup); GrpcClientChannel consumerChannel = GrpcClientChannel.getChannel(channelManager, consumerGroupName, clientId); if (consumerChannel == null) { future.complete(PollCommandResponse.newBuilder() @@ -464,7 +464,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo ReportThreadStackTraceRequest request) { String commandId = request.getCommandId(); String threadStack = request.getThreadStackTrace(); - PollCommandResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(commandId); + PollResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(commandId); if (pollCommandResponseFuture != null) { RemotingServer remotingServer = this.brokerController.getRemotingServer(); if (remotingServer instanceof NettyRemotingAbstract) { @@ -487,14 +487,14 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo ReportMessageConsumptionResultRequest request) { String commandId = request.getCommandId(); - PollCommandResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(commandId); + PollResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(commandId); if (pollCommandResponseFuture != null) { RemotingServer remotingServer = this.brokerController.getRemotingServer(); if (remotingServer instanceof NettyRemotingAbstract) { NettyRemotingAbstract nettyRemotingAbstract = (NettyRemotingAbstract) remotingServer; RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "From gRPC client"); remotingCommand.setOpaque(pollCommandResponseFuture.getOpaque()); - ConsumeMessageDirectlyResult result = Converter.buildConsumeMessageDirectlyResult(request); + ConsumeMessageDirectlyResult result = GrpcConverter.buildConsumeMessageDirectlyResult(request); remotingCommand.setBody(result.encode()); nettyRemotingAbstract.processResponseCommand(new SimpleChannelHandlerContext(channelManager.createChannel()), remotingCommand); } @@ -509,7 +509,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo NotifyClientTerminationRequest request) { Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel); - UnregisterClientRequestHeader header = Converter.buildUnregisterClientRequestHeader(request); + UnregisterClientRequestHeader header = GrpcConverter.buildUnregisterClientRequestHeader(request); RemotingCommand remotingCommand = RemotingCommand.createRequestCommand(RequestCode.UNREGISTER_CLIENT, header); remotingCommand.makeCustomHeaderToNet(); @@ -526,7 +526,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo Channel channel = channelManager.createChannel(); SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); - ChangeInvisibleTimeRequestHeader requestHeader = Converter.buildChangeInvisibleTimeRequestHeader(request); + ChangeInvisibleTimeRequestHeader requestHeader = GrpcConverter.buildChangeInvisibleTimeRequestHeader(request); ReceiptHandle receiptHandle = ReceiptHandle.decode(request.getReceiptHandle()); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, requestHeader); command.makeCustomHeaderToNet(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java index fd2854da15..2c7bb94d3d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java @@ -21,7 +21,7 @@ import io.grpc.Context; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.consumer.ReceiptHandle; import org.apache.rocketmq.proxy.connector.ConnectorManager; -import org.apache.rocketmq.proxy.grpc.common.ProxyException; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyException; public class BaseService { 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 d6975e6892..56c207a295 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 @@ -44,10 +44,10 @@ import org.apache.rocketmq.proxy.connector.ForwardProducer; import org.apache.rocketmq.proxy.connector.ForwardReadConsumer; import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; -import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.DelayPolicy; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.common.ResponseHook; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.adapter.DelayPolicy; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; import java.util.ArrayList; import java.util.List; @@ -115,7 +115,7 @@ public class ConsumerService extends BaseService { protected PopMessageRequestHeader convertToPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) { // check filterExpression is correct or not - Converter.buildSubscriptionData(Converter.getResourceNameWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); + GrpcConverter.buildSubscriptionData(GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); long timeRemaining = ctx.getDeadline() .timeRemaining(TimeUnit.MILLISECONDS); @@ -124,12 +124,12 @@ public class ConsumerService extends BaseService { pollTime = timeRemaining; } - return Converter.buildPopMessageRequestHeader(request, pollTime); + return GrpcConverter.buildPopMessageRequestHeader(request, pollTime); } protected ReceiveMessageResponse convertToReceiveMessageResponse(Context ctx, ReceiveMessageRequest request, PopResult result) { - SubscriptionData subscriptionData = Converter.buildSubscriptionData( - Converter.getResourceNameWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); + SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData( + GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); PopStatus status = result.getPopStatus(); switch (status) { case FOUND: @@ -152,7 +152,7 @@ public class ConsumerService extends BaseService { this.ackNoMatchedMessage(ctx, request, messageExt); continue; } - messages.add(Converter.buildMessage(messageExt)); + messages.add(GrpcConverter.buildMessage(messageExt)); } return ReceiveMessageResponse.newBuilder() @@ -170,7 +170,7 @@ public class ConsumerService extends BaseService { return; } String brokerAddr = this.getBrokerAddr(ctx, handle.getBrokerName()); - ackMessageRequestHeader.setConsumerGroup(Converter.getResourceNameWithNamespace(request.getGroup())); + ackMessageRequestHeader.setConsumerGroup(GrpcConverter.wrapResourceWithNamespace(request.getGroup())); ackMessageRequestHeader.setTopic(messageExt.getTopic()); ackMessageRequestHeader.setQueueId(handle.getQueueId()); ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle()); @@ -219,7 +219,7 @@ public class ConsumerService extends BaseService { } protected AckMessageRequestHeader convertToAckMessageRequestHeader(Context ctx, AckMessageRequest request) { - return Converter.buildAckMessageRequestHeader(request); + return GrpcConverter.buildAckMessageRequestHeader(request); } protected AckMessageResponse convertToAckMessageResponse(Context ctx, AckMessageRequest request, AckResult ackResult) { @@ -286,11 +286,11 @@ public class ConsumerService extends BaseService { } protected ChangeInvisibleTimeRequestHeader convertToChangeInvisibleTimeRequestHeader(Context ctx, NackMessageRequest request) { - return Converter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); + return GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); } protected ConsumerSendMsgBackRequestHeader convertToConsumerSendMsgBackToDLQRequestHeader(Context ctx, NackMessageRequest request) { - return Converter.buildConsumerSendMsgBackToDLQRequestHeader(request); + return GrpcConverter.buildConsumerSendMsgBackToDLQRequestHeader(request); } protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request, AckResult ackResult) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java index 92ea26eb5a..d81f9c92b5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultAssignmentQueueSelector.java @@ -22,7 +22,7 @@ import java.util.List; import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; -import org.apache.rocketmq.proxy.grpc.common.Converter; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; public class DefaultAssignmentQueueSelector implements AssignmentQueueSelector { @@ -34,7 +34,7 @@ public class DefaultAssignmentQueueSelector implements AssignmentQueueSelector { @Override public List getAssignment(Context ctx, QueryAssignmentRequest request) throws Exception { - String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); MessageQueueWrapper messageQueueWrapper = topicRouteCache.getMessageQueue(topicName); return messageQueueWrapper.getReadSelector().getBrokerActingQueues(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java similarity index 75% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java index 69aa281186..6c2c66268e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java @@ -24,6 +24,7 @@ import apache.rocketmq.v1.PollCommandRequest; import apache.rocketmq.v1.PollCommandResponse; import apache.rocketmq.v1.Resource; import io.grpc.Context; +import java.time.Duration; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @@ -33,35 +34,40 @@ import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.common.MQVersion; import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; -import org.apache.rocketmq.proxy.grpc.common.Converter; -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.slf4j.Logger; import org.slf4j.LoggerFactory; -public class ClientService extends BaseService { - - private static final Logger log = LoggerFactory.getLogger(ClientService.class); +public class ForwardClientService extends BaseService { + private static final Logger LOGGER = LoggerFactory.getLogger(ForwardClientService.class); private final ChannelManager channelManager; - private final ConsumerManager consumerManager = new ConsumerManager((event, group, args) -> { - }); + private final ConsumerManager consumerManager; private final ProducerManager producerManager; - private final PollCommandResponseManager pollCommandResponseManager; + private final PollResponseManager pollCommandResponseManager; - public ClientService( + public ForwardClientService( ConnectorManager connectorManager, ScheduledExecutorService scheduledExecutorService, ChannelManager channelManager, - PollCommandResponseManager pollCommandResponseManager + PollResponseManager pollCommandResponseManager ) { super(connectorManager); - scheduledExecutorService.scheduleWithFixedDelay(this::scanNotActiveChannel, 1000 * 10, 1000 * 10, TimeUnit.MILLISECONDS); + scheduledExecutorService.scheduleWithFixedDelay( + this::scanNotActiveChannel, + Duration.ofSeconds(10).toMillis(), + Duration.ofSeconds(10).toMillis(), + TimeUnit.MILLISECONDS); this.channelManager = channelManager; this.pollCommandResponseManager = pollCommandResponseManager; + this.consumerManager = new ConsumerManager((event, group, args) -> { + // nothing to do in handler. + }); this.producerManager = new ProducerManager(); this.producerManager.setProducerOfflineListener(connectorManager.getTransactionHeartbeatRegisterService()::onProducerGroupOffline); } @@ -72,7 +78,7 @@ public class ClientService extends BaseService { String clientId = request.getClientId(); if (request.hasProducerData()) { - String producerGroup = Converter.getResourceNameWithNamespace(request.getProducerData().getGroup()); + String producerGroup = GrpcConverter.wrapResourceWithNamespace(request.getProducerData().getGroup()); GrpcClientChannel channel = GrpcClientChannel.create(channelManager, producerGroup, clientId, pollCommandResponseManager); ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal()); producerManager.registerProducer(producerGroup, clientChannelInfo); @@ -80,17 +86,17 @@ public class ClientService extends BaseService { if (request.hasConsumerData()) { ConsumerData consumerData = request.getConsumerData(); - String consumerGroup = Converter.getResourceNameWithNamespace(consumerData.getGroup()); + String consumerGroup = GrpcConverter.wrapResourceWithNamespace(consumerData.getGroup()); GrpcClientChannel channel = GrpcClientChannel.create(channelManager, consumerGroup, clientId, pollCommandResponseManager); ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal()); consumerManager.registerConsumer( consumerGroup, clientChannelInfo, - Converter.buildConsumeType(consumerData.getConsumeType()), - Converter.buildMessageModel(consumerData.getConsumeModel()), - Converter.buildConsumeFromWhere(consumerData.getConsumePolicy()), - Converter.buildSubscriptionDataSet(consumerData.getSubscriptionsList()), + GrpcConverter.buildConsumeType(consumerData.getConsumeType()), + GrpcConverter.buildMessageModel(consumerData.getConsumeModel()), + GrpcConverter.buildConsumeFromWhere(consumerData.getConsumePolicy()), + GrpcConverter.buildSubscriptionDataSet(consumerData.getSubscriptionsList()), false ); } @@ -100,7 +106,7 @@ public class ClientService extends BaseService { String clientId = request.getClientId(); if (request.hasProducerGroup()) { - String producerGroup = Converter.getResourceNameWithNamespace(request.getProducerGroup()); + String producerGroup = GrpcConverter.wrapResourceWithNamespace(request.getProducerGroup()); GrpcClientChannel channel = GrpcClientChannel.removeChannel(channelManager, producerGroup, clientId); if (channel != null) { producerManager.doChannelCloseEvent(producerGroup, channel); @@ -108,7 +114,7 @@ public class ClientService extends BaseService { } if (request.hasConsumerGroup()) { - String consumerGroup = Converter.getResourceNameWithNamespace(request.getConsumerGroup()); + String consumerGroup = GrpcConverter.wrapResourceWithNamespace(request.getConsumerGroup()); GrpcClientChannel channel = GrpcClientChannel.removeChannel(channelManager, consumerGroup, clientId); if (channel != null) { consumerManager.doChannelCloseEvent(consumerGroup, channel); @@ -118,13 +124,15 @@ public class ClientService extends BaseService { public CompletableFuture pollCommand(Context ctx, PollCommandRequest request) { CompletableFuture future = new CompletableFuture<>(); - String clientId = request.getClientId(); - PollCommandResponse noopCommandResponse = PollCommandResponse.newBuilder().setNoopCommand(NoopCommand.newBuilder().build()).build(); + PollCommandResponse noopCommandResponse = PollCommandResponse.newBuilder().setNoopCommand( + NoopCommand.newBuilder().build() + ).build(); + String clientId = request.getClientId(); switch (request.getGroupCase()) { case PRODUCER_GROUP: Resource producerGroup = request.getProducerGroup(); - String producerGroupName = Converter.getResourceNameWithNamespace(producerGroup); + String producerGroupName = GrpcConverter.wrapResourceWithNamespace(producerGroup); GrpcClientChannel producerChannel = GrpcClientChannel.getChannel(this.channelManager, producerGroupName, clientId); if (producerChannel == null) { future.complete(noopCommandResponse); @@ -134,7 +142,7 @@ public class ClientService extends BaseService { break; case CONSUMER_GROUP: Resource consumerGroup = request.getConsumerGroup(); - String consumerGroupName = Converter.getResourceNameWithNamespace(consumerGroup); + String consumerGroupName = GrpcConverter.wrapResourceWithNamespace(consumerGroup); GrpcClientChannel consumerChannel = GrpcClientChannel.getChannel(this.channelManager, consumerGroupName, clientId); if (consumerChannel == null) { future.complete(noopCommandResponse); @@ -153,7 +161,7 @@ public class ClientService extends BaseService { this.consumerManager.scanNotActiveChannel(); this.producerManager.scanNotActiveChannel(); } catch (Exception e) { - log.error("error occurred when scan not active client channels.", e); + LOGGER.error("error occurred when scan not active client channels.", e); } } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java index 0435c8f794..dfb8837799 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java @@ -34,10 +34,10 @@ import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; -import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.ProxyException; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.common.ResponseHook; +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; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class ProducerService extends BaseService { @@ -110,7 +110,7 @@ public class ProducerService extends BaseService { protected Pair convertSendMessageRequest( Context ctx, SendMessageRequest request) { - return Pair.of(Converter.buildSendMessageRequestHeader(request), Converter.buildMessage(request.getMessage())); + return Pair.of(GrpcConverter.buildSendMessageRequestHeader(request), GrpcConverter.buildMessage(request.getMessage())); } protected SendMessageResponse convertToSendMessageResponse(Context ctx, SendMessageRequest request, @@ -123,8 +123,8 @@ public class ProducerService extends BaseService { if (StringUtils.isNotBlank(sendResult.getTransactionId())) { Message message = request.getMessage(); - String group = Converter.getResourceNameWithNamespace(message.getSystemAttribute().getProducerGroup()); - String topic = Converter.getResourceNameWithNamespace(message.getTopic()); + String group = GrpcConverter.wrapResourceWithNamespace(message.getSystemAttribute().getProducerGroup()); + String topic = GrpcConverter.wrapResourceWithNamespace(message.getTopic()); this.connectorManager.getTransactionHeartbeatRegisterService().addProducerGroup(group, topic); } @@ -169,6 +169,6 @@ public class ProducerService extends BaseService { protected ConsumerSendMsgBackRequestHeader convertToConsumerSendMsgBackRequestHeader(Context ctx, ForwardMessageToDeadLetterQueueRequest request) { - return Converter.buildConsumerSendMsgBackRequestHeader(request); + return GrpcConverter.buildConsumerSendMsgBackRequestHeader(request); } } 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 6d7adab3da..93b87b9ad2 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 @@ -39,10 +39,10 @@ import org.apache.rocketmq.proxy.common.utils.ProxyUtils; 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.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.ProxyException; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.common.ResponseHook; +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; public class PullMessageService extends BaseService { @@ -66,7 +66,7 @@ public class PullMessageService extends BaseService { }); try { Partition partition = request.getPartition(); - String topic = Converter.getResourceNameWithNamespace(partition.getTopic()); + String topic = GrpcConverter.wrapResourceWithNamespace(partition.getTopic()); String brokerName = partition.getBroker().getName(); int queueId = partition.getId(); @@ -133,14 +133,14 @@ public class PullMessageService extends BaseService { protected PullMessageRequestHeader convertToPullMessageRequestHeader(Context ctx, PullMessageRequest request) { // check filterExpression is correct or not - Converter.buildSubscriptionData(Converter.getResourceNameWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); + GrpcConverter.buildSubscriptionData(GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); long pollTime = ctx.getDeadline() .timeRemaining(TimeUnit.MILLISECONDS) - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); if (pollTime <= 0) { throw new ProxyException(Code.DEADLINE_EXCEEDED, "request has been canceled due to timeout"); } - return Converter.buildPullMessageRequestHeader(request, pollTime); + return GrpcConverter.buildPullMessageRequestHeader(request, pollTime); } protected PullMessageResponse convertToPullMessageResponse(Context ctx, PullMessageRequest request, PullResult result) { @@ -150,14 +150,14 @@ public class PullMessageService extends BaseService { .setMaxOffset(result.getMaxOffset()) .setNextOffset(result.getNextBeginOffset()); - SubscriptionData subscriptionData = Converter.buildSubscriptionData( - Converter.getResourceNameWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); + SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData( + GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); PullStatus status = result.getPullStatus(); if (status.equals(PullStatus.FOUND)) { List messageList = result.getMsgFoundList().stream() .filter(msg -> FilterUtils.isTagMatched(subscriptionData.getTagsSet(), msg.getTags())) // only return tag matched messages. - .map(Converter::buildMessage) + .map(GrpcConverter::buildMessage) .collect(Collectors.toList()); return responseBuilder.addAllMessages(messageList).build(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java index 37d4acb047..6513dc68b4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java @@ -46,11 +46,11 @@ import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper; -import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.ParameterConverter; -import org.apache.rocketmq.proxy.grpc.common.ProxyMode; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.common.ResponseHook; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.adapter.ParameterConverter; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; public class RouteService extends BaseService { private final ProxyMode mode; @@ -103,7 +103,7 @@ public class RouteService extends BaseService { try { MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache() - .getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic())); + .getMessageQueue(GrpcConverter.wrapResourceWithNamespace(request.getTopic())); TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData(); List queueDataList = topicRouteData.getQueueDatas(); List brokerDataList = topicRouteData.getBrokerDatas(); @@ -218,7 +218,7 @@ public class RouteService extends BaseService { List messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request); if (ProxyMode.isLocalMode(mode)) { MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache() - .getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic())); + .getMessageQueue(GrpcConverter.wrapResourceWithNamespace(request.getTopic())); TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData(); Map> brokerMap = buildBrokerMap(topicRouteData.getBrokerDatas()); for (SelectableMessageQueue messageQueue : messageQueueList) { 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 3484ddfa01..48ffd2e402 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 @@ -34,9 +34,9 @@ import org.apache.rocketmq.proxy.connector.ForwardProducer; import org.apache.rocketmq.proxy.connector.transaction.TransactionStateCheckRequest; import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; -import org.apache.rocketmq.proxy.grpc.common.Converter; -import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.common.ResponseHook; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; +import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; public class TransactionService extends BaseService implements TransactionStateChecker { @@ -62,7 +62,7 @@ public class TransactionService extends BaseService implements TransactionStateC GrpcClientChannel channel = GrpcClientChannel.getChannel(this.channelManager, checkData.getGroupId(), clientId); String transactionId = checkData.getTransactionId().getProxyTransactionId(); - Message message = Converter.buildMessage(checkData.getMessageExt()); + Message message = GrpcConverter.buildMessage(checkData.getMessageExt()); PollCommandResponse response = PollCommandResponse.newBuilder() .setRecoverOrphanedTransactionCommand( RecoverOrphanedTransactionCommand.newBuilder() @@ -102,7 +102,7 @@ public class TransactionService extends BaseService implements TransactionStateC } protected EndTransactionRequestHeader toEndTransactionRequestHeader(Context ctx, EndTransactionRequest request) { - return Converter.buildEndTransactionRequestHeader(request); + return GrpcConverter.buildEndTransactionRequestHeader(request); } public void setCheckTransactionStateHook( diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/common/utils/FilterUtilTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/common/utils/FilterUtilTest.java index 92ad3a362a..0d36a23c73 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/common/utils/FilterUtilTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/common/utils/FilterUtilTest.java @@ -17,7 +17,6 @@ package org.apache.rocketmq.proxy.common.utils; -import java.util.concurrent.ThreadLocalRandom; import org.apache.rocketmq.common.filter.FilterAPI; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.junit.Test; @@ -26,25 +25,25 @@ import static org.assertj.core.api.Assertions.assertThat; public class FilterUtilTest { @Test - public void testIsTagMatched() throws Exception { + public void testTagMatched() throws Exception { SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA"); assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue(); } @Test - public void testIsTagNotMatched() throws Exception { + public void testTagNotMatched() throws Exception { SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA"); assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagB")).isFalse(); } @Test - public void testIsTagMatchedStar() throws Exception { + public void testTagMatchedStar() throws Exception { SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "*"); assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), "tagA")).isTrue(); } @Test - public void testIsTagNotMatchedNull() throws Exception { + public void testTagNotMatchedNull() throws Exception { SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("", "tagA"); assertThat(FilterUtils.isTagMatched(subscriptionData.getTagsSet(), null)).isFalse(); } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/config/ConfigurationManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/config/ConfigurationManagerTest.java index 669efe8ca2..ea3943a838 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/config/ConfigurationManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/config/ConfigurationManagerTest.java @@ -17,7 +17,7 @@ package org.apache.rocketmq.proxy.config; -import org.apache.rocketmq.proxy.grpc.common.ProxyMode; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; import org.junit.Test; import static org.assertj.core.api.Assertions.assertThat; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/connector/ForwardClientManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/connector/ForwardClientManagerTest.java index b2061744cb..01d89f4472 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/connector/ForwardClientManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/connector/ForwardClientManagerTest.java @@ -28,24 +28,26 @@ import static org.assertj.core.api.Assertions.assertThat; public class ForwardClientManagerTest extends InitConfigAndLoggerTest { @Test - public void testClientManager() throws Exception { + public void testConnectorManager() throws Exception { TransactionStateChecker mockedTransactionStateChecker = Mockito.mock(TransactionStateChecker.class); - ConnectorManager clientManager = new ConnectorManager(mockedTransactionStateChecker); - clientManager.start(); + ConnectorManager connectorManager = new ConnectorManager(mockedTransactionStateChecker); + connectorManager.start(); - assertThat(clientManager.getDefaultForwardClient()).isNotNull(); - assertThat(clientManager.getDefaultForwardClient().getClientNum()) + assertThat(connectorManager.getDefaultForwardClient()).isNotNull(); + assertThat(connectorManager.getDefaultForwardClient().getClientNum()) .isEqualTo(ConfigurationManager.getProxyConfig().getDefaultForwardClientNum()); - assertThat(clientManager.getForwardProducer()).isNotNull(); - assertThat(clientManager.getForwardProducer().getClientNum()) + assertThat(connectorManager.getForwardProducer()).isNotNull(); + assertThat(connectorManager.getForwardProducer().getClientNum()) .isEqualTo(ConfigurationManager.getProxyConfig().getForwardProducerNum()); - assertThat(clientManager.getForwardReadConsumer()).isNotNull(); - assertThat(clientManager.getForwardReadConsumer().getClientNum()) + assertThat(connectorManager.getForwardReadConsumer()).isNotNull(); + assertThat(connectorManager.getForwardReadConsumer().getClientNum()) .isEqualTo(ConfigurationManager.getProxyConfig().getForwardConsumerNum()); - + assertThat(connectorManager.getForwardWriteConsumer()).isNotNull(); + assertThat(connectorManager.getForwardWriteConsumer().getClientNum()) + .isEqualTo(ConfigurationManager.getProxyConfig().getForwardConsumerNum()); } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java index e473b039c9..bbe999cecb 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java @@ -77,7 +77,7 @@ import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader; import org.apache.rocketmq.common.protocol.header.PullMessageResponseHeader; import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; -import org.apache.rocketmq.proxy.grpc.common.Converter; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.remoting.exception.RemotingCommandException; import org.apache.rocketmq.remoting.protocol.RemotingCommand; @@ -263,7 +263,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); assertThat(r.getMessagesCount()).isEqualTo(1); assertThat(Durations.toMillis(r.getInvisibleDuration())).isEqualTo(invisibleTime); - assertThat(Converter.getResourceNameWithNamespace(r.getMessages(0).getTopic())).isEqualTo(topic); + assertThat(GrpcConverter.wrapResourceWithNamespace(r.getMessages(0).getTopic())).isEqualTo(topic); assertThat(r.getMessages(0).getBody().toByteArray()).isEqualTo(body); } @@ -563,7 +563,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { PullMessageResponse r = grpcFuture.get(); assertThat(r.getCommon().getStatus().getCode()).isEqualTo(Code.OK.getNumber()); assertThat(r.getMessagesCount()).isEqualTo(1); - assertThat(Converter.getResourceNameWithNamespace(r.getMessages(0).getTopic())).isEqualTo(topic); + assertThat(GrpcConverter.wrapResourceWithNamespace(r.getMessages(0).getTopic())).isEqualTo(topic); assertThat(r.getMessages(0).getBody().toByteArray()).isEqualTo(body); assertThat(r.getMinOffset()).isEqualTo(minOffset); assertThat(r.getNextOffset()).isEqualTo(nextOffset); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelectorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelectorTest.java index 21accfcf1e..72324a7a2f 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelectorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultProducerQueueSelectorTest.java @@ -12,7 +12,7 @@ import java.nio.charset.StandardCharsets; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; -import org.apache.rocketmq.proxy.grpc.common.Converter; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.junit.Test; import static org.junit.Assert.assertEquals; @@ -61,8 +61,8 @@ public class DefaultProducerQueueSelectorTest extends BaseServiceTest { .build(); WriteQueueSelector queueSelector = new DefaultWriteQueueSelector(this.topicRouteCache); SelectableMessageQueue queue = queueSelector.selectQueue(Context.current(), request, - Converter.buildSendMessageRequestHeader(request), - Converter.buildMessage(request.getMessage())); + GrpcConverter.buildSendMessageRequestHeader(request), + GrpcConverter.buildMessage(request.getMessage())); assertEquals("selectOrderQueue", queue.getBrokerName()); assertEquals("selectOrderQueueAddr", queue.getBrokerAddr()); @@ -85,8 +85,8 @@ public class DefaultProducerQueueSelectorTest extends BaseServiceTest { .build(); WriteQueueSelector queueSelector = new DefaultWriteQueueSelector(this.topicRouteCache); SelectableMessageQueue queue = queueSelector.selectQueue(Context.current(), request, - Converter.buildSendMessageRequestHeader(request), - Converter.buildMessage(request.getMessage())); + GrpcConverter.buildSendMessageRequestHeader(request), + GrpcConverter.buildMessage(request.getMessage())); assertEquals("selectOrderQueue", queue.getBrokerName()); assertEquals("selectOrderQueueAddr", queue.getBrokerAddr()); @@ -108,8 +108,8 @@ public class DefaultProducerQueueSelectorTest extends BaseServiceTest { .build(); WriteQueueSelector queueSelector = new DefaultWriteQueueSelector(this.topicRouteCache); SelectableMessageQueue queue = queueSelector.selectQueue(Context.current(), request, - Converter.buildSendMessageRequestHeader(request), - Converter.buildMessage(request.getMessage())); + GrpcConverter.buildSendMessageRequestHeader(request), + GrpcConverter.buildMessage(request.getMessage())); assertEquals("selectNormalQueue", queue.getBrokerName()); assertEquals("selectNormalQueueAddr", queue.getBrokerAddr()); @@ -136,8 +136,8 @@ public class DefaultProducerQueueSelectorTest extends BaseServiceTest { .build(); WriteQueueSelector queueSelector = new DefaultWriteQueueSelector(this.topicRouteCache); SelectableMessageQueue queue = queueSelector.selectQueue(Context.current(), request, - Converter.buildSendMessageRequestHeader(request), - Converter.buildMessage(request.getMessage())); + GrpcConverter.buildSendMessageRequestHeader(request), + GrpcConverter.buildMessage(request.getMessage())); assertEquals("selectTargetQueue", queue.getBrokerName()); assertEquals("selectTargetQueueAddr", queue.getBrokerAddr()); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java index 708d683d76..6a56bfd544 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java @@ -31,7 +31,7 @@ import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; -import org.apache.rocketmq.proxy.grpc.common.ProxyException; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyException; import org.junit.Test; import static org.junit.Assert.assertEquals; 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 73761cd9b3..4c8a3a4380 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 @@ -38,7 +38,7 @@ import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.constant.PermName; import org.apache.rocketmq.common.protocol.route.BrokerData; import org.apache.rocketmq.common.protocol.route.QueueData; -import org.apache.rocketmq.proxy.grpc.common.ProxyMode; +import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; import org.junit.Test; import static org.assertj.core.api.Assertions.assertThat;