[ISSUE #3949] Move common package, using adapter.

This commit is contained in:
Jixiang.jjx
2022-07-13 11:29:13 +08:00
committed by zhouxiang
parent 76cf9dbdc1
commit 36b1e9fa2a
43 changed files with 252 additions and 248 deletions
@@ -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;
@@ -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<ProxyConfig> 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;
}
@@ -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";
@@ -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);
}
}
@@ -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<List<String>> 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<Long> 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<Long> 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);
}
}
@@ -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<Integer> 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<SendResult> sendMessage(String address, String brokerName, Message msg,
SendMessageRequestHeader requestHeader, long timeoutMillis) {
public CompletableFuture<SendResult> sendMessage(
String address,
String brokerName,
Message msg,
SendMessageRequestHeader requestHeader,
long timeoutMillis
) {
CompletableFuture<SendResult> future = this.getClient().sendMessage(address, brokerName, msg, requestHeader, timeoutMillis);
return future.thenApply(sendResult -> {
int tranType = MessageSysFlag.getTransactionValue(requestHeader.getSysFlag());
@@ -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<PopResult> popMessage(String address, String brokerName, PopMessageRequestHeader requestHeader,
long timeoutMillis) {
return getClient().popMessage(address, brokerName, requestHeader, timeoutMillis);
@@ -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<AckResult> ackMessage(
String address,
AckMessageRequestHeader requestHeader,
@@ -78,7 +78,7 @@ public abstract class AbstractClientFactory<T> {
try {
this.shutdown(v);
} catch (Exception e) {
LOGGER.warn("RocketMQClientConstructor shutdown all err.", e);
LOGGER.warn("try to shutdown client err.", e);
}
});
}
@@ -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<MQClientAPIExt> {
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<MQCl
if (!client.updateNameServerAddressList()) {
this.scheduledExecutorService.scheduleAtFixedRate(
client::fetchNameServerAddr,
1000 * 10,
1000 * 60 * 2,
Duration.ofSeconds(10).toMillis(),
Duration.ofMinutes(2).toMillis(),
TimeUnit.MILLISECONDS
);
}
@@ -56,7 +56,7 @@ import io.grpc.Context;
import io.grpc.stub.StreamObserver;
import java.util.concurrent.CompletableFuture;
import org.apache.rocketmq.common.constant.LoggerName;
import org.apache.rocketmq.proxy.grpc.common.ResponseWriter;
import org.apache.rocketmq.proxy.grpc.adapter.ResponseWriter;
import org.apache.rocketmq.proxy.grpc.service.GrpcForwardService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -15,7 +15,7 @@
* limitations under the License.
*/
package org.apache.rocketmq.proxy.grpc.common;
package org.apache.rocketmq.proxy.grpc.adapter;
import com.google.common.base.Splitter;
import com.google.common.collect.Lists;
@@ -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.AckMessageRequest;
import apache.rocketmq.v1.ChangeInvisibleDurationRequest;
@@ -94,10 +94,10 @@ import org.apache.rocketmq.proxy.connector.transaction.TransactionId;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class Converter {
public class GrpcConverter {
private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
public static String getResourceNameWithNamespace(Resource resource) {
public static String wrapResourceWithNamespace(Resource resource) {
return NamespaceUtil.wrapNamespace(resource.getResourceNamespace(), resource.getName());
}
@@ -108,8 +108,8 @@ public class Converter {
SystemAttribute systemAttribute = message.getSystemAttribute();
Map<String, String> 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<SubscriptionData> buildSubscriptionDataSet(List<SubscriptionEntry> subscriptionEntryList) {
Set<SubscriptionData> 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;
}
@@ -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;
@@ -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;
}
@@ -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<String, PollCommandResponseFuture> futureTable = new ConcurrentHashMap<>();
public class PollResponseManager {
private final ConcurrentMap<String, PollResponseFuture> 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);
}
}
@@ -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;
@@ -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"),
@@ -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,
@@ -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;
@@ -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<T, R> {
@@ -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;
@@ -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) {
}
}
@@ -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<PullMessageRequest, PullMessageResponse> {
@Override public void handle(RemotingCommand responseCommand,
@Override
public void handle(RemotingCommand responseCommand,
InvocationContext<PullMessageRequest, PullMessageResponse> context) {
try {
PullMessageResponseHeader responseHeader = (PullMessageResponseHeader) responseCommand.readCustomHeader();
@@ -40,7 +41,7 @@ public class PullMessageResponseHandler implements ResponseHandler<PullMessageRe
ByteBuffer byteBuffer = ByteBuffer.wrap(responseCommand.getBody());
List<MessageExt> 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()))
@@ -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<ReceiveMes
Resource topic = context.getRequest()
.getPartition()
.getTopic();
String topicName = Converter.getResourceNameWithNamespace(topic);
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
messageExt.setTopic(topicName);
messageExt.setBrokerName(brokerName);
messageExt.getProperties().computeIfAbsent(MessageConst.PROPERTY_FIRST_POP_TIME,
@@ -130,7 +130,7 @@ public class ReceiveMessageResponseHandler implements ResponseHandler<ReceiveMes
}
for (MessageExt messageExt : msgFoundList) {
builder.addMessages(Converter.buildMessage(messageExt));
builder.addMessages(GrpcConverter.buildMessage(messageExt));
}
}
response = builder.build();
@@ -20,7 +20,7 @@ package org.apache.rocketmq.proxy.grpc.adapter.handler;
import apache.rocketmq.v1.SendMessageRequest;
import apache.rocketmq.v1.SendMessageResponse;
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
public class SendMessageResponseHandler implements ResponseHandler<SendMessageRequest, SendMessageResponse> {
@@ -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);
@@ -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<HeartbeatResponse> 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<SendMessageResponse> 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<AckMessageResponse> 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<QueryOffsetResponse> 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();
@@ -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 {
@@ -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) {
@@ -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<SelectableMessageQueue> 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();
}
@@ -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<PollCommandResponse> pollCommand(Context ctx, PollCommandRequest request) {
CompletableFuture<PollCommandResponse> 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);
}
}
}
@@ -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<SendMessageRequestHeader, org.apache.rocketmq.common.message.Message> 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);
}
}
@@ -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<Message> 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();
@@ -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<QueueData> queueDataList = topicRouteData.getQueueDatas();
List<BrokerData> brokerDataList = topicRouteData.getBrokerDatas();
@@ -218,7 +218,7 @@ public class RouteService extends BaseService {
List<SelectableMessageQueue> 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<String, Map<Long, Broker>> brokerMap = buildBrokerMap(topicRouteData.getBrokerDatas());
for (SelectableMessageQueue messageQueue : messageQueueList) {
@@ -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(
@@ -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();
}
@@ -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;
@@ -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());
}
}
@@ -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);
@@ -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());
@@ -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;
@@ -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;