[ISSUE #3949] Use MessageExt property to get TopicMessageType

This commit is contained in:
zhouxiang
2022-07-13 11:29:40 +08:00
parent cc827021c4
commit 8bb6af15e3
7 changed files with 26 additions and 10 deletions
@@ -20,5 +20,4 @@ package org.apache.rocketmq.proxy.common;
public class ContextVariable {
public final static String REMOTE_ADDRESS = "remote-address";
public final static String LOCAL_ADDRESS = "local-address";
public final static String MESSAGE_TYPE = "message-type";
}
@@ -23,5 +23,5 @@ public enum ProxyExceptionCode {
INVALID_RECEIPT_HANDLE,
ILLEGAL_MESSAGE,
INTERNAL_SERVER_ERROR,
TOPIC_MESSAGE_TYPE_NOT_MATCH,
MESSAGE_PROPERTY_DOES_NOT_MATCH_MESSAGE_TYPE,
}
@@ -36,6 +36,7 @@ public class GrpcProxyException extends RuntimeException {
CODE_MAPPING.put(ProxyExceptionCode.FORBIDDEN, Code.FORBIDDEN);
CODE_MAPPING.put(ProxyExceptionCode.ILLEGAL_MESSAGE, Code.ILLEGAL_MESSAGE);
CODE_MAPPING.put(ProxyExceptionCode.INTERNAL_SERVER_ERROR, Code.INTERNAL_SERVER_ERROR);
CODE_MAPPING.put(ProxyExceptionCode.MESSAGE_PROPERTY_DOES_NOT_MATCH_MESSAGE_TYPE, Code.MESSAGE_PROPERTY_DOES_NOT_MATCH_MESSAGE_TYPE);
}
public GrpcProxyException(Code code, String message) {
@@ -45,10 +45,7 @@ import org.apache.rocketmq.common.message.MessageAccessor;
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
import org.apache.rocketmq.proxy.common.ContextVariable;
import org.apache.rocketmq.proxy.common.ProxyContext;
import org.apache.rocketmq.proxy.common.ProxyException;
import org.apache.rocketmq.proxy.common.ProxyExceptionCode;
import org.apache.rocketmq.proxy.grpc.v2.AbstractMessingActivity;
import org.apache.rocketmq.proxy.grpc.v2.common.GrpcClientSettingsManager;
import org.apache.rocketmq.proxy.grpc.v2.common.GrpcConverter;
@@ -85,7 +82,7 @@ public class SendMessageActivity extends AbstractMessingActivity {
List<Message> messageList = request.getMessagesList();
Resource topic = messageList.get(0).getTopic();
future = this.messagingProcessor.sendMessage(
context.withVal(ContextVariable.MESSAGE_TYPE, topicMessageType.getValue()),
context,
new SendMessageQueueSelector(request),
GrpcConverter.wrapResourceWithNamespace(topic),
buildMessage(context, request.getMessagesList(), topic)
@@ -16,7 +16,10 @@
*/
package org.apache.rocketmq.proxy.processor;
import org.apache.rocketmq.common.attribute.TopicMessageType;
import org.apache.rocketmq.common.consumer.ReceiptHandle;
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.proxy.common.ProxyException;
import org.apache.rocketmq.proxy.common.ProxyExceptionCode;
import org.apache.rocketmq.proxy.service.ServiceManager;
@@ -37,4 +40,20 @@ public abstract class AbstractProcessor {
throw new ProxyException(ProxyExceptionCode.RECEIPT_HANDLE_EXPIRED, "receipt handle is expired");
}
}
protected TopicMessageType parseFromMessageExt(MessageExt messageExt) {
String isTrans = messageExt.getProperty(MessageConst.PROPERTY_TRANSACTION_PREPARED);
String isTransValue = "true";
if (isTransValue.equals(isTrans)) {
return TopicMessageType.TRANSACTION;
} else if (messageExt.getProperty(MessageConst.PROPERTY_DELAY_TIME_LEVEL) != null
|| messageExt.getProperty(MessageConst.PROPERTY_TIMER_DELIVER_MS) != null
|| messageExt.getProperty(MessageConst.PROPERTY_TIMER_DELAY_SEC) != null) {
return TopicMessageType.DELAY;
} else if (messageExt.getProperty(MessageConst.PROPERTY_SHARDING_KEY) != null) {
return TopicMessageType.FIFO;
} else {
return TopicMessageType.NORMAL;
}
}
}
@@ -33,7 +33,6 @@ import org.apache.rocketmq.common.protocol.ResponseCode;
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader;
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
import org.apache.rocketmq.proxy.common.ContextVariable;
import org.apache.rocketmq.proxy.common.ProxyContext;
import org.apache.rocketmq.proxy.common.ProxyException;
import org.apache.rocketmq.proxy.common.ProxyExceptionCode;
@@ -62,11 +61,12 @@ public class ProducerProcessor extends AbstractProcessor {
String producerGroup, List<MessageExt> messageExtList, long timeoutMillis) {
CompletableFuture<List<SendResult>> future = new CompletableFuture<>();
try {
String topic = messageExtList.get(0).getTopic();
MessageExt messageExt0 = messageExtList.get(0);
String topic = messageExt0.getTopic();
if (ConfigurationManager.getProxyConfig().isEnableTopicMessageTypeCheck()) {
if (topicMessageTypeValidator != null) {
TopicMessageType topicMessageType = serviceManager.getMetadataService().getTopicMessageType(topic);
TopicMessageType messageType = TopicMessageType.valueOf(ctx.getVal(ContextVariable.MESSAGE_TYPE));
TopicMessageType messageType = parseFromMessageExt(messageExt0);
topicMessageTypeValidator.validate(topicMessageType, messageType);
}
}
@@ -25,7 +25,7 @@ public class DefaultTopicMessageTypeValidator implements TopicMessageTypeValidat
public void validate(TopicMessageType topicMessageType, TopicMessageType messageType) {
if (messageType.equals(TopicMessageType.UNSPECIFIED) || !messageType.equals(topicMessageType)) {
throw new ProxyException(ProxyExceptionCode.TOPIC_MESSAGE_TYPE_NOT_MATCH, messageType.name() + " " + topicMessageType.name());
throw new ProxyException(ProxyExceptionCode.MESSAGE_PROPERTY_DOES_NOT_MATCH_MESSAGE_TYPE, messageType.name() + " " + topicMessageType.name());
}
}
}